diff --git a/FreeDMR.cfg b/FreeDMR.cfg index a6d4331..c06fe80 100644 --- a/FreeDMR.cfg +++ b/FreeDMR.cfg @@ -10,6 +10,8 @@ SERVER_ID: 0 # Passive DMR data logging: decoded content at INFO, raw bursts at DEBUG. # May expose private message or location content. DATA_PACKET_LOGGING: False +# Use the dedicated FBP5 Data Gateway link instead of the legacy D-APRS peer. +DATA_GATEWAY: False [REPORTS] @@ -19,6 +21,24 @@ DATA_PACKET_LOGGING: False [ALLSTAR] +# Native Data Gateway FBP5 relationship. The installer enables this for new +# deployments. HDStack gives each worker successive PORT and TARGET_PORT values. +[DATA-GATEWAY] +MODE: OPENBRIDGE +ENABLED: False +IP: 0.0.0.0 +PORT: 62041 +NETWORK_ID: 0 +PASSPHRASE: change-this +TARGET_IP: +TARGET_PORT: 62031 +USE_ACL: False +SUB_ACL: PERMIT:ALL +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 + # Example OpenBridge Protocol (OBP) or FreeDMR Bridge Protocol (FBP) # configuration. If you join FreeDMR, you will be given a config like this to # paste in. diff --git a/README.md b/README.md index 4a63c5e..94dde3e 100755 --- a/README.md +++ b/README.md @@ -8,4 +8,5 @@ For the single-client SQL importer and dashboard reporting relay, see [docs/reporting-relay.md](docs/reporting-relay.md). For the CPU-sized one-to-five-worker HDStack container deployment and reporting -MUX, see [docs/hdstack.md](docs/hdstack.md). +MUX, including the default separate Data Gateway service, see +[docs/hdstack.md](docs/hdstack.md). diff --git a/bridge_master.py b/bridge_master.py index bb8c013..0721203 100644 --- a/bridge_master.py +++ b/bridge_master.py @@ -1011,6 +1011,61 @@ def SubMapTrimmer(): subMapWrite() +def dataGatewayEnabled(): + """Return true when the dedicated native Data Gateway link is active.""" + return ( + CONFIG['GLOBAL'].get('DATA_GATEWAY', False) + and 'DATA-GATEWAY' in CONFIG['SYSTEMS'] + and CONFIG['SYSTEMS']['DATA-GATEWAY']['MODE'] == 'OPENBRIDGE' + and CONFIG['SYSTEMS']['DATA-GATEWAY']['ENABLED'] + and 'DATA-GATEWAY' in systems + ) + + +def isHdstackSubscriberObservation(_system, _hops, _source_server, _source_rptr): + """Return true for HBP-originated traffic relayed by this HDStack.""" + _instances = CONFIG['GLOBAL'].get('HDSTACK_INSTANCES', 0) + _base_server_id = CONFIG['GLOBAL'].get('HDSTACK_BASE_SERVER_ID', 0) + if not _instances or not _base_server_id: + return False + if _system != 'OBP-HDSTACK-AGGREGATOR': + return False + if not isinstance(_hops, bytes) or len(_hops) != 1 or int_id(_hops) != 3: + return False + if not isinstance(_source_rptr, bytes) or len(_source_rptr) != 4 or not int_id(_source_rptr): + return False + + _source_server_id = int_id(_source_server) + _local_server_id = int_id(CONFIG['GLOBAL']['SERVER_ID']) + return ( + _base_server_id < _source_server_id <= _base_server_id + _instances + and _source_server_id != _local_server_id + ) + + +def subMapTarget(_subscriber): + """Resolve a subscriber's last observed access peer on this process.""" + if _subscriber not in SUB_MAP: + return False + + (_location, _slot, _observed) = SUB_MAP[_subscriber] + if isinstance(_location, bytes) and len(_location) == 4: + for _system in systems: + _system_config = CONFIG['SYSTEMS'][_system] + if ( + _system_config['MODE'] == 'MASTER' + and _location in _system_config['PEERS'] + ): + return (_system, _slot, _observed) + return False + + # Retain old single-process persisted maps. HDStack system names are local + # to each worker and therefore cannot safely identify another worker's peer. + if not CONFIG['GLOBAL'].get('HDSTACK_INSTANCES', 0) and _location in systems: + return (_location, _slot, _observed) + return False + + # run this every 10 seconds to trim stream ids def stream_trimmer_loop(): logger.debug('(ROUTER) Trimming inactive stream IDs from system lists') @@ -2020,6 +2075,9 @@ class routerOBP(OPENBRIDGE): log_dmr_data_observation(self, self._system, 'FBP/OBP', _peer_id, _rf_src, _dst_id, _seq, _slot, _call_type, _frame_type, _dtype_vseq, _stream_id, dmrpkt) + if isHdstackSubscriberObservation(self._system, _hops, _source_server, _source_rptr): + SUB_MAP[_rf_src] = (_source_rptr, _slot, pkt_time) + #pkt_crc = Crc32.calc(_data[4:53]) #_pkt_crc = Crc32.calc(dmrpkt) @@ -2121,12 +2179,6 @@ class routerOBP(OPENBRIDGE): logger.info('(%s) *UNKNOWN DATA TYPE* STREAM ID: %s, RPTR: %s, SUB: %s (%s) PEER: %s (%s) TGID %s (%s), TS %s, SRC: %s, RPTR: %s', \ self._system, int_id(_stream_id), self.get_rptr(_source_rptr), get_alias(_rf_src, subscriber_ids), int_id(_rf_src), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot,int_id(_source_server),int_id(_source_rptr)) - #Send all data to DATA-GATEWAY if enabled and valid - if CONFIG['GLOBAL']['DATA_GATEWAY'] and 'DATA-GATEWAY' in CONFIG['SYSTEMS'] and CONFIG['SYSTEMS']['DATA-GATEWAY']['MODE'] == 'OPENBRIDGE' and CONFIG['SYSTEMS']['DATA-GATEWAY']['ENABLED']: - logger.debug('(%s) DATA packet sent to DATA-GATEWAY',self._system) - self.sendDataToOBP('DATA-GATEWAY',_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_source_rptr,_ber,_rssi) - - #Send other openbridges for system in systems: if system == self._system: @@ -2139,8 +2191,9 @@ class routerOBP(OPENBRIDGE): self.sendDataToOBP(system,_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_hops,_source_server,_ber,_rssi,_source_rptr) #If destination ID is in the Subscriber Map - if _dst_id in SUB_MAP: - (_d_system,_d_slot,_d_time) = SUB_MAP[_dst_id] + _sub_map_target = subMapTarget(_dst_id) + if _sub_map_target: + (_d_system,_d_slot,_d_time) = _sub_map_target _dst_slot = systems[_d_system].STATUS[_d_slot] logger.info('(%s) SUB_MAP matched, System: %s Slot: %s, Time: %s',self._system, _d_system,_d_slot,_d_time) #If slot is idle for RX and TX @@ -2669,7 +2722,7 @@ class routerHBP(HBSYSTEM): if CONFIG['REPORTS']['REPORT']: systems[_d_system]._report.send_bridgeEvent('UNIT DATA,DATA,TX,{},{},{},{},{},{}'.format(_d_system, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), _d_slot, _int_dst_id).encode(encoding='utf-8', errors='ignore')) - def sendDataToOBP(self,_target,_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_hops = b'',_ber = b'\x00', _rssi = b'\x00',_source_server = b'\x00\x00\x00\x00', _source_rptr = b'\x00\x00\x00\x00'): + def sendDataToOBP(self,_target,_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_hops = b'',_ber = b'\x00', _rssi = b'\x00',_source_server = b'\x00\x00\x00\x00', _source_rptr = b'\x00\x00\x00\x00', _report_data = True): # _sysIgnore = sysIgnore _source_server = self._CONFIG['GLOBAL']['SERVER_ID'] _source_rptr = _peer_id @@ -2680,9 +2733,13 @@ class routerHBP(HBSYSTEM): #We want to ignore this system and TS combination if it's called again for this packet # _sysIgnore.append((_target,_target['TS'])) + #If target has quenched us, don't send + if ('_bcsq' in _target_system) and (_dst_id in _target_system['_bcsq']) and (_target_system['_bcsq'][_dst_id] == _stream_id): + return False + #If target has missed 6 (in 1 min) of keepalives, don't send if _target_system['ENHANCED_OBP'] and ('_bcka' not in _target_system or _target_system['_bcka'] < pkt_time - 60): - return + return False if (_stream_id not in _target_status): # This is a new call stream on the target @@ -2707,10 +2764,38 @@ class routerHBP(HBSYSTEM): _tmp_data = b''.join([_data[:15], _tmp_bits.to_bytes(1, 'big'), _data[16:20]]) _tmp_data = b''.join([_tmp_data, dmrpkt]) systems[_target].send_system(_tmp_data,b'',_ber,_rssi,_source_server,_source_rptr) - logger.debug('(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s', self._system, _target,_int_dst_id) - if CONFIG['REPORTS']['REPORT']: + if _report_data: + logger.debug('(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s', self._system, _target,_int_dst_id) + if _report_data and CONFIG['REPORTS']['REPORT']: systems[_target]._report.send_bridgeEvent('UNIT DATA,DATA,TX,{},{},{},{},{},{}'.format(_target, int_id(_stream_id), int_id(_peer_id), int_id(_rf_src), 1, _int_dst_id).encode(encoding='utf-8', errors='ignore')) + return True + def sendToDataGateway(self,_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_ber,_rssi): + if not dataGatewayEnabled(): + return + if self.sendDataToOBP( + 'DATA-GATEWAY', + _data, + dmrpkt, + pkt_time, + _stream_id, + _dst_id, + _peer_id, + _rf_src, + _bits, + _slot, + _ber=_ber, + _rssi=_rssi, + _report_data=False, + ): + logger.debug( + '(%s) Locally originated HBP packet sent to DATA-GATEWAY, ' + 'STREAM ID: %s DST_ID: %s', + self._system, + int_id(_stream_id), + int_id(_dst_id), + ) + def dmrd_received(self, _peer_id, _rf_src, _dst_id, _seq, _slot, _call_type, _frame_type, _dtype_vseq, _stream_id, _data): @@ -2753,8 +2838,24 @@ class routerHBP(HBSYSTEM): _data_call = False _voice_call = False - #Add system to SUB_MAP - SUB_MAP[_rf_src] = (self._system,_slot,pkt_time) + # Record the access peer so the location remains valid across HDStack workers. + SUB_MAP[_rf_src] = (_peer_id,_slot,pkt_time) + + # The Data Gateway observes locally originated HBP traffic only. It owns + # data decoding and any consent or application policy above this link. + self.sendToDataGateway( + _data, + dmrpkt, + pkt_time, + _stream_id, + _dst_id, + _peer_id, + _rf_src, + _bits, + _slot, + _ber, + _rssi, + ) def resetallStarMode(): self.STATUS[_slot]['_allStarMode'] = False @@ -2804,11 +2905,6 @@ class routerHBP(HBSYSTEM): self._system, int_id(_stream_id), get_alias(_rf_src, subscriber_ids), int_id(_rf_src), get_alias(_peer_id, peer_ids), int_id(_peer_id), get_alias(_dst_id, talkgroup_ids), int_id(_dst_id), _slot) - #Send all data to DATA-GATEWAY if enabled and valid - if CONFIG['GLOBAL']['DATA_GATEWAY'] and 'DATA-GATEWAY' in CONFIG['SYSTEMS'] and CONFIG['SYSTEMS']['DATA-GATEWAY']['MODE'] == 'OPENBRIDGE' and CONFIG['SYSTEMS']['DATA-GATEWAY']['ENABLED']: - logger.debug('(%s) DATA packet sent to DATA-GATEWAY',self._system) - self.sendDataToOBP('DATA-GATEWAY',_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_ber=_ber,_rssi=_rssi) - #Send to all openbridges # sysIgnore = [] for system in systems: @@ -2821,8 +2917,9 @@ class routerHBP(HBSYSTEM): self.sendDataToOBP(system,_data,dmrpkt,pkt_time,_stream_id,_dst_id,_peer_id,_rf_src,_bits,_slot,_ber=_ber,_rssi=_rssi) #If destination ID is in the Subscriber Map - if _dst_id in SUB_MAP: - (_d_system,_d_slot,_d_time) = SUB_MAP[_dst_id] + _sub_map_target = subMapTarget(_dst_id) + if _sub_map_target: + (_d_system,_d_slot,_d_time) = _sub_map_target _dst_slot = systems[_d_system].STATUS[_d_slot] logger.info('(%s) SUB_MAP matched, System: %s Slot: %s, Time: %s',self._system, _d_system,_d_slot,_d_time) #If slot is idle for RX and TX @@ -2837,7 +2934,7 @@ class routerHBP(HBSYSTEM): else: logger.debug('(%s) UNIT Data not bridged to HBP on slot 1 - target busy: %s DST_ID: %s',self._system,_d_system,_int_dst_id) - elif _int_dst_id == 900999: + elif _int_dst_id == 900999 and not dataGatewayEnabled(): if 'D-APRS' in systems and CONFIG['SYSTEMS']['D-APRS']['MODE'] == 'MASTER': _d_system = 'D-APRS' _d_slot = _slot diff --git a/config.py b/config.py index 0758321..3f176da 100755 --- a/config.py +++ b/config.py @@ -194,6 +194,8 @@ def build_config(_config_file): 'ALLOW_NULL_PASSPHRASE': config.getboolean(section, 'ALLOW_NULL_PASSPHRASE', fallback=True), 'ANNOUNCEMENT_LANGUAGES': config.get(section, 'ANNOUNCEMENT_LANGUAGES', fallback=''), 'SERVER_ID': config.getint(section, 'SERVER_ID', fallback=0).to_bytes(4, 'big'), + 'HDSTACK_INSTANCES': config.getint(section, 'HDSTACK_INSTANCES', fallback=0), + 'HDSTACK_BASE_SERVER_ID': config.getint(section, 'HDSTACK_BASE_SERVER_ID', fallback=0), 'DATA_GATEWAY': config.getboolean(section, 'DATA_GATEWAY', fallback=False), 'DATA_PACKET_LOGGING': config.getboolean(section, 'DATA_PACKET_LOGGING', fallback=False), 'VALIDATE_SERVER_IDS': config.getboolean(section, 'VALIDATE_SERVER_IDS', fallback=True), diff --git a/docker-configs/docker-compose.yml b/docker-configs/docker-compose.yml index 526fb01..0e2ffd7 100644 --- a/docker-configs/docker-compose.yml +++ b/docker-configs/docker-compose.yml @@ -31,8 +31,8 @@ services: ports: - '62031:62031/udp' - #Change the below to inlude ports used for your OBP(s) - - '62041:62041/udp' + # Native Data Gateway worker ports; use other ports for external peers. + - '62041-62045:62041-62045/udp' image: 'gitlab.hacknix.net:5050/hacknix/freedmr:latest' restart: "unless-stopped" networks: @@ -63,6 +63,35 @@ services: tmpfs: /tmp read_only: "true" + + data-gateway: + container_name: freedmr-data-gateway + depends_on: + - freedmr + image: 'gitlab.hacknix.net:5050/freedmr/freedmr-data-gateway:development-latest' + restart: "unless-stopped" + read_only: true + tmpfs: + - /tmp:size=16m,mode=1777 + cap_drop: + - ALL + security_opt: + - no-new-privileges:true + environment: + FREEDMR_GATEWAY_INGRESS_MODE: fbp + FREEDMR_GATEWAY_FBP_RELATIONSHIPS: '__FREEDMR_GATEWAY_FBP_RELATIONSHIPS__' + FREEDMR_GATEWAY_RADIOID_PATH: /data/subscriber_ids.json + FREEDMR_GATEWAY_APRS_HOST: rotate.aprs2.net + FREEDMR_GATEWAY_APRS_PORT: '14580' + FREEDMR_GATEWAY_APRS_LOGIN_CALL: '__FREEDMR_GATEWAY_APRS_LOGIN_CALL__' + FREEDMR_GATEWAY_APRS_PASSCODE: '__FREEDMR_GATEWAY_APRS_PASSCODE__' + FREEDMR_GATEWAY_LOG_LEVEL: INFO + volumes: + - '/etc/freedmr/json/:/data/:ro' + networks: + app_net: + ipv4_address: 172.16.238.40 + freedmrmonitor2: container_name: freedmrmonitor2 cpu_shares: 512 diff --git a/docker-configs/docker-compose_install.sh b/docker-configs/docker-compose_install.sh index eaa064b..4524050 100644 --- a/docker-configs/docker-compose_install.sh +++ b/docker-configs/docker-compose_install.sh @@ -65,6 +65,33 @@ validate_server_id() fi } +build_data_gateway_relationships() +{ + count="$1" + base_id="$2" + key="$3" + relationships='[' + number=1 + while [ "$number" -le "$count" ] + do + if [ "$count" -eq 1 ] + then + server_id="$base_id" + else + server_id=$((base_id + number)) + fi + bind_port=$((62030 + number)) + peer_port=$((62040 + number)) + if [ "$number" -gt 1 ] + then + relationships="${relationships}," + fi + relationships="${relationships}{\"link_id\":\"hdstack-${number}\",\"bind_host\":\"0.0.0.0\",\"bind_port\":${bind_port},\"peer_host\":\"freedmr\",\"peer_port\":${peer_port},\"peer_network_id\":${server_id},\"direct_source_server_id\":${server_id},\"key\":\"${key}\"}" + number=$((number + 1)) + done + printf '%s]' "$relationships" +} + CPU_COUNT=$(getconf _NPROCESSORS_ONLN 2>/dev/null) || CPU_COUNT=2 if ! is_unsigned_integer "$CPU_COUNT" || [ "$CPU_COUNT" -lt 1 ] then @@ -144,6 +171,55 @@ then fi fi +echo +echo 'The separate Data Gateway provides D-APRS and future packet-data services.' +FREEDMR_GATEWAY_APRS_LOGIN_CALL="${FREEDMR_GATEWAY_APRS_LOGIN_CALL:-}" +if [ -z "$FREEDMR_GATEWAY_APRS_LOGIN_CALL" ] +then + FREEDMR_GATEWAY_APRS_LOGIN_CALL=$(prompt_value 'APRS-IS login callsign' '') +fi +case "$FREEDMR_GATEWAY_APRS_LOGIN_CALL" in + ''|*[!A-Za-z0-9-]*) + echo 'APRS-IS login callsign may contain only letters, digits and a hyphen.' >&2 + exit 2 + ;; +esac +if [ "${#FREEDMR_GATEWAY_APRS_LOGIN_CALL}" -gt 9 ] +then + echo 'APRS-IS login callsign must not exceed 9 characters.' >&2 + exit 2 +fi + +FREEDMR_GATEWAY_APRS_PASSCODE="${FREEDMR_GATEWAY_APRS_PASSCODE:-}" +if [ -z "$FREEDMR_GATEWAY_APRS_PASSCODE" ] +then + FREEDMR_GATEWAY_APRS_PASSCODE=$(prompt_value 'APRS-IS passcode' '') +fi +if ! is_unsigned_integer "$FREEDMR_GATEWAY_APRS_PASSCODE" +then + echo 'APRS-IS passcode must be numeric.' >&2 + exit 2 +fi + +FREEDMR_GATEWAY_FBP_KEY="${FREEDMR_GATEWAY_FBP_KEY:-}" +if [ -z "$FREEDMR_GATEWAY_FBP_KEY" ] +then + FREEDMR_GATEWAY_FBP_KEY=$(od -An -N10 -tx1 /dev/urandom | tr -d ' \n') +fi +case "$FREEDMR_GATEWAY_FBP_KEY" in + ''|*[!A-Za-z0-9._-]*) + echo 'Data Gateway FBP key may contain only letters, digits, dot, underscore and hyphen.' >&2 + exit 2 + ;; +esac +if [ "${#FREEDMR_GATEWAY_FBP_KEY}" -gt 20 ] +then + echo 'Data Gateway FBP key must not exceed 20 characters.' >&2 + exit 2 +fi +FREEDMR_GATEWAY_FBP_RELATIONSHIPS=$(build_data_gateway_relationships \ + "$HDSTACK_COUNT" "$HDSTACK_BASEID" "$FREEDMR_GATEWAY_FBP_KEY") + echo echo 'Installing required packages and Docker Community Edition...' apt-get -y remove docker docker-engine docker.io @@ -183,7 +259,12 @@ chmod -R 755 /etc/freedmr chown 54000:54000 /etc/freedmr/json curl -fsSL "$BASE_URL/FreeDMR.cfg" -o /etc/freedmr/freedmr.cfg -sed -i "0,/^SERVER_ID:/s//SERVER_ID: $HDSTACK_BASEID/" /etc/freedmr/freedmr.cfg +sed -i "0,/^SERVER_ID:/s/^SERVER_ID:.*/SERVER_ID: $HDSTACK_BASEID/" /etc/freedmr/freedmr.cfg +sed -i '/^\[GLOBAL\]/,/^\[/ s/^DATA_GATEWAY:.*/DATA_GATEWAY: True/' /etc/freedmr/freedmr.cfg +sed -i '/^\[DATA-GATEWAY\]/,/^\[/ s/^ENABLED:.*/ENABLED: True/' /etc/freedmr/freedmr.cfg +sed -i "/^\[DATA-GATEWAY\]/,/^\[/ s|^PASSPHRASE:.*|PASSPHRASE: $FREEDMR_GATEWAY_FBP_KEY|" /etc/freedmr/freedmr.cfg +sed -i '/^\[DATA-GATEWAY\]/,/^\[/ s/^TARGET_IP:.*/TARGET_IP: freedmr-data-gateway/' /etc/freedmr/freedmr.cfg + if [ ! -e /etc/freedmr/freedmr-bridge.cfg ] then @@ -191,7 +272,7 @@ then fi if [ "$HDSTACK_BRIDGE_VALUE" -eq 1 ] then - sed -i "0,/^SERVER_ID:/s//SERVER_ID: $HDSTACK_BRIDGE_ID/" /etc/freedmr/freedmr-bridge.cfg + sed -i "0,/^SERVER_ID:/s/^SERVER_ID:.*/SERVER_ID: $HDSTACK_BRIDGE_ID/" /etc/freedmr/freedmr-bridge.cfg fi if [ ! -e /etc/freedmr/rules-bridge.py ] then @@ -202,6 +283,9 @@ curl -fsSL "$BASE_URL/docker-configs/docker-compose.yml" -o /etc/freedmr/docker- sed -i "s|^[[:space:]]*#*- HDSTACK=.*| - HDSTACK=$HDSTACK_COUNT|" /etc/freedmr/docker-compose.yml sed -i "s|^[[:space:]]*#*- HDSTACK_BASEID=.*| - HDSTACK_BASEID=$HDSTACK_BASEID|" /etc/freedmr/docker-compose.yml sed -i "s|^[[:space:]]*- HDSTACK_BRIDGE=.*| - HDSTACK_BRIDGE=$HDSTACK_BRIDGE_VALUE|" /etc/freedmr/docker-compose.yml +sed -i "s|__FREEDMR_GATEWAY_FBP_RELATIONSHIPS__|$FREEDMR_GATEWAY_FBP_RELATIONSHIPS|" /etc/freedmr/docker-compose.yml +sed -i "s|__FREEDMR_GATEWAY_APRS_LOGIN_CALL__|$FREEDMR_GATEWAY_APRS_LOGIN_CALL|" /etc/freedmr/docker-compose.yml +sed -i "s|__FREEDMR_GATEWAY_APRS_PASSCODE__|$FREEDMR_GATEWAY_APRS_PASSCODE|" /etc/freedmr/docker-compose.yml chown -R 54000:54000 /etc/freedmr if [ -e /etc/cron.daily/lastheard ] diff --git a/docker-configs/freedmr.cfg b/docker-configs/freedmr.cfg index a6d4331..c06fe80 100644 --- a/docker-configs/freedmr.cfg +++ b/docker-configs/freedmr.cfg @@ -10,6 +10,8 @@ SERVER_ID: 0 # Passive DMR data logging: decoded content at INFO, raw bursts at DEBUG. # May expose private message or location content. DATA_PACKET_LOGGING: False +# Use the dedicated FBP5 Data Gateway link instead of the legacy D-APRS peer. +DATA_GATEWAY: False [REPORTS] @@ -19,6 +21,24 @@ DATA_PACKET_LOGGING: False [ALLSTAR] +# Native Data Gateway FBP5 relationship. The installer enables this for new +# deployments. HDStack gives each worker successive PORT and TARGET_PORT values. +[DATA-GATEWAY] +MODE: OPENBRIDGE +ENABLED: False +IP: 0.0.0.0 +PORT: 62041 +NETWORK_ID: 0 +PASSPHRASE: change-this +TARGET_IP: +TARGET_PORT: 62031 +USE_ACL: False +SUB_ACL: PERMIT:ALL +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 + # Example OpenBridge Protocol (OBP) or FreeDMR Bridge Protocol (FBP) # configuration. If you join FreeDMR, you will be given a config like this to # paste in. diff --git a/docs/hdstack.md b/docs/hdstack.md index 001e51a..5007d51 100644 --- a/docs/hdstack.md +++ b/docs/hdstack.md @@ -15,6 +15,7 @@ multi-process layout: hotspots -> proxy -> HBP workers -> no-HBP aggregator -> external FBP/OBP | | +-> Loro +-> optional bridge.py over FBP v5/FBCP + +-> Data Gateway over FBP v5/FBCP worker reports -> reporting MUX -> existing dashboard/reporting consumer ``` @@ -77,14 +78,16 @@ environment: The main `freedmr.cfg` remains the sysop-facing aggregator and worker-template configuration. It must contain one enabled `[SYSTEM]` stanza in `MASTER` mode. -All enabled and disabled `OPENBRIDGE` sections are copied to the aggregator; -this includes external FBP/OBP links. Other protocol systems, including -`[SYSTEM]` and `[ECHO]`, are not copied to the aggregator. +All enabled and disabled `OPENBRIDGE` sections except `[DATA-GATEWAY]` are +copied to the aggregator; this includes external FBP/OBP links. Other protocol +systems, including `[SYSTEM]` and `[ECHO]`, are not copied to the aggregator. Each worker receives the same `[SYSTEM]` policy. Only its base HBP port is -derived. Worker and aggregator configuration files are generated privately in -`/dev/shm/freedmr-hdstack` when the container starts. The mounted source -configuration and bridge files are not changed. +derived. When the native Data Gateway is enabled, each worker also receives a +direct `[DATA-GATEWAY]` FBP v5/FBCP relationship. Worker and aggregator +configuration files are generated privately in `/dev/shm/freedmr-hdstack` when +the container starts. The mounted source configuration and bridge files are +not changed. When `[ECHO]` is enabled, each worker receives its own generated ECHO peer and Loro playback process on a private loopback port pair. Calls cannot leak @@ -113,6 +116,10 @@ worker 5 -> 54400..54499 The generated ranges must fit within the UDP port range and must not collide with configured bridge, external OpenBridge or reporting ports. +The internal worker path consumes one additional FBP hop. FreeDMR therefore +accepts a maximum of 11 FBP hops before discarding traffic; this ceiling applies +to all FBP relationships rather than being a special HDStack exception. + Workers connect to the aggregator over generated loopback FBP v5 links with FBCP enabled. FBP v5 preserves origin metadata and carries unit data as well as group traffic; FBCP provides source quench and the existing link-control @@ -134,6 +141,41 @@ relationships remain on the aggregator. Each generated Loro pair uses successive private ports beginning with the existing `[ECHO]` peer and master ports, normally `54916` and `54915`. +## Data Gateway + +New installations include a separate FreeDMR Data Gateway container for +D-APRS and future packet-data services. The installer asks for the APRS-IS +login callsign and passcode, generates a private FBP key, enables the +`[DATA-GATEWAY]` relationship, and writes matching relationships into the +Compose service. Existing configurations keep `DATA_GATEWAY: False` and the +disabled section by default, so replacing only the FreeDMR image does not +remove the legacy D-APRS path. + +Each worker connects directly to the gateway. The aggregator and optional +`bridge.py` process do not. The default ports are: + +```text +worker 1 62041 -> gateway 62031 +worker 2 62042 -> gateway 62032 +worker 3 62043 -> gateway 62033 +worker 4 62044 -> gateway 62034 +worker 5 62045 -> gateway 62035 +``` + +FreeDMR sends every admitted, locally originated HBP DMR packet to this link, +including unit data, group data and voice packets. Peer or mesh ingress is not +copied back to the gateway. The complete 33-byte DMR burst is preserved; the +gateway owns decoding, APRS consent and application policy. Enabling the native +gateway disables legacy destination-`900999` D-APRS delivery to avoid duplicate +processing. It is not an automatic fallback if the native link is unavailable. + +The default installation runs the gateway locally, but this is not a required +one-gateway-per-server architecture. A sysop may point each worker relationship +at a shared remote gateway by configuring matching contiguous target ports, +peer IDs and keys, then removing or disabling the local Compose service. The +gateway endpoint must be reachable from the FreeDMR container; `localhost` +would address the FreeDMR container itself, not a separate gateway container. + ## Session assignment The proxy rotates new DMR IDs through the configured worker ranges. With two @@ -148,6 +190,19 @@ new session 4 -> worker 2 Within the selected worker range, the proxy retains its existing random choice of an available generated system port. +Subscriber location remains best-effort state. A worker records the originating +HBP repeater or hotspot ID, slot and observation time rather than a worker-local +generated system name. Authenticated FBP v5 traffic received through the +private aggregator relationship also updates this map when its source server +and hop path identify another worker in the same HDStack deployment. + +Each worker resolves a unit-data destination against its currently connected +HBP peers before transmitting. This allows a hotspot or repeater to reconnect +on another worker without requiring the subscriber to transmit again first, +provided the subscriber-to-access-peer observation has already propagated. +Workers retain separate persistence files; they do not concurrently write one +shared pickle file. + A DMR ID remains pinned to that exact generated port for the lifetime of its active proxy session. Moving it would start a new HBP login and reset its diff --git a/freedmr.cfg b/freedmr.cfg index a6d4331..c06fe80 100644 --- a/freedmr.cfg +++ b/freedmr.cfg @@ -10,6 +10,8 @@ SERVER_ID: 0 # Passive DMR data logging: decoded content at INFO, raw bursts at DEBUG. # May expose private message or location content. DATA_PACKET_LOGGING: False +# Use the dedicated FBP5 Data Gateway link instead of the legacy D-APRS peer. +DATA_GATEWAY: False [REPORTS] @@ -19,6 +21,24 @@ DATA_PACKET_LOGGING: False [ALLSTAR] +# Native Data Gateway FBP5 relationship. The installer enables this for new +# deployments. HDStack gives each worker successive PORT and TARGET_PORT values. +[DATA-GATEWAY] +MODE: OPENBRIDGE +ENABLED: False +IP: 0.0.0.0 +PORT: 62041 +NETWORK_ID: 0 +PASSPHRASE: change-this +TARGET_IP: +TARGET_PORT: 62031 +USE_ACL: False +SUB_ACL: PERMIT:ALL +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 + # Example OpenBridge Protocol (OBP) or FreeDMR Bridge Protocol (FBP) # configuration. If you join FreeDMR, you will be given a config like this to # paste in. diff --git a/hblink.py b/hblink.py index 5023d6a..d2292ab 100755 --- a/hblink.py +++ b/hblink.py @@ -52,6 +52,7 @@ import pickle from reporting_const import * MAX_OPTIONS_LENGTH = 255 +MAX_FBP_HOPS = 11 # The module needs logging logging, but handlers, etc. are controlled by the parent import logging @@ -491,7 +492,7 @@ class OPENBRIDGE(DatagramProtocol): #Increment max hops _inthops = _hops +1 - if _inthops > 10: + if _inthops > MAX_FBP_HOPS: logger.warning('(%s) MAX HOPS exceed, dropping. Hops: %s, DST: %s, SRC: %s', self._system, _inthops, _int_dst_id, int.from_bytes(_source_server,'big')) self.send_bcsq(_dst_id,_stream_id) return diff --git a/hdstack_runtime.py b/hdstack_runtime.py index 56a78fe..c08ef88 100644 --- a/hdstack_runtime.py +++ b/hdstack_runtime.py @@ -142,6 +142,10 @@ def generate_hdstack( for port in worker_report_ports: _validate_port(port, 'worker report port') + data_gateway_enabled, data_gateway_ports, data_gateway_target_ports = ( + _data_gateway_layout(source, instances) + ) + _validate_generated_port_collisions( source, bridge_source, @@ -153,14 +157,19 @@ def generate_hdstack( bridge_enabled, loro_port_pairs, loro_report_ports, + data_gateway_ports, + data_gateway_target_ports, ) aggregator = _support_config(source) aggregator['GLOBAL']['SERVER_ID'] = str(base_server_id) + aggregator['GLOBAL']['HDSTACK_INSTANCES'] = str(instances) + aggregator['GLOBAL']['HDSTACK_BASE_SERVER_ID'] = str(base_server_id) aggregator['GLOBAL']['GEN_STAT_BRIDGES'] = 'True' aggregator['REPORTS']['REPORT'] = 'False' + aggregator['GLOBAL']['DATA_GATEWAY'] = 'False' for section in source.sections(): - if _mode(source, section) == 'OPENBRIDGE': + if section != 'DATA-GATEWAY' and _mode(source, section) == 'OPENBRIDGE': _copy_section(source, aggregator, section) for number in range(1, instances + 1): name = 'OBP-HDSTACK-WORKER-{}'.format(number) @@ -209,9 +218,11 @@ def generate_hdstack( for number, (start, _end) in enumerate(backend_ranges, 1): worker = _support_config(source) worker['GLOBAL']['SERVER_ID'] = str(base_server_id + number) + worker['GLOBAL']['HDSTACK_INSTANCES'] = str(instances) + worker['GLOBAL']['HDSTACK_BASE_SERVER_ID'] = str(base_server_id) worker['GLOBAL']['GEN_STAT_BRIDGES'] = 'False' worker['GLOBAL']['ENABLE_API'] = 'False' - worker['GLOBAL']['DATA_GATEWAY'] = 'False' + worker['GLOBAL']['DATA_GATEWAY'] = str(data_gateway_enabled) worker['ALIASES']['TRY_DOWNLOAD'] = 'False' _number_mutable_alias_file(worker, 'SUB_MAP_FILE', number) _number_mutable_alias_file(worker, 'KEYS_FILE', number) @@ -220,6 +231,11 @@ def generate_hdstack( worker['REPORTS']['REPORT_CLIENTS'] = '127.0.0.1' _copy_section(source, worker, 'SYSTEM') worker['SYSTEM']['PORT'] = str(start) + if data_gateway_enabled: + _copy_section(source, worker, 'DATA-GATEWAY') + worker['DATA-GATEWAY']['PORT'] = str(data_gateway_ports[number - 1]) + worker['DATA-GATEWAY']['TARGET_PORT'] = str(data_gateway_target_ports[number - 1]) + name = 'OBP-HDSTACK-AGGREGATOR' if echo_enabled: echo_peer_port, echo_master_port = loro_port_pairs[number - 1] @@ -476,6 +492,34 @@ def _enabled_bind_ports(source, excluded_sections=()): return ports +def _data_gateway_layout(source, instances): + if not _boolean(source, 'GLOBAL', 'DATA_GATEWAY', False): + return False, (), () + + if ( + not source.has_section('DATA-GATEWAY') + or not _boolean(source, 'DATA-GATEWAY', 'ENABLED', True) + ): + return False, (), () + + if _mode(source, 'DATA-GATEWAY') != 'OPENBRIDGE': + raise HDStackConfigurationError('[DATA-GATEWAY] must use MODE OPENBRIDGE') + if _section_integer(source, 'DATA-GATEWAY', 'PROTO_VER', 5) != 5: + raise HDStackConfigurationError('[DATA-GATEWAY] must use PROTO_VER 5') + if not _boolean(source, 'DATA-GATEWAY', 'ENHANCED_OBP', True): + raise HDStackConfigurationError('[DATA-GATEWAY] must enable FBCP') + + bind_port = _section_integer(source, 'DATA-GATEWAY', 'PORT', 0) + target_port = _section_integer(source, 'DATA-GATEWAY', 'TARGET_PORT', 0) + bind_ports = tuple( + bind_port + number - 1 for number in range(1, instances + 1) + ) + target_ports = tuple( + target_port + number - 1 for number in range(1, instances + 1) + ) + return True, bind_ports, target_ports + + def _validate_generated_port_collisions( source, bridge_source, @@ -487,6 +531,8 @@ def _validate_generated_port_collisions( bridge_enabled, loro_port_pairs, loro_report_ports, + data_gateway_ports, + data_gateway_target_ports, ): generated = {} @@ -511,6 +557,12 @@ def _validate_generated_port_collisions( INTERNAL_AGGREGATOR_PORT_BASE + number, 'aggregator worker {} FBP'.format(number), ) + for number, (port, target_port) in enumerate( + zip(data_gateway_ports, data_gateway_target_ports), + 1, + ): + reserve(port, 'worker {} DATA-GATEWAY'.format(number)) + _validate_port(target_port, 'worker {} DATA-GATEWAY target'.format(number)) if bridge_enabled: reserve(INTERNAL_AGGREGATOR_BRIDGE_PORT, 'aggregator bridge FBP') reserve(INTERNAL_BRIDGE_PORT, 'bridge aggregator FBP') @@ -524,7 +576,10 @@ def _validate_generated_port_collisions( for number, port in enumerate(worker_report_ports, 1): reserve(port, 'worker {} reporting'.format(number)) - configured = _enabled_bind_ports(source, excluded_sections=('SYSTEM', 'ECHO')) + configured = _enabled_bind_ports( + source, + excluded_sections=('SYSTEM', 'ECHO', 'DATA-GATEWAY'), + ) bridge_ports = {} if bridge_enabled: bridge_ports = _enabled_bind_ports(bridge_source) diff --git a/tests/harness/udp_blackbox.py b/tests/harness/udp_blackbox.py index c4fb3c2..b063629 100644 --- a/tests/harness/udp_blackbox.py +++ b/tests/harness/udp_blackbox.py @@ -506,11 +506,13 @@ def write_bridge_master_config( network_direct_dial_slot: int = 1, data_packet_logging: bool = False, master_extra_config: str = "", + data_gateway: bool = False, ) -> None: global_use_acl_text = "True" if global_use_acl else "False" dial_a_tg_text = "True" if dial_a_tg else "False" dynamic_tg_routing_text = "True" if dynamic_tg_routing else "False" data_packet_logging_text = "True" if data_packet_logging else "False" + data_gateway_text = "True" if data_gateway else "False" systems = [] for name, port in system_ports.items(): systems.append( @@ -587,7 +589,7 @@ GEN_STAT_BRIDGES: False ALLOW_NULL_PASSPHRASE: True ANNOUNCEMENT_LANGUAGES: en_GB SERVER_ID: 9990 -DATA_GATEWAY: False +DATA_GATEWAY: {data_gateway_text} DATA_PACKET_LOGGING: {data_packet_logging_text} VALIDATE_SERVER_IDS: False DEBUG_BRIDGES: False @@ -1075,6 +1077,7 @@ class UdpBlackBoxScenario: network_direct_dial_slot: int = 1, data_packet_logging: bool = False, master_extra_config: str = "", + data_gateway: bool = False, fbp_systems: dict[str, int] | None = None, fbp_proto_versions: dict[str, int] | None = None, enable_proxy: bool = False, @@ -1124,6 +1127,7 @@ class UdpBlackBoxScenario: "dynamic_tg_routing": dynamic_tg_routing, "network_direct_dial_slot": network_direct_dial_slot, "data_packet_logging": data_packet_logging, + "data_gateway": data_gateway, "master_extra_config": master_extra_config, } diff --git a/tests/test_deterministic_harness.py b/tests/test_deterministic_harness.py index 2a51dbd..78cf94e 100644 --- a/tests/test_deterministic_harness.py +++ b/tests/test_deterministic_harness.py @@ -3917,6 +3917,319 @@ NETWORK_DIRECT_DIAL_SLOT: {value} self.assertEqual(scenario.capture.for_system("OBP-1"), []) + def test_native_data_gateway_observes_all_local_hbp_packet_classes(self): + config = minimal_config(("MASTER-A",)) + config["GLOBAL"]["DATA_GATEWAY"] = True + add_openbridge_system(config, "DATA-GATEWAY", network_id=0) + gateway = config["SYSTEMS"]["DATA-GATEWAY"] + gateway["ENHANCED_OBP"] = True + + packets = ( + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + slot=2, + stream_id=0x01020310, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + payload=b"\x11" * 33, + ber=b"\x01", + rssi=b"\x02", + ), + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=9, + slot=2, + stream_id=0x01020311, + call_type="group", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + payload=b"\x22" * 33, + ), + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=9, + slot=2, + stream_id=0x01020312, + call_type="group", + frame_type=HBPF_VOICE, + dtype_vseq=0, + payload=b"\x33" * 33, + ), + ) + + with DeterministicScenario(config=config) as scenario: + scenario.register_peer("MASTER-A", 4000001) + gateway["_bcka"] = scenario.clock.time() + for packet in packets: + scenario.inject_hbp("MASTER-A", packet) + + captured = scenario.capture.for_system("DATA-GATEWAY") + self.assertEqual(len(captured), 3) + self.assertEqual( + [item.fields["dmr_payload"] for item in captured], + [packet.payload for packet in packets], + ) + self.assertEqual( + [item.fields["call_type"] for item in captured], + ["unit", "group", "group"], + ) + self.assertTrue(all(item.fields["slot"] == 1 for item in captured)) + self.assertTrue(all(item.hops == b"" for item in captured)) + self.assertTrue(all(item.source_server == bytes_4(9990) for item in captured)) + self.assertTrue(all(item.source_rptr == bytes_4(4000001) for item in captured)) + self.assertEqual(captured[0].ber, b"\x01") + self.assertEqual(captured[0].rssi, b"\x02") + + def test_native_data_gateway_does_not_observe_fbp_ingress(self): + config = minimal_config(("MASTER-A",)) + config["GLOBAL"]["DATA_GATEWAY"] = True + add_openbridge_system(config, "OBP-1", network_id=3001) + add_openbridge_system(config, "DATA-GATEWAY", network_id=0) + + with DeterministicScenario(config=config) as scenario: + packet = PacketSpec( + peer_id=3001, + rf_src=1234567, + dst_id=900999, + stream_id=0x01020320, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + + scenario.inject_obp("OBP-1", packet) + + self.assertEqual(scenario.capture.for_system("DATA-GATEWAY"), []) + + def test_native_data_gateway_honours_existing_fbcp_source_quench(self): + config = minimal_config(("MASTER-A",)) + config["GLOBAL"]["DATA_GATEWAY"] = True + add_openbridge_system(config, "DATA-GATEWAY", network_id=0) + gateway = config["SYSTEMS"]["DATA-GATEWAY"] + gateway["ENHANCED_OBP"] = True + quenched_stream = bytes_4(0x01020321) + gateway["_bcsq"] = {bytes_3(900999): quenched_stream} + + with DeterministicScenario(config=config) as scenario: + scenario.register_peer("MASTER-A", 4000001) + gateway["_bcka"] = scenario.clock.time() + scenario.inject_hbp( + "MASTER-A", + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + stream_id=quenched_stream, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ), + ) + self.assertEqual(scenario.capture.for_system("DATA-GATEWAY"), []) + + scenario.inject_hbp( + "MASTER-A", + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + stream_id=0x01020322, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ), + ) + self.assertEqual(len(scenario.capture.for_system("DATA-GATEWAY")), 1) + + def test_legacy_daprs_is_used_only_without_native_data_gateway(self): + legacy_config = minimal_config(("MASTER-A", "D-APRS")) + + with DeterministicScenario(config=legacy_config) as scenario: + scenario.register_peer("MASTER-A", 4000001) + scenario.clock.advance(1) + scenario.inject_hbp( + "MASTER-A", + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + stream_id=0x01020330, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ), + ) + self.assertEqual(len(scenario.capture.for_system("D-APRS")), 1) + + native_config = minimal_config(("MASTER-A", "D-APRS")) + native_config["GLOBAL"]["DATA_GATEWAY"] = True + add_openbridge_system(native_config, "DATA-GATEWAY", network_id=0) + gateway = native_config["SYSTEMS"]["DATA-GATEWAY"] + gateway["ENHANCED_OBP"] = True + + with DeterministicScenario(config=native_config) as scenario: + scenario.register_peer("MASTER-A", 4000001) + scenario.clock.advance(1) + gateway["_bcka"] = scenario.clock.time() + scenario.inject_hbp( + "MASTER-A", + PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + stream_id=0x01020331, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ), + ) + self.assertEqual(len(scenario.capture.for_system("DATA-GATEWAY")), 1) + self.assertEqual(scenario.capture.for_system("D-APRS"), []) + + def test_hbp_sub_map_records_the_connected_access_peer(self): + config = minimal_config(("MASTER-A",)) + + with DeterministicScenario(config=config) as scenario: + access_peer = scenario.register_peer("MASTER-A", 4000001) + packet = PacketSpec( + peer_id=4000001, + rf_src=1234567, + dst_id=900999, + slot=2, + stream_id=0x01020304, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + + scenario.inject_hbp("MASTER-A", packet) + + self.assertEqual( + scenario.bm.SUB_MAP[bytes_3(1234567)][:2], + (access_peer, 2), + ) + + def test_hdstack_worker_learns_sibling_hbp_location_and_routes_to_connected_peer(self): + config = minimal_config(("MASTER-B",)) + config["GLOBAL"]["SERVER_ID"] = bytes_4(23402) + config["GLOBAL"]["HDSTACK_INSTANCES"] = 2 + config["GLOBAL"]["HDSTACK_BASE_SERVER_ID"] = 23400 + add_openbridge_system( + config, + "OBP-HDSTACK-AGGREGATOR", + network_id=23400, + ) + + with DeterministicScenario(config=config) as scenario: + access_peer = scenario.register_peer("MASTER-B", 4000001) + scenario.clock.advance(10) + observation = PacketSpec( + peer_id=23400, + rf_src=1234567, + dst_id=900999, + slot=2, + stream_id=0x01020304, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + internal = scenario.systems["OBP-HDSTACK-AGGREGATOR"] + + internal.dmrd_received( + *observation.decoded_obp_args( + hops=b"\x03", + source_server=23401, + source_rptr=access_peer, + ) + ) + + self.assertEqual( + scenario.bm.SUB_MAP[bytes_3(1234567)][:2], + (access_peer, 2), + ) + + inbound = PacketSpec( + peer_id=23400, + rf_src=7654321, + dst_id=1234567, + slot=1, + stream_id=0x01020305, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + internal.dmrd_received( + *inbound.decoded_obp_args( + hops=b"\x03", + source_server=23401, + source_rptr=4000002, + ) + ) + + captured = scenario.capture.for_system("MASTER-B") + self.assertEqual(len(captured), 1) + self.assertEqual(captured[0].fields["dst_id"], bytes_3(1234567)) + self.assertEqual(captured[0].fields["slot"], 2) + + del config["SYSTEMS"]["MASTER-B"]["PEERS"][access_peer] + later = PacketSpec( + peer_id=23400, + rf_src=7654321, + dst_id=1234567, + slot=1, + stream_id=0x01020306, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + internal.dmrd_received( + *later.decoded_obp_args( + hops=b"\x03", + source_server=23401, + source_rptr=4000002, + ) + ) + + self.assertEqual(len(scenario.capture.for_system("MASTER-B")), 1) + + def test_hdstack_worker_ignores_non_sibling_subscriber_observations(self): + config = minimal_config(("MASTER-B",)) + config["GLOBAL"]["SERVER_ID"] = bytes_4(23402) + config["GLOBAL"]["HDSTACK_INSTANCES"] = 2 + config["GLOBAL"]["HDSTACK_BASE_SERVER_ID"] = 23400 + add_openbridge_system(config, "OBP-HDSTACK-AGGREGATOR", network_id=23400) + + with DeterministicScenario(config=config) as scenario: + internal = scenario.systems["OBP-HDSTACK-AGGREGATOR"] + for index, (hops, source_server, source_rptr) in enumerate(( + (b"\x02", 23401, 4000001), + (b"\x03", 23500, 4000001), + (b"\x03", 23401, 0), + )): + subscriber = 1234567 + index + packet = PacketSpec( + peer_id=23400, + rf_src=subscriber, + dst_id=900999, + stream_id=0x01020400 + index, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + ) + internal.dmrd_received( + *packet.decoded_obp_args( + hops=hops, + source_server=source_server, + source_rptr=source_rptr, + ) + ) + self.assertNotIn(bytes_3(subscriber), scenario.bm.SUB_MAP) + def test_hbp_unit_data_to_hbp_reports_actual_target_slot(self): config = minimal_config(("MASTER-A", "MASTER-B")) config["REPORTS"]["REPORT"] = True diff --git a/tests/test_hdstack.py b/tests/test_hdstack.py index 9ca7a65..d331b84 100644 --- a/tests/test_hdstack.py +++ b/tests/test_hdstack.py @@ -45,6 +45,16 @@ class HDStackRuntimeTests(unittest.TestCase): [worker.getint('GLOBAL', 'SERVER_ID') for worker in workers], [23401, 23402], ) + self.assertEqual(aggregator.getint('GLOBAL', 'HDSTACK_INSTANCES'), 2) + self.assertEqual( + aggregator.getint('GLOBAL', 'HDSTACK_BASE_SERVER_ID'), + 23400, + ) + self.assertTrue(all( + worker.getint('GLOBAL', 'HDSTACK_INSTANCES') == 2 + and worker.getint('GLOBAL', 'HDSTACK_BASE_SERVER_ID') == 23400 + for worker in workers + )) self.assertEqual( [worker.getint('SYSTEM', 'PORT') for worker in workers], [54000, 54100], @@ -112,6 +122,63 @@ class HDStackRuntimeTests(unittest.TestCase): self.assertFalse(worker.getboolean('GLOBAL', 'ENABLE_API')) self.assertFalse(worker.getboolean('GLOBAL', 'DATA_GATEWAY')) + def test_data_gateway_is_connected_directly_to_each_worker(self): + self.main_config.write_text( + MAIN_CONFIG + DATA_GATEWAY_CONFIG, + encoding='utf-8', + ) + layout = self.generate() + aggregator = _read_config(layout.aggregator_config) + workers = [_read_config(path) for path in layout.worker_configs] + + self.assertFalse(aggregator.getboolean('GLOBAL', 'DATA_GATEWAY')) + self.assertNotIn('DATA-GATEWAY', aggregator.sections()) + for number, worker in enumerate(workers, 1): + self.assertTrue(worker.getboolean('GLOBAL', 'DATA_GATEWAY')) + self.assertIn('DATA-GATEWAY', worker.sections()) + self.assertEqual( + worker.getint('DATA-GATEWAY', 'PORT'), + 62040 + number, + ) + self.assertEqual( + worker.getint('DATA-GATEWAY', 'TARGET_PORT'), + 62030 + number, + ) + self.assertEqual( + worker.get('DATA-GATEWAY', 'TARGET_IP'), + 'data-gateway', + ) + self.assertEqual(worker.getint('DATA-GATEWAY', 'PROTO_VER'), 5) + self.assertTrue(worker.getboolean('DATA-GATEWAY', 'ENHANCED_OBP')) + self.assertNotIn('EXTERNAL-FBP', worker.sections()) + + def test_data_gateway_requires_fbp5_fbcp_and_valid_derived_ports(self): + invalid_profiles = ( + DATA_GATEWAY_CONFIG.replace("MODE: OPENBRIDGE", "MODE: MASTER"), + DATA_GATEWAY_CONFIG.replace("PROTO_VER: 5", "PROTO_VER: 1"), + DATA_GATEWAY_CONFIG.replace("ENHANCED_OBP: True", "ENHANCED_OBP: False"), + ) + for data_gateway_config in invalid_profiles: + with self.subTest(config=data_gateway_config): + self.main_config.write_text( + MAIN_CONFIG + data_gateway_config, + encoding="utf-8", + ) + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate() + + for data_gateway_config in ( + DATA_GATEWAY_CONFIG.replace("PORT: 62041", "PORT: 54000"), + DATA_GATEWAY_CONFIG.replace("TARGET_PORT: 62031", "TARGET_PORT: 65535"), + ): + with self.subTest(config=data_gateway_config): + self.main_config.write_text( + MAIN_CONFIG + data_gateway_config, + encoding="utf-8", + ) + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate() + def test_internal_fbp_links_are_reciprocal_and_use_expected_ids(self): layout = self.generate() aggregator = _read_config(layout.aggregator_config) @@ -430,6 +497,12 @@ class HDStackContainerTests(unittest.TestCase): ' tmpfs: /tmp\n read_only: "true"', compose, ) + self.assertIn("62041-62045:62041-62045/udp", compose) + self.assertIn("container_name: freedmr-data-gateway", compose) + self.assertIn("freedmr-data-gateway:development-latest", compose) + self.assertIn("FREEDMR_GATEWAY_INGRESS_MODE: fbp", compose) + self.assertIn("FREEDMR_GATEWAY_FBP_RELATIONSHIPS:", compose) + self.assertIn("/etc/freedmr/json/:/data/:ro", compose) def test_installer_records_guided_hdstack_and_bridge_choices(self): installer = ( @@ -447,8 +520,32 @@ class HDStackContainerTests(unittest.TestCase): self.assertIn('- HDSTACK=$HDSTACK_COUNT', installer) self.assertIn('- HDSTACK_BASEID=$HDSTACK_BASEID', installer) self.assertIn('- HDSTACK_BRIDGE=$HDSTACK_BRIDGE_VALUE', installer) + self.assertIn("prompt_value 'APRS-IS login callsign'", installer) + self.assertIn("prompt_value 'APRS-IS passcode'", installer) + self.assertIn('build_data_gateway_relationships', installer) + self.assertIn('\\"bind_port\\":${bind_port}', installer) + self.assertIn('\\"peer_port\\":${peer_port}', installer) + self.assertIn('DATA_GATEWAY: True', installer) + self.assertIn('ENABLED: True', installer) + self.assertIn('TARGET_IP: freedmr-data-gateway', installer) + self.assertIn('s/^SERVER_ID:.*/SERVER_ID: $HDSTACK_BASEID/', installer) + self.assertIn('s/^SERVER_ID:.*/SERVER_ID: $HDSTACK_BRIDGE_ID/', installer) self.assertNotIn('/etc/freedmr/.env', installer) + def test_shipped_config_keeps_native_gateway_disabled_until_install(self): + for path in ( + ROOT / 'FreeDMR.cfg', + ROOT / 'freedmr.cfg', + DOCKER_CONFIGS / 'freedmr.cfg', + ): + with self.subTest(path=path): + config = _read_config(path) + self.assertFalse(config.getboolean('GLOBAL', 'DATA_GATEWAY')) + self.assertFalse(config.getboolean('DATA-GATEWAY', 'ENABLED')) + self.assertEqual(config.get('DATA-GATEWAY', 'MODE'), 'OPENBRIDGE') + self.assertEqual(config.getint('DATA-GATEWAY', 'PROTO_VER'), 5) + self.assertTrue(config.getboolean('DATA-GATEWAY', 'ENHANCED_OBP')) + def test_bridge_templates_are_valid_and_idle(self): bridge = _read_config(DOCKER_CONFIGS / 'freedmr-bridge.cfg') rules_text = ( @@ -593,6 +690,25 @@ ANNOUNCEMENT_LANGUAGE: en_GB """.lstrip() +DATA_GATEWAY_CONFIG = """ +[DATA-GATEWAY] +MODE: OPENBRIDGE +ENABLED: True +IP: 0.0.0.0 +PORT: 62041 +NETWORK_ID: 0 +PASSPHRASE: gateway-test +TARGET_IP: data-gateway +TARGET_PORT: 62031 +USE_ACL: False +SUB_ACL: PERMIT:ALL +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 +""".lstrip() + + BRIDGE_CONFIG = """ [GLOBAL] SERVER_ID: 23500 diff --git a/tests/test_udp_blackbox_harness.py b/tests/test_udp_blackbox_harness.py index 3d4c3ba..055d190 100644 --- a/tests/test_udp_blackbox_harness.py +++ b/tests/test_udp_blackbox_harness.py @@ -806,6 +806,54 @@ class UdpBlackBoxHarnessTest(unittest.TestCase): self.assertEqual(captured.fields["dst_id"], bytes_3(91)) self.assertEqual(captured.fields["slot"], 1) + def test_local_hbp_data_reaches_native_gateway_as_fbp_v5(self): + require_udp_integration_enabled() + + with UdpBlackBoxScenario( + ts2_static="", + data_gateway=True, + fbp_systems={"DATA-GATEWAY": 0}, + ) as scenario: + master_a = scenario.repeater("MASTER-A", 1001) + gateway = scenario.fbp_peer("DATA-GATEWAY") + payload = bytes(range(33)) + try: + master_a.login() + gateway.send_bcka() + gateway.send_bcve() + gateway.drain(seconds=0.2) + + master_a.send_dmr( + PacketSpec( + peer_id=1001, + rf_src=3120001, + dst_id=900999, + slot=2, + stream_id=0x01020324, + call_type="unit", + frame_type=HBPF_DATA_SYNC, + dtype_vseq=6, + payload=payload, + ber=bytes((2,)), + rssi=bytes((3,)), + ) + ) + captured = gateway.recv_dmre(timeout=2.0) + finally: + master_a.close() + + self.assertEqual(captured.packet[:4], b"DMRE") + self.assertEqual(captured.fields["fbp_version"], 5) + self.assertEqual(captured.fields["peer_id"], bytes_4(9990)) + self.assertEqual(captured.fields["rf_src"], bytes_3(3120001)) + self.assertEqual(captured.fields["dst_id"], bytes_3(900999)) + self.assertEqual(captured.fields["source_server"], bytes_4(9990)) + self.assertEqual(captured.fields["source_rptr"], bytes_4(1001)) + self.assertEqual(captured.fields["slot"], 1) + self.assertEqual(captured.fields["dmr_payload"], payload) + self.assertEqual(captured.fields["ber"], bytes((2,))) + self.assertEqual(captured.fields["rssi"], bytes((3,))) + def test_fbp_enhanced_keepalive_gates_hbp_to_fbp_forwarding(self): require_udp_integration_enabled() @@ -1523,7 +1571,7 @@ class UdpBlackBoxHarnessTest(unittest.TestCase): self.assertEqual(quench.packet[4:7], bytes_3(91)) self.assertEqual(quench.packet[7:11], bytes_4(stream_id)) - def test_fbp_max_hops_is_source_quenched_without_hbp_leak(self): + def test_fbp_eleven_hops_is_accepted(self): require_udp_integration_enabled() stream_id = 0x01020310 @@ -1547,6 +1595,37 @@ class UdpBlackBoxHarnessTest(unittest.TestCase): source_rptr=1001, hops=10, ) + captured = master_b.recv(timeout=2.0) + finally: + master_b.close() + + self.assertEqual(captured.fields["dst_id"], bytes_3(91)) + self.assertEqual(captured.fields["stream_id"], bytes_4(stream_id)) + + def test_fbp_above_max_hops_is_source_quenched_without_hbp_leak(self): + require_udp_integration_enabled() + + stream_id = 0x01020310 + with UdpBlackBoxScenario(fbp_systems={"OBP-1": 3001}) as scenario: + master_b = scenario.repeater("MASTER-B", 1002) + fbp_peer = scenario.fbp_peer("OBP-1") + try: + master_b.login() + fbp_peer.send_bcka() + fbp_peer.send_bcve() + fbp_peer.drain(seconds=0.2) + + fbp_peer.send_fbp( + PacketSpec( + peer_id=3001, + dst_id=91, + slot=1, + stream_id=stream_id, + ), + source_server=9991, + source_rptr=1001, + hops=11, + ) quench = fbp_peer.recv_opcode(b"BCSQ", timeout=2.0) with self.assertRaises(socket.timeout): master_b.recv(timeout=0.5)