diff --git a/README.md b/README.md index 3639cfc..6da7427 100755 --- a/README.md +++ b/README.md @@ -6,3 +6,6 @@ Please see the wiki for documentation. For the single-client SQL importer and dashboard reporting relay, see [docs/reporting-relay.md](docs/reporting-relay.md). + +For the optional two-to-five-worker container deployment and reporting MUX, see +[docs/hdstack.md](docs/hdstack.md). diff --git a/docker-configs/docker-compose.yml b/docker-configs/docker-compose.yml index b2fe633..ce0fe2d 100644 --- a/docker-configs/docker-compose.yml +++ b/docker-configs/docker-compose.yml @@ -24,6 +24,9 @@ services: mem_reservation: 600m volumes: - '/etc/freedmr/freedmr.cfg:/opt/freedmr/freedmr.cfg' + # HDSTACK=2..5 requires the separate bridge inputs below. + #- '/etc/freedmr/freedmr-bridge.cfg:/opt/freedmr/freedmr-bridge.cfg:ro' + #- '/etc/freedmr/rules-bridge.py:/opt/freedmr/rules.py:ro' #Write JSON files outside of container - '/etc/freedmr/json/:/opt/freedmr/json/' @@ -48,6 +51,11 @@ services: - FDPROXY_DEBUG=0 #Override proxy external port #- FDPROXY_LISTENPORT=62031 + # Optional 2..5 worker HDStack deployment. Omit for current mode. + #- HDSTACK=2 + #- HDSTACK_BASEID=23400 + #- HDSTACK_BRIDGE_CONFIG=/opt/freedmr/freedmr-bridge.cfg + #- HDSTACK_BRIDGE_RULES=/opt/freedmr/rules.py read_only: "true" freedmrmonitor2: diff --git a/docker-configs/entrypoint-proxy b/docker-configs/entrypoint-proxy index 02d77a2..9ef1d4a 100755 --- a/docker-configs/entrypoint-proxy +++ b/docker-configs/entrypoint-proxy @@ -19,14 +19,48 @@ cd /opt/freedmr +HDSTACK_COUNT="${HDSTACK:-1}" +HDSTACK_OUTPUT_DIR="${HDSTACK_OUTPUT_DIR:-/dev/shm/freedmr-hdstack}" +HDSTACK_BRIDGE_CONFIG="${HDSTACK_BRIDGE_CONFIG:-/opt/freedmr/freedmr-bridge.cfg}" +HDSTACK_BRIDGE_RULES="${HDSTACK_BRIDGE_RULES:-/opt/freedmr/rules.py}" + if [ "$BRIDGE_SERVER" == 1 ] then echo 'Starting in Bridge mode...' exec python /opt/freedmr/bridge.py -c freedmr.cfg -r rules.py + +elif [ "$HDSTACK_COUNT" = 1 ] +then + echo 'Starting in FreeDMR mode...' + exec /usr/bin/supervisord -c /etc/supervisor/conf.d/supervisord.conf + #python /opt/freedmr/hotspot_proxy_v2.py & + #python /opt/freedmr/playback.py -c loro.cfg & + #exec python /opt/freedmr/bridge_master.py -c freedmr.cfg + +elif [ "$HDSTACK_COUNT" -ge 2 ] 2>/dev/null && [ "$HDSTACK_COUNT" -le 5 ] 2>/dev/null +then + if [ -z "$HDSTACK_BASEID" ] + then + echo 'HDSTACK_BASEID is required when HDSTACK is 2..5' >&2 + exit 2 + fi + + echo "Starting HDStack with $HDSTACK_COUNT workers..." + python /opt/freedmr/hdstack_runtime.py \ + --main-config /opt/freedmr/freedmr.cfg \ + --bridge-config "$HDSTACK_BRIDGE_CONFIG" \ + --bridge-rules "$HDSTACK_BRIDGE_RULES" \ + --output-dir "$HDSTACK_OUTPUT_DIR" \ + --instances "$HDSTACK_COUNT" \ + --base-id "$HDSTACK_BASEID" + status=$? + if [ "$status" -ne 0 ] + then + exit "$status" + fi + exec /usr/bin/supervisord -c "$HDSTACK_OUTPUT_DIR/supervisord.conf" + else - echo 'Starting in FreeDMR mode...' - exec /usr/bin/supervisord -c /etc/supervisor/conf.d/supervisord.conf - #python /opt/freedmr/hotspot_proxy_v2.py & - #python /opt/freedmr/playback.py -c loro.cfg & - #exec python /opt/freedmr/bridge_master.py -c freedmr.cfg + echo 'HDSTACK must be omitted or set to an integer from 1 to 5' >&2 + exit 2 fi diff --git a/docs/hdstack.md b/docs/hdstack.md new file mode 100644 index 0000000..fc04ed9 --- /dev/null +++ b/docs/hdstack.md @@ -0,0 +1,183 @@ +# HDStack multi-process deployment + +HDStack raises the aggregate hotspot capacity of one FreeDMR container by +running two to five independent FreeDMR HBP workers. It does not raise the +connection limit of an individual worker. + +Omitting `HDSTACK`, or setting `HDSTACK=1`, preserves the existing standard +single-FreeDMR process layout. Setting `HDSTACK` to a value from 2 through 5 +selects the generated multi-process layout: + +```text +hotspots -> proxy -> HBP workers -> no-HBP aggregator -> external FBP/OBP + | + +-> separate bridge.py instance + +worker reports -> reporting MUX -> existing dashboard/reporting consumer +``` + +The aggregator owns all OpenBridge/FBP links from the main `freedmr.cfg` and +has no HBP `MASTER`, `PEER` or `XLXPEER` system. The separate `bridge.py` +process retains its own configuration and rules. The reporting MUX reads only +worker reports; aggregator and bridge reporting are not included. + +## Container configuration + +Add these environment values to the standard `freedmr` service: + +```yaml +environment: + - HDSTACK=2 + - HDSTACK_BASEID=23400 +``` + +`HDSTACK_BASEID` is the aggregator server ID. Worker IDs are allocated by +adding their one-based worker number. For the example above, the aggregator is +`23400` and the two workers are `23401` and `23402`. With `HDSTACK=5`, worker +IDs continue through `23405`. + +Multi-process mode also requires the existing separate bridge inputs: + +```yaml +volumes: + - '/etc/freedmr/freedmr.cfg:/opt/freedmr/freedmr.cfg' + - '/etc/freedmr/freedmr-bridge.cfg:/opt/freedmr/freedmr-bridge.cfg:ro' + - '/etc/freedmr/rules-bridge.py:/opt/freedmr/rules.py:ro' +``` + +The paths can be changed with: + +```yaml +environment: + - HDSTACK_BRIDGE_CONFIG=/opt/freedmr/freedmr-bridge.cfg + - HDSTACK_BRIDGE_RULES=/opt/freedmr/rules.py +``` + +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 and the configured link to `bridge.py`. +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. + +The bridge configuration must contain the reciprocal link to the aggregator +and any `PEER` or `XLXPEER` systems used by `bridge.py`. Its `rules.py` remains +the authority for bridge routing. Bridge reporting must either be disabled or +use a port which does not collide with the worker reports or MUX. + +## Worker ports + +The `[SYSTEM]` `PORT` and `GENERATOR` values define each worker range. With the +standard values `PORT=54000` and `GENERATOR=100`: + +```text +worker 1 -> 54000..54099 +worker 2 -> 54100..54199 +worker 3 -> 54200..54299 +worker 4 -> 54300..54399 +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. + +Workers connect to the aggregator over generated loopback OBP v1 links. Worker +ports are `7001` through `7005`; reciprocal aggregator ports are `7101` through +`7105`. These links are private to the container. + +## Session assignment + +The proxy rotates new DMR IDs through the configured worker ranges. With two +workers the assignments are: + +```text +new session 1 -> worker 1 +new session 2 -> worker 2 +new session 3 -> worker 1 +new session 4 -> worker 2 +``` + +Within the selected worker range, the proxy retains its existing random choice +of an available generated system port. + +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 +BRIDGES and OPTIONS state. The proxy does not rebalance or fail over an active +session. Once the session expires, a later login by the same DMR ID is a new +assignment. Affinity is not persisted across proxy restarts. + +Existing explicit proxy configurations using `DestportStart` and `DestportEnd` +continue to describe one range. The equivalent direct multi-range setting is: + +```ini +[PROXY] +BackendPortRanges: [[54000, 54099], [54100, 54199]] +``` + +or: + +```sh +FDPROXY_BACKEND_PORT_RANGES='[[54000,54099],[54100,54199]]' +``` + +The standard HDStack entrypoint supplies this setting automatically. + +## Proxy diagnostics + +In multi-worker mode the proxy logs each backend number and port range during +startup. With the standard two-worker layout this includes: + +```text +(PROXY)(HDSTACK) Backend:1 ports:54000-54099. +(PROXY)(HDSTACK) Backend:2 ports:54100-54199. +``` + +When `FDPROXY_CLIENTINFO=1`, as in the shipped Compose configuration, client +assignment and removal messages include both the backend number and exact +generated port. This makes round-robin assignment and retained session +affinity visible without enabling packet-level debug logging. Single-backend +startup output remains unchanged. + +## Reporting MUX + +When reporting is enabled in the main configuration, worker report ports are +derived from its `REPORT_PORT`. With the default output port `4321`: + +```text +MUX output -> 4321 +worker 1 -> 4322 +worker 2 -> 4323 +worker 3 -> 4324 +worker 4 -> 4325 +worker 5 -> 4326 +``` + +The MUX preserves the existing netstring reporting protocol and the main +configuration's `REPORT_CLIENTS` allow-list. It merges worker configuration +and bridge snapshots and relays worker bridge events. System names are +prefixed with the stable worker number to prevent collisions: + +```text +1:SYSTEM-001 +2:SYSTEM-001 +``` + +Only the system-name field is changed. The reporting protocol contains trusted +Python pickle data and must remain on a trusted network. + +A failed or malformed worker reporting source is isolated from healthy +sources. Reporting failure does not affect the proxy or DMR traffic. If main +reporting is disabled, worker reporting and the MUX are both disabled. + +## Limitations + +HDStack deliberately provides no least-connections balancing, active backend +failover, live migration, persistent affinity or automatic reconfiguration. +A failed worker does not cause its active sessions to move to another worker. +The separate bridge config and rules remain sysop-managed; the runtime +generator does not create routing policy. diff --git a/docs/reporting-relay.md b/docs/reporting-relay.md index 94ede4b..aa47bd8 100644 --- a/docs/reporting-relay.md +++ b/docs/reporting-relay.md @@ -30,3 +30,21 @@ END: type,event,trx,system,streamid,peerid,subid,slot,dstid,duration,source_se Legacy events remain valid and are stored with null origin columns. The relay is optional; omitting `--relay-port` preserves the previous importer behavior. + +## Multiple FreeDMR backends + +`report_mux.py` combines reporting from multiple FreeDMR backend processes +into one compatible reporting socket. It is used by the HDStack container +profile and can also be run directly: + +```sh +python report_mux.py --listen-interface 127.0.0.1 --listen-port 4321 \ + --client 127.0.0.1 \ + --backend 1,127.0.0.1,4322 --backend 2,127.0.0.1,4323 +``` + +Backend numbers are explicit and must be unique. Configuration keys, bridge +member system fields and bridge-event system fields are prefixed with that +number, for example `1:SYSTEM-001`. Other fields are preserved. The MUX +reconnects to each backend independently and a reporting-source failure does +not affect the radio proxy. See [hdstack.md](hdstack.md) for the generated container topology and configuration. diff --git a/hdstack_runtime.py b/hdstack_runtime.py new file mode 100644 index 0000000..121cecc --- /dev/null +++ b/hdstack_runtime.py @@ -0,0 +1,563 @@ +############################################################################### +# Copyright (C) 2026 Simon Adlem, G7RZU +# +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +############################################################################### + +"""Generate private FreeDMR HDStack process configuration.""" + +import argparse +import configparser +import json +import shlex +from dataclasses import dataclass +from pathlib import Path + + +MIN_INSTANCES = 2 +MAX_INSTANCES = 5 +MAX_SERVER_ID = 0xFFFFFFFF +SUPPORT_SECTIONS = ('GLOBAL', 'REPORTS', 'LOGGER', 'ALIASES', 'ALLSTAR') +INTERNAL_WORKER_PORT_BASE = 7000 +INTERNAL_AGGREGATOR_PORT_BASE = 7100 + + +class HDStackConfigurationError(ValueError): + """Raised when an HDStack deployment cannot be generated safely.""" + + +@dataclass(frozen=True) +class HDStackLayout: + """Paths and derived addresses for one generated deployment.""" + + instances: int + base_server_id: int + aggregator_config: Path + worker_configs: tuple + supervisor_config: Path + backend_ranges: tuple + worker_report_ports: tuple + + +def generate_hdstack( + main_config, + bridge_config, + bridge_rules, + output_dir, + instances, + base_server_id, +): + """Generate configs for an aggregator, workers, bridge, proxy and MUX.""" + instances = _validated_instances(instances) + base_server_id = _validated_base_id(base_server_id, instances) + main_path = _required_file(main_config, 'main FreeDMR configuration') + bridge_path = _required_file(bridge_config, 'bridge configuration') + rules_path = _required_file(bridge_rules, 'bridge rules') + output_path = Path(output_dir) + + source = _read_config(main_path) + bridge_source = _read_config(bridge_path) + system = _system_template(source) + hbp_port = _integer(system, 'PORT', 54000, 'SYSTEM PORT') + generator_size = _integer(system, 'GENERATOR', 100, 'SYSTEM GENERATOR') + if generator_size < 1: + raise HDStackConfigurationError('SYSTEM GENERATOR must be positive') + + backend_ranges = tuple( + ( + hbp_port + (number - 1) * generator_size, + hbp_port + number * generator_size - 1, + ) + for number in range(1, instances + 1) + ) + for start, end in backend_ranges: + _validate_port(start, 'generated HBP range start') + _validate_port(end, 'generated HBP range end') + + report_enabled = _boolean(source, 'REPORTS', 'REPORT', True) + report_port = _section_integer(source, 'REPORTS', 'REPORT_PORT', 4321) + _validate_port(report_port, 'report MUX port') + worker_report_ports = tuple( + report_port + number for number in range(1, instances + 1) + ) + if report_enabled: + for port in worker_report_ports: + _validate_port(port, 'worker report port') + + _validate_generated_port_collisions( + source, + bridge_source, + backend_ranges, + report_enabled, + report_port, + worker_report_ports, + instances, + ) + + aggregator = _support_config(source) + aggregator['GLOBAL']['SERVER_ID'] = str(base_server_id) + aggregator['GLOBAL']['GEN_STAT_BRIDGES'] = 'True' + aggregator['REPORTS']['REPORT'] = 'False' + for section in source.sections(): + if _mode(source, section) == 'OPENBRIDGE': + _copy_section(source, aggregator, section) + for number in range(1, instances + 1): + name = 'OBP-HDSTACK-WORKER-{}'.format(number) + if source.has_section(name): + raise HDStackConfigurationError( + 'main configuration already contains reserved section {}'.format(name) + ) + _add_internal_openbridge( + aggregator, + name, + INTERNAL_AGGREGATOR_PORT_BASE + number, + INTERNAL_WORKER_PORT_BASE + number, + base_server_id + number, + ) + + workers = [] + for number, (start, _end) in enumerate(backend_ranges, 1): + worker = _support_config(source) + worker['GLOBAL']['SERVER_ID'] = str(base_server_id + number) + worker['GLOBAL']['GEN_STAT_BRIDGES'] = 'False' + worker['GLOBAL']['ENABLE_API'] = 'False' + worker['GLOBAL']['DATA_GATEWAY'] = 'False' + worker['ALIASES']['TRY_DOWNLOAD'] = 'False' + _number_mutable_alias_file(worker, 'SUB_MAP_FILE', number) + _number_mutable_alias_file(worker, 'KEYS_FILE', number) + worker['REPORTS']['REPORT'] = str(report_enabled) + worker['REPORTS']['REPORT_PORT'] = str(worker_report_ports[number - 1]) + worker['REPORTS']['REPORT_CLIENTS'] = '127.0.0.1' + _copy_section(source, worker, 'SYSTEM') + worker['SYSTEM']['PORT'] = str(start) + name = 'OBP-HDSTACK-AGGREGATOR' + if source.has_section(name): + raise HDStackConfigurationError( + 'main configuration already contains reserved section {}'.format(name) + ) + _add_internal_openbridge( + worker, + name, + INTERNAL_WORKER_PORT_BASE + number, + INTERNAL_AGGREGATOR_PORT_BASE + number, + base_server_id, + ) + workers.append(worker) + + output_path.mkdir(parents=True, exist_ok=True) + aggregator_path = output_path / 'aggregator.cfg' + worker_paths = tuple( + output_path / 'worker-{}.cfg'.format(number) + for number in range(1, instances + 1) + ) + supervisor_path = output_path / 'supervisord.conf' + + _write_config(aggregator, aggregator_path) + for worker, path in zip(workers, worker_paths): + _write_config(worker, path) + supervisor_path.write_text( + _supervisor_config( + aggregator_path, + worker_paths, + bridge_path, + rules_path, + backend_ranges, + report_enabled, + report_port, + worker_report_ports, + _report_clients(source), + ), + encoding='utf-8', + ) + + return HDStackLayout( + instances=instances, + base_server_id=base_server_id, + aggregator_config=aggregator_path, + worker_configs=worker_paths, + supervisor_config=supervisor_path, + backend_ranges=backend_ranges, + worker_report_ports=worker_report_ports, + ) + + +def _validated_instances(value): + if isinstance(value, bool) or not isinstance(value, int): + raise HDStackConfigurationError('HDSTACK must be an integer') + if value < MIN_INSTANCES or value > MAX_INSTANCES: + raise HDStackConfigurationError( + 'generated HDStack requires HDSTACK in {}..{}'.format( + MIN_INSTANCES, + MAX_INSTANCES, + ) + ) + return value + + +def _validated_base_id(value, instances): + if isinstance(value, bool) or not isinstance(value, int): + raise HDStackConfigurationError('HDSTACK_BASEID must be an integer') + if value < 1 or value + instances > MAX_SERVER_ID: + raise HDStackConfigurationError( + 'HDSTACK_BASEID and worker IDs must fit an unsigned 32-bit ID' + ) + return value + + +def _required_file(value, description): + path = Path(value) + if not path.is_file(): + raise HDStackConfigurationError('{} not found: {}'.format(description, path)) + return path + + +def _read_config(path): + parser = configparser.ConfigParser(interpolation=None) + try: + loaded = parser.read(path) + except configparser.Error as err: + raise HDStackConfigurationError( + 'invalid configuration {}: {}'.format(path, err) + ) from err + if not loaded: + raise HDStackConfigurationError('could not read configuration {}'.format(path)) + return parser + + +def _system_template(source): + if not source.has_section('SYSTEM'): + raise HDStackConfigurationError('main configuration requires [SYSTEM]') + if _mode(source, 'SYSTEM') != 'MASTER': + raise HDStackConfigurationError('[SYSTEM] must use MODE MASTER') + if not _boolean(source, 'SYSTEM', 'ENABLED', True): + raise HDStackConfigurationError('[SYSTEM] must be enabled') + return source['SYSTEM'] + + +def _support_config(source): + result = configparser.ConfigParser(interpolation=None) + for section in SUPPORT_SECTIONS: + result.add_section(section) + if source.has_section(section): + for key, value in source.items(section, raw=True): + result[section][key] = value + return result + + +def _copy_section(source, destination, section): + if destination.has_section(section): + destination.remove_section(section) + destination.add_section(section) + for key, value in source.items(section, raw=True): + destination[section][key] = value + + +def _mode(source, section): + return source.get(section, 'MODE', fallback='').strip().upper() + + +def _boolean(source, section, option, fallback): + if not source.has_section(section): + return fallback + try: + return source.getboolean(section, option, fallback=fallback) + except ValueError as err: + raise HDStackConfigurationError( + '{} {} must be boolean'.format(section, option) + ) from err + + +def _integer(section, option, fallback, description): + try: + return section.getint(option, fallback=fallback) + except ValueError as err: + raise HDStackConfigurationError('{} must be an integer'.format(description)) from err + + +def _section_integer(source, section, option, fallback): + if not source.has_section(section): + return fallback + try: + return source.getint(section, option, fallback=fallback) + except ValueError as err: + raise HDStackConfigurationError( + '{} {} must be an integer'.format(section, option) + ) from err + + +def _validate_port(port, description): + if port < 1 or port > 65535: + raise HDStackConfigurationError('{} must be within 1..65535'.format(description)) + + +def _enabled_bind_ports(source, excluded_sections=()): + excluded = set(excluded_sections) + ports = {} + for section in source.sections(): + if section in excluded or not _mode(source, section): + continue + if not _boolean(source, section, 'ENABLED', True): + continue + if source.has_option(section, 'PORT'): + port = _section_integer(source, section, 'PORT', 0) + _validate_port(port, '{} PORT'.format(section)) + ports.setdefault(port, []).append(section) + return ports + + +def _validate_generated_port_collisions( + source, + bridge_source, + backend_ranges, + report_enabled, + report_port, + worker_report_ports, + instances, +): + generated = {} + + def reserve(port, owner): + _validate_port(port, owner) + if port in generated: + raise HDStackConfigurationError( + 'generated port {} is used by both {} and {}'.format( + port, + generated[port], + owner, + ) + ) + generated[port] = owner + + for number, (start, end) in enumerate(backend_ranges, 1): + for port in range(start, end + 1): + reserve(port, 'worker {} HBP'.format(number)) + for number in range(1, instances + 1): + reserve(INTERNAL_WORKER_PORT_BASE + number, 'worker {} OBP'.format(number)) + reserve( + INTERNAL_AGGREGATOR_PORT_BASE + number, + 'aggregator worker {} OBP'.format(number), + ) + if report_enabled: + reserve(report_port, 'report MUX') + for number, port in enumerate(worker_report_ports, 1): + reserve(port, 'worker {} reporting'.format(number)) + + configured = _enabled_bind_ports(source, excluded_sections=('SYSTEM', 'ECHO')) + bridge_ports = _enabled_bind_ports(bridge_source) + for port, sections in configured.items(): + if port in generated: + raise HDStackConfigurationError( + 'main section {} conflicts with {} on port {}'.format( + ', '.join(sections), + generated[port], + port, + ) + ) + for port, sections in bridge_ports.items(): + if port in generated: + raise HDStackConfigurationError( + 'bridge section {} conflicts with {} on port {}'.format( + ', '.join(sections), + generated[port], + port, + ) + ) + + if _boolean(bridge_source, 'REPORTS', 'REPORT', True): + bridge_report_port = _section_integer( + bridge_source, + 'REPORTS', + 'REPORT_PORT', + 4321, + ) + _validate_port(bridge_report_port, 'bridge report port') + if bridge_report_port in generated: + raise HDStackConfigurationError( + 'bridge reporting conflicts with {} on port {}'.format( + generated[bridge_report_port], + bridge_report_port, + ) + ) + + +def _add_internal_openbridge( + config, + name, + local_port, + target_port, + expected_network_id, +): + config.add_section(name) + config[name].update({ + 'MODE': 'OPENBRIDGE', + 'ENABLED': 'True', + 'IP': '127.0.0.1', + 'PORT': str(local_port), + 'NETWORK_ID': str(expected_network_id), + 'PASSPHRASE': 'internal', + 'TARGET_IP': '127.0.0.1', + 'TARGET_PORT': str(target_port), + 'USE_ACL': 'False', + 'SUB_ACL': 'PERMIT:ALL', + 'TGID_ACL': 'PERMIT:ALL', + 'RELAX_CHECKS': 'False', + 'ENHANCED_OBP': 'False', + 'PROTO_VER': '1', + }) + + +def _number_mutable_alias_file(config, option, number): + value = config['ALIASES'].get(option, '').strip() + if not value: + return + path = Path(value) + suffix = ''.join(path.suffixes) + if suffix: + stem = path.name[:-len(suffix)] + else: + stem = path.name + name = '{}-hdstack-{}{}'.format(stem, number, suffix) + config['ALIASES'][option] = str(path.with_name(name)) + + +def _report_clients(source): + if not source.has_section('REPORTS'): + return ('*',) + value = source.get('REPORTS', 'REPORT_CLIENTS', fallback='*') + clients = tuple(item.strip() for item in value.split(',') if item.strip()) + return clients or ('*',) + + +def _write_config(config, path): + with path.open('w', encoding='utf-8') as config_file: + config.write(config_file) + + +def _program(name, command, priority, environment=None): + lines = [ + '[program:{}]'.format(name), + 'directory=/opt/freedmr', + 'stdout_logfile=/dev/fd/1', + 'stdout_logfile_maxbytes=0', + 'redirect_stderr=true', + 'command={}'.format(command), + 'stopwaitsecs=30', + 'autorestart=true', + 'priority={}'.format(priority), + ] + if environment: + lines.append('environment={}'.format(environment)) + return '\n'.join(lines) + + +def _supervisor_config( + aggregator_path, + worker_paths, + bridge_path, + rules_path, + backend_ranges, + report_enabled, + report_port, + worker_report_ports, + report_clients, +): + parts = [ + '[supervisord]\n' + 'nodaemon=true\n' + 'logfile=/dev/null\n' + 'logfile_maxbytes=0\n' + 'pidfile={}'.format(aggregator_path.parent / 'supervisord.pid'), + _program( + 'hdstack-aggregator', + 'python /opt/freedmr/bridge_master.py -c {}'.format( + shlex.quote(str(aggregator_path)) + ), + 10, + ), + ] + for number, path in enumerate(worker_paths, 1): + parts.append( + _program( + 'hdstack-worker-{}'.format(number), + 'python /opt/freedmr/bridge_master.py -c {}'.format( + shlex.quote(str(path)) + ), + 20, + ) + ) + parts.extend([ + _program( + 'bridge', + 'python /opt/freedmr/bridge.py -c {} -r {}'.format( + shlex.quote(str(bridge_path)), + shlex.quote(str(rules_path)), + ), + 30, + ), + _program( + 'playback', + '/opt/freedmr/playback.py -c /opt/freedmr/loro.cfg', + 30, + ), + _program( + 'proxy', + 'python /opt/freedmr/hotspot_proxy_v2.py', + 40, + 'FDPROXY_BACKEND_PORT_RANGES="{}"'.format( + json.dumps(backend_ranges, separators=(',', ':')) + ), + ), + ]) + if report_enabled: + command = [ + 'python', + '/opt/freedmr/report_mux.py', + '--listen-interface', + '0.0.0.0', + '--listen-port', + str(report_port), + ] + for client in report_clients: + command.extend(('--client', client)) + for number, port in enumerate(worker_report_ports, 1): + command.extend(('--backend', '{},127.0.0.1,{}'.format(number, port))) + parts.append( + _program( + 'report-mux', + ' '.join(shlex.quote(item) for item in command), + 50, + ) + ) + return '\n\n'.join(parts) + '\n' + + +def main(argv=None): + parser = argparse.ArgumentParser( + description='Generate an HDStack FreeDMR container configuration', + ) + parser.add_argument('--main-config', required=True) + parser.add_argument('--bridge-config', required=True) + parser.add_argument('--bridge-rules', required=True) + parser.add_argument('--output-dir', required=True) + parser.add_argument('--instances', required=True, type=int) + parser.add_argument('--base-id', required=True, type=int) + args = parser.parse_args(argv) + try: + generate_hdstack( + args.main_config, + args.bridge_config, + args.bridge_rules, + args.output_dir, + args.instances, + args.base_id, + ) + except HDStackConfigurationError as err: + parser.error(str(err)) + return 0 + + +if __name__ == '__main__': + raise SystemExit(main()) diff --git a/hotspot_proxy_v2.py b/hotspot_proxy_v2.py index 0fbf1bd..84b964d 100644 --- a/hotspot_proxy_v2.py +++ b/hotspot_proxy_v2.py @@ -20,9 +20,10 @@ from twisted.internet.protocol import DatagramProtocol from twisted.internet import reactor, task from time import time from utils import int_id -import random import ipaddress +import json import os +import random from setproctitle import setproctitle from datetime import datetime import Pyro5.api @@ -43,6 +44,62 @@ DEFAULT_DESTPORT_END = DEFAULT_DESTPORT_START + DEFAULT_DESTPORT_COUNT - 1 def bool_from_env(value): return str(value).strip().lower() in ('1', 'true', 'yes', 'on') +def backend_port_ranges(value, default_start, default_end): + """Return validated inclusive backend port ranges in configured order.""" + if value in (None, ''): + ranges = [[default_start, default_end]] + elif isinstance(value, str): + ranges = json.loads(value) + else: + ranges = value + + if not isinstance(ranges, (list, tuple)) or not ranges: + raise ValueError('BackendPortRanges must be a non-empty JSON list') + + normalised = [] + allocated = set() + for item in ranges: + if not isinstance(item, (list, tuple)) or len(item) != 2: + raise ValueError('each backend port range must contain start and end') + start, end = item + if ( + isinstance(start, bool) + or isinstance(end, bool) + or not isinstance(start, int) + or not isinstance(end, int) + ): + raise ValueError('backend port range values must be integers') + if start < 1 or end > 65535 or start > end: + raise ValueError( + 'backend port range must be within 1..65535 and start at or before end' + ) + ports = set(range(start, end + 1)) + if allocated.intersection(ports): + raise ValueError('backend port ranges must not overlap') + allocated.update(ports) + normalised.append((start, end)) + return tuple(normalised) + +def backend_number_for_port(port, ranges): + """Return the one-based backend number owning a generated port.""" + for number, (start, end) in enumerate(ranges, 1): + if start <= port <= end: + return number + raise ValueError('port {} is outside the configured backend ranges'.format(port)) + +def log_backend_port_ranges(ranges): + """Log the generated backend map when more than one backend is active.""" + if len(ranges) < 2: + return + for number, (start, end) in enumerate(ranges, 1): + print( + '(PROXY)(HDSTACK) Backend:{} ports:{}-{}.'.format( + number, + start, + end, + ) + ) + def IsIPv4Address(ip): try: ipaddress.IPv4Address(ip) @@ -94,7 +151,7 @@ class privHelper(): class Proxy(DatagramProtocol): - def __init__(self,Master,ListenPort,connTrack,peerTrack,blackList,IPBlackList,Timeout,Debug,ClientInfo,DestportStart,DestPortEnd,privHelper,rptlTrack): + def __init__(self,Master,ListenPort,connTrack,peerTrack,blackList,IPBlackList,Timeout,Debug,ClientInfo,DestportStart,DestPortEnd,privHelper,rptlTrack,BackendPortRanges=None): self.master = Master self.ListenPort = ListenPort self.connTrack = connTrack @@ -106,10 +163,32 @@ class Proxy(DatagramProtocol): self.IPBlackList = IPBlackList self.destPortStart = DestportStart self.destPortEnd = DestPortEnd - self.numPorts = DestPortEnd - DestportStart + 1 + self.backendPortRanges = backend_port_ranges( + BackendPortRanges, + DestportStart, + DestPortEnd, + ) + self.numPorts = sum(end - start + 1 for start, end in self.backendPortRanges) + self._nextBackendRange = 0 self.privHelper = privHelper self.rptlTrack = rptlTrack + def allocate_port(self): + """Allocate one free port, rotating new sessions between backends.""" + range_count = len(self.backendPortRanges) + for offset in range(range_count): + range_index = (self._nextBackendRange + offset) % range_count + start, end = self.backendPortRanges[range_index] + available = [ + port + for port in range(start, end + 1) + if port in self.connTrack and not self.connTrack[port] + ] + if available: + self._nextBackendRange = (range_index + 1) % range_count + return random.choice(available) + return None + def packet_too_short(self, data, minimum, command): if len(data) >= minimum: @@ -124,7 +203,14 @@ class Proxy(DatagramProtocol): if self.debug: print("dead",_peer_id) if self.clientinfo and _peer_id != b'\xff\xff\xff\xff': - print(f"{datetime.now().replace(microsecond=0)} Client: ID:{str(int_id(_peer_id)).rjust(9)} IP:{self.peerTrack[_peer_id]['shost'].rjust(15)} Port:{self.peerTrack[_peer_id]['sport']} Removed.") + _backend = '' + if len(self.backendPortRanges) > 1: + _dport = self.peerTrack[_peer_id]['dport'] + _backend = ' from backend:{} port:{}'.format( + backend_number_for_port(_dport, self.backendPortRanges), + _dport, + ) + print(f"{datetime.now().replace(microsecond=0)} Client: ID:{str(int_id(_peer_id)).rjust(9)} IP:{self.peerTrack[_peer_id]['shost'].rjust(15)} Port:{self.peerTrack[_peer_id]['sport']} Removed{_backend}.") self.transport.write(b'RPTCL'+_peer_id, (self.master,self.peerTrack[_peer_id]['dport'])) #Tell client we have closed the session - 3 times, in case they are on a lossy network self.transport.write(b'MSTCL',(self.peerTrack[_peer_id]['shost'],self.peerTrack[_peer_id]['sport'])) @@ -292,6 +378,8 @@ class Proxy(DatagramProtocol): return if _peer_id in self.peerTrack: + # Keep the DMR ID on its exact generated port; moving it resets + # the backend login, BRIDGES and OPTIONS state. _dport = self.peerTrack[_peer_id]['dport'] self.peerTrack[_peer_id]['sport'] = port self.peerTrack[_peer_id]['shost'] = host @@ -304,11 +392,8 @@ class Proxy(DatagramProtocol): else: if int_id(_peer_id) in self.blackList: return - # Make a list with the available ports - _ports_avail = [port for port in self.connTrack if not self.connTrack[port]] - if _ports_avail: - _dport = random.choice(_ports_avail) - else: + _dport = self.allocate_port() + if _dport is None: return self.connTrack[_dport] = _peer_id self.peerTrack[_peer_id] = {} @@ -322,7 +407,15 @@ class Proxy(DatagramProtocol): self.transport.write(pripacket, (self.master,_dport)) if self.clientinfo and _peer_id != b'\xff\xff\xff\xff': - print(f'{datetime.now().replace(microsecond=0)} New client: ID:{str(int_id(_peer_id)).rjust(9)} IP:{host.rjust(15)} Port:{port}, assigned to port:{_dport}.') + _backend = '' + if len(self.backendPortRanges) > 1: + _backend = ', backend:{}'.format( + backend_number_for_port( + _dport, + self.backendPortRanges, + ) + ) + print(f'{datetime.now().replace(microsecond=0)} New client: ID:{str(int_id(_peer_id)).rjust(9)} IP:{host.rjust(15)} Port:{port}, assigned to port:{_dport}{_backend}.') if self.debug: print(data) return @@ -334,7 +427,6 @@ if __name__ == '__main__': import configparser import argparse import sys - import json import stat import functools @@ -376,6 +468,11 @@ if __name__ == '__main__': ClientInfo = config.getboolean('PROXY','ClientInfo') BlackList = json.loads(config.get('PROXY','BlackList')) IPBlackList = json.loads(config.get('PROXY','IPBlackList')) + BackendPortRanges = backend_port_ranges( + config.get('PROXY','BackendPortRanges', fallback=''), + DestportStart, + DestPortEnd, + ) except configparser.Error as err: print('(PROXY)Error processing configuration file -- {}'.format(err)) @@ -396,6 +493,7 @@ if __name__ == '__main__': BlackList = [1234567] #e.g. {10.0.0.1: 0, 10.0.0.2: 0} IPBlackList = {} + BackendPortRanges = backend_port_ranges(None, DestportStart, DestPortEnd) #******************* @@ -428,6 +526,12 @@ if __name__ == '__main__': ClientInfo = bool_from_env(os.environ['FDPROXY_CLIENTINFO']) if 'FDPROXY_LISTENPORT' in os.environ: ListenPort = int(os.environ['FDPROXY_LISTENPORT']) + if 'FDPROXY_BACKEND_PORT_RANGES' in os.environ: + BackendPortRanges = backend_port_ranges( + os.environ['FDPROXY_BACKEND_PORT_RANGES'], + DestportStart, + DestPortEnd, + ) unixSocket = '/run/priv_control/priv_control.unixsocket' @@ -440,15 +544,17 @@ if __name__ == '__main__': PRIV_HELPER.blocklistFlush() - for port in range(DestportStart,DestPortEnd+1,1): - CONNTRACK[port] = False + for start, end in BackendPortRanges: + for port in range(start,end+1,1): + CONNTRACK[port] = False + log_backend_port_ranges(BackendPortRanges) #If we are listening IPv6 and Master is an IPv4 IPv4Address #IPv6ify the address. if ListenIP == '::' and IsIPv4Address(Master): Master = '::ffff:' + Master - reactor.listenUDP(ListenPort,Proxy(Master,ListenPort,CONNTRACK,PEERTRACK,BlackList,IPBlackList,Timeout,Debug,ClientInfo,DestportStart,DestPortEnd,PRIV_HELPER, RPTLTRACK),interface=ListenIP) + reactor.listenUDP(ListenPort,Proxy(Master,ListenPort,CONNTRACK,PEERTRACK,BlackList,IPBlackList,Timeout,Debug,ClientInfo,DestportStart,DestPortEnd,PRIV_HELPER, RPTLTRACK, BackendPortRanges),interface=ListenIP) def loopingErrHandle(failure): print('(PROXY)(GLOBAL) STOPPING REACTOR TO AVOID MEMORY LEAK: Unhandled error innowtimed loop.\n {}'.format(failure)) @@ -461,7 +567,7 @@ if __name__ == '__main__': if CONNTRACK[port]: count = count+1 - totalPorts = DestPortEnd - DestportStart + 1 + totalPorts = len(CONNTRACK) freePorts = totalPorts - count print("{} ports out of {} in use ({} free)".format(count,totalPorts,freePorts)) diff --git a/report_mux.py b/report_mux.py new file mode 100644 index 0000000..b860301 --- /dev/null +++ b/report_mux.py @@ -0,0 +1,282 @@ +############################################################################### +# Copyright (C) 2026 Simon Adlem, G7RZU +# +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Combine trusted FreeDMR reporting sockets into one compatible feed.""" + +import argparse +import copy +import pickle +import signal + +from twisted.internet import reactor +from twisted.internet.protocol import Factory, ReconnectingClientFactory +from twisted.protocols.basic import NetstringReceiver + +from reporting_const import REPORT_OPCODES + + +def prefixed_system_name(backend_number, system_name): + if not isinstance(system_name, str): + raise ValueError('report system name must be text') + return '{}:{}'.format(backend_number, system_name) + + +def prefix_config(backend_number, config): + if not isinstance(config, dict): + raise ValueError('report configuration must be a dictionary') + return { + prefixed_system_name(backend_number, system_name): system_config + for system_name, system_config in config.items() + } + + +def prefix_bridges(backend_number, bridges): + if not isinstance(bridges, dict): + raise ValueError('report bridge state must be a dictionary') + + rewritten = copy.deepcopy(bridges) + for bridge_systems in rewritten.values(): + if not isinstance(bridge_systems, list): + raise ValueError('report bridge members must be a list') + for system in bridge_systems: + if not isinstance(system, dict) or 'SYSTEM' not in system: + raise ValueError('report bridge member has no SYSTEM field') + system['SYSTEM'] = prefixed_system_name( + backend_number, + system['SYSTEM'], + ) + return rewritten + + +def prefix_bridge_event(backend_number, event): + try: + fields = event.decode('UTF-8').split(',') + except UnicodeDecodeError as err: + raise ValueError('report event is not UTF-8') from err + if len(fields) < 4: + raise ValueError('report event has fewer than four fields') + fields[3] = prefixed_system_name(backend_number, fields[3]) + return ','.join(fields).encode('UTF-8') + + +class ReportMuxClient(NetstringReceiver): + """Downstream endpoint compatible with the FreeDMR report socket.""" + + def connectionMade(self): + self.factory.add_client(self) + + def connectionLost(self, reason): + self.factory.clients.discard(self) + + def stringReceived(self, data): + opcode = data[:1] + if opcode == REPORT_OPCODES['CONFIG_REQ']: + self.factory.send_cached(self, REPORT_OPCODES['CONFIG_SND']) + elif opcode == REPORT_OPCODES['BRIDGE_REQ']: + self.factory.send_cached(self, REPORT_OPCODES['BRIDGE_SND']) + + +class ReportMuxFactory(Factory): + protocol = ReportMuxClient + + def __init__(self, allowed_clients=('*',)): + self.allowed_clients = frozenset(allowed_clients) + self.clients = set() + self._config_by_backend = {} + self._bridges_by_backend = {} + self._cached = {} + + def buildProtocol(self, addr): + if addr.host in self.allowed_clients or '*' in self.allowed_clients: + protocol = ReportMuxClient() + protocol.factory = self + return protocol + print('(REPORT-MUX) rejecting report client {}'.format(addr.host)) + return None + + def add_client(self, client): + self.clients.add(client) + self.send_cached(client, REPORT_OPCODES['CONFIG_SND']) + self.send_cached(client, REPORT_OPCODES['BRIDGE_SND']) + + def send_cached(self, client, opcode): + message = self._cached.get(opcode) + if message is not None: + self._send(client, message) + + def ingest(self, backend_number, message): + opcode = message[:1] + try: + if opcode == REPORT_OPCODES['CONFIG_SND']: + config = pickle.loads(message[1:]) + self._config_by_backend[backend_number] = prefix_config( + backend_number, + config, + ) + self._publish_config() + elif opcode == REPORT_OPCODES['BRIDGE_SND']: + bridges = pickle.loads(message[1:]) + self._bridges_by_backend[backend_number] = prefix_bridges( + backend_number, + bridges, + ) + self._publish_bridges() + elif opcode == REPORT_OPCODES['BRDG_EVENT']: + event = prefix_bridge_event(backend_number, message[1:]) + self._broadcast(REPORT_OPCODES['BRDG_EVENT'] + event) + else: + print( + '(REPORT-MUX) ignoring unsupported opcode from backend {}'. + format(backend_number) + ) + except Exception as err: + print( + '(REPORT-MUX) ignoring malformed report from backend {}: {}'. + format(backend_number, err) + ) + + def backend_disconnected(self, backend_number): + config_removed = self._config_by_backend.pop(backend_number, None) + bridges_removed = self._bridges_by_backend.pop(backend_number, None) + if config_removed is not None: + self._publish_config() + if bridges_removed is not None: + self._publish_bridges() + + def _publish_config(self): + combined = {} + for backend_number in sorted(self._config_by_backend): + combined.update(self._config_by_backend[backend_number]) + self._publish_snapshot(REPORT_OPCODES['CONFIG_SND'], combined) + + def _publish_bridges(self): + combined = {} + for backend_number in sorted(self._bridges_by_backend): + for bridge_name, systems in self._bridges_by_backend[backend_number].items(): + combined.setdefault(bridge_name, []).extend(copy.deepcopy(systems)) + self._publish_snapshot(REPORT_OPCODES['BRIDGE_SND'], combined) + + def _publish_snapshot(self, opcode, value): + message = opcode + pickle.dumps(value, protocol=2) + self._cached[opcode] = message + self._broadcast(message) + + def _broadcast(self, message): + for client in tuple(self.clients): + self._send(client, message) + + def _send(self, client, message): + try: + client.sendString(message) + except Exception as err: + self.clients.discard(client) + print( + '(REPORT-MUX) dropping report client after send failure: {}'. + format(err) + ) + + +class ReportMuxSource(NetstringReceiver): + """One upstream FreeDMR reporting connection.""" + + def __init__(self, factory): + self._factory = factory + + def connectionMade(self): + self.sendString(REPORT_OPCODES['CONFIG_REQ']) + self.sendString(REPORT_OPCODES['BRIDGE_REQ']) + + def stringReceived(self, data): + self._factory.mux.ingest(self._factory.backend_number, data) + + +class ReportMuxSourceFactory(ReconnectingClientFactory): + def __init__(self, backend_number, mux): + self.backend_number = backend_number + self.mux = mux + + def buildProtocol(self, addr): + self.resetDelay() + return ReportMuxSource(self) + + def clientConnectionLost(self, connector, reason): + self.mux.backend_disconnected(self.backend_number) + ReconnectingClientFactory.clientConnectionLost(self, connector, reason) + + def clientConnectionFailed(self, connector, reason): + self.mux.backend_disconnected(self.backend_number) + ReconnectingClientFactory.clientConnectionFailed(self, connector, reason) + + +def parse_backend(value): + try: + number_text, host, port_text = value.split(',', 2) + number = int(number_text) + port = int(port_text) + except ValueError as err: + raise argparse.ArgumentTypeError( + 'backend must be NUMBER,HOST,PORT' + ) from err + if number < 1: + raise argparse.ArgumentTypeError('backend number must be positive') + if not host: + raise argparse.ArgumentTypeError('backend host must not be empty') + if port < 1 or port > 65535: + raise argparse.ArgumentTypeError('backend port must be within 1..65535') + return number, host, port + + +if __name__ == '__main__': + parser = argparse.ArgumentParser( + description='Combine FreeDMR reporting sockets', + ) + parser.add_argument( + '--backend', + action='append', + required=True, + type=parse_backend, + help='numbered source as NUMBER,HOST,PORT; repeat for each backend', + ) + parser.add_argument('--listen-interface', default='127.0.0.1') + parser.add_argument('--listen-port', type=int, default=4321) + parser.add_argument( + '--client', + action='append', + help='permitted reporting client IP; repeat as needed (default: *)', + ) + args = parser.parse_args() + + backend_numbers = [number for number, _host, _port in args.backend] + if len(set(backend_numbers)) != len(backend_numbers): + parser.error('backend numbers must be unique') + + def sig_handler(_signal, _frame): + reactor.stop() + + signal.signal(signal.SIGINT, sig_handler) + signal.signal(signal.SIGTERM, sig_handler) + + mux = ReportMuxFactory(args.client or ('*',)) + reactor.listenTCP(args.listen_port, mux, interface=args.listen_interface) + for backend_number, host, port in args.backend: + reactor.connectTCP( + host, + port, + ReportMuxSourceFactory(backend_number, mux), + ) + reactor.run() diff --git a/tests/harness/fdmr_monitor_compat.py b/tests/harness/fdmr_monitor_compat.py new file mode 100644 index 0000000..31fa224 --- /dev/null +++ b/tests/harness/fdmr_monitor_compat.py @@ -0,0 +1,202 @@ +"""Execute pinned FDMR Monitor consumer paths for compatibility tests.""" + +import ast +import hashlib +import pickle +import subprocess +from time import localtime, strftime, time + + +PINNED_FDMR_MONITOR_COMMIT = 'cd57a18f351bbd370688c902efda4ae7efaf8136' +PINNED_PATCHED_MONITOR_SHA256 = ( + 'ed8445200760d0e1712e67c35347ea49fcb77991c4a2ed42242d841955e685f9' +) +PINNED_BRIDGE_HELPER_SHA256 = ( + 'e485d1d4856b96721565f761ea8db00f96dc6d711e72cd6bf534b171880bfead' +) +PINNED_REPORTING_HELPER_SHA256 = ( + '00993184de75527279b69375ba96ba91f854f7ea8ba0003f477f27d9618fafdb' +) +_REQUIRED_FUNCTIONS = { + 'time_str', + 'add_hb_peer', + 'build_hblink_table', + 'build_bridge_table', + 'load_dictionary', + 'process_message', +} + + +class MonitorCompatibilityError(RuntimeError): + """Raised when pinned monitor sources cannot perform a compatibility check.""" + + +def consume_config_message(monitor_source, message): + """Return the monitor connection table produced from one CONFIG_SND message.""" + namespace = _monitor_namespace(monitor_source) + namespace['process_message'](bytes(message)) + return namespace['CTABLE'] + + +def consume_bridge_message(monitor_source, message): + """Return dashboard bridge rows produced from one BRIDGE_SND message.""" + namespace = _monitor_namespace(monitor_source) + bridges = namespace['load_dictionary'](bytes(message)) + return namespace['build_bridge_table'](bridges) + + +def consume_event_message(monitor_source, message): + """Return dashboard event observations produced from one BRDG_EVENT message.""" + namespace = _monitor_namespace(monitor_source) + dashboard = _DashboardCapture() + namespace.update({ + 'CONF': { + 'GLOBAL': {'TGC_INC': False, 'LH_INC': False}, + 'OPB_FLTR': {'OPB_FILTER': []}, + }, + 'LOGBUF': [], + 'alias_call': _alias, + 'alias_short': _alias, + 'alias_tgid': _alias, + 'dashboard_server': dashboard, + 'db2dict': _discard, + 'rts_update': _discard, + 'subscriber_ids': {}, + 'sys_dict': {'lst_clean': 0}, + 'talkgroup_ids': {}, + }) + namespace['process_message'](bytes(message)) + return namespace['sys_dict'], dashboard.messages + + +def _monitor_namespace(monitor_source): + source = monitor_source.resolve() + if not source.is_file(): + raise MonitorCompatibilityError('FDMR Monitor source not found: {}'.format(source)) + + head = subprocess.run( + ['git', '-C', str(source.parent), 'rev-parse', 'HEAD'], + check=True, + capture_output=True, + text=True, + ).stdout.strip() + if head != PINNED_FDMR_MONITOR_COMMIT: + raise MonitorCompatibilityError( + 'FDMR Monitor source must be at {}, found {}'.format( + PINNED_FDMR_MONITOR_COMMIT, + head, + ) + ) + + source_text = _verified_text( + source, + PINNED_PATCHED_MONITOR_SHA256, + 'patched monitor.py', + ) + tree = ast.parse(source_text, filename=str(source)) + functions = [ + node + for node in tree.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + and node.name in _REQUIRED_FUNCTIONS + ] + found = {node.name for node in functions} + if found != _REQUIRED_FUNCTIONS: + raise MonitorCompatibilityError( + 'pinned monitor decoder functions changed: {}'.format( + ', '.join(sorted(_REQUIRED_FUNCTIONS - found)) + ) + ) + + connection_table = { + 'MASTERS': {}, + 'PEERS': {}, + 'OPENBRIDGES': {}, + 'SETUP': {'LASTHEARD': False}, + } + namespace = { + 'BRIDGES': {}, + 'BRIDGES_RX': '', + 'BTABLE': {'BRIDGES': {}, 'SETUP': {}}, + 'CONFIG': {}, + 'CONFIG_RX': '', + 'CONF': {'GLOBAL': {'BRDG_INC': False, 'HB_INC': True}}, + 'CTABLE': connection_table, + 'OPCODE': { + 'CONFIG_SND': '\x01', + 'BRIDGE_SND': '\x03', + 'LINK_EVENT': '\x06', + 'BRDG_EVENT': '\x07', + 'SERVER_MSG': 'b', + }, + 'int_id': _int_id, + 'loads': pickle.loads, + 'localtime': localtime, + 'logger': _SilentLogger(), + 'strftime': strftime, + 'time': time, + } + _load_helper_functions( + source.parent / 'fdmr_monitor_bridge_state.py', + PINNED_BRIDGE_HELPER_SHA256, + namespace, + ) + _load_helper_functions( + source.parent / 'fdmr_monitor_reporting.py', + PINNED_REPORTING_HELPER_SHA256, + namespace, + ) + extracted = ast.Module(body=functions, type_ignores=[]) + exec(compile(extracted, str(source), 'exec'), namespace) + return namespace + + +def _load_helper_functions(path, expected_digest, namespace): + source_text = _verified_text(path, expected_digest, path.name) + tree = ast.parse(source_text, filename=str(path)) + functions = [ + node + for node in tree.body + if isinstance(node, (ast.FunctionDef, ast.AsyncFunctionDef)) + ] + extracted = ast.Module(body=functions, type_ignores=[]) + exec(compile(extracted, str(path), 'exec'), namespace) + + +def _verified_text(path, expected_digest, description): + if not path.is_file(): + raise MonitorCompatibilityError('{} not found: {}'.format(description, path)) + source_text = path.read_text(encoding='utf-8') + digest = hashlib.sha256(source_text.encode('utf-8')).hexdigest() + if digest != expected_digest: + raise MonitorCompatibilityError( + '{} does not match the deployed pinned source'.format(description) + ) + return source_text + + +def _int_id(value): + if isinstance(value, bytes): + return int.from_bytes(value, 'big') + return int(value) + + +def _alias(value, _aliases): + return str(value) + + +class _DashboardCapture: + def __init__(self): + self.messages = [] + + def broadcast(self, message, group): + self.messages.append((message, group)) + + +class _SilentLogger: + def __getattr__(self, _name): + return _discard + + +def _discard(*_args, **_kwargs): + return None diff --git a/tests/test_auxiliary_tools.py b/tests/test_auxiliary_tools.py index d446529..dfac571 100644 --- a/tests/test_auxiliary_tools.py +++ b/tests/test_auxiliary_tools.py @@ -577,6 +577,236 @@ class AuxiliaryToolTests(unittest.TestCase): finally: self._restore_modules(saved_modules) + def test_proxy_interleaves_new_sessions_across_two_backend_ranges(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + ranges = ((54000, 54001), (54100, 54101)) + conn_track = { + port: False + for limits in ranges + for port in range(limits[0], limits[1] + 1) + } + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54101, + None, + {}, + ranges, + ) + proxy.transport = _FakeTransport() + + peer_ids = [(1000 + index).to_bytes(4, "big") for index in range(4)] + for index, peer_id in enumerate(peer_ids): + proxy.datagramReceived( + b"RPTL" + peer_id, + ("198.51.100.{}".format(index + 1), 40000 + index), + ) + + assignments = [peer_track[peer_id]["dport"] for peer_id in peer_ids] + self.assertEqual( + [ + next( + range_index + for range_index, limits in enumerate(ranges) + if limits[0] <= port <= limits[1] + ) + for port in assignments + ], + [0, 1, 0, 1], + ) + self.assertEqual(len(set(assignments)), 4) + + proxy.datagramReceived( + b"RPTPING" + peer_ids[0], + ("198.51.100.1", 41000), + ) + self.assertEqual(peer_track[peer_ids[0]]["dport"], assignments[0]) + self.assertEqual( + proxy.transport.writes[-1][1], + ("127.0.0.1", assignments[0]), + ) + finally: + self._restore_modules(saved_modules) + + def test_proxy_interleaves_more_than_two_backends_and_reuses_reaped_port(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + ranges = ((54000, 54001), (54100, 54101), (54200, 54201)) + conn_track = { + port: False + for limits in ranges + for port in range(limits[0], limits[1] + 1) + } + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54201, + None, + {}, + ranges, + ) + proxy.transport = _FakeTransport() + + peer_ids = [(2000 + index).to_bytes(4, "big") for index in range(6)] + for index, peer_id in enumerate(peer_ids): + proxy.datagramReceived( + b"RPTL" + peer_id, + ("203.0.113.{}".format(index + 1), 42000 + index), + ) + + assignments = [peer_track[peer_id]["dport"] for peer_id in peer_ids] + self.assertEqual( + [ + next( + range_index + for range_index, limits in enumerate(ranges) + if limits[0] <= port <= limits[1] + ) + for port in assignments + ], + [0, 1, 2, 0, 1, 2], + ) + self.assertEqual(len(set(assignments)), 6) + + released_port = assignments[0] + proxy.reaper(peer_ids[0]) + proxy.datagramReceived( + b"RPTL" + peer_ids[0], + ("203.0.113.20", 43000), + ) + self.assertEqual(peer_track[peer_ids[0]]["dport"], released_port) + finally: + self._restore_modules(saved_modules) + + def test_proxy_backend_ranges_reject_overlap(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + with self.assertRaises(ValueError): + hotspot_proxy_v2.backend_port_ranges( + [[54000, 54010], [54010, 54020]], + 54000, + 54099, + ) + finally: + self._restore_modules(saved_modules) + + def test_proxy_logs_hdstack_ranges_assignment_and_removal(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + ranges = ((54000, 54000), (54100, 54100)) + conn_track = {54000: False, 54100: False} + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + True, + 54000, + 54100, + None, + {}, + ranges, + ) + proxy.transport = _FakeTransport() + first_peer_id = (3001).to_bytes(4, "big") + second_peer_id = (3002).to_bytes(4, "big") + output = io.StringIO() + + with redirect_stdout(output): + hotspot_proxy_v2.log_backend_port_ranges(ranges) + proxy.datagramReceived( + b"RPTL" + first_peer_id, + ("198.51.100.30", 44000), + ) + proxy.datagramReceived( + b"RPTL" + second_peer_id, + ("198.51.100.31", 44001), + ) + proxy.reaper(first_peer_id) + proxy.reaper(second_peer_id) + + messages = output.getvalue() + self.assertIn( + "(PROXY)(HDSTACK) Backend:1 ports:54000-54000.", + messages, + ) + self.assertIn( + "(PROXY)(HDSTACK) Backend:2 ports:54100-54100.", + messages, + ) + self.assertIn("assigned to port:54000, backend:1.", messages) + self.assertIn("assigned to port:54100, backend:2.", messages) + self.assertIn("Removed from backend:1 port:54000.", messages) + self.assertIn("Removed from backend:2 port:54100.", messages) + finally: + self._restore_modules(saved_modules) + + def test_proxy_preserves_single_backend_logging(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: False}, + {}, + [], + {}, + 30, + False, + True, + 54000, + 54000, + None, + {}, + ) + proxy.transport = _FakeTransport() + peer_id = (3001).to_bytes(4, "big") + output = io.StringIO() + + with redirect_stdout(output): + hotspot_proxy_v2.log_backend_port_ranges(((54000, 54000),)) + proxy.datagramReceived( + b"RPTL" + peer_id, + ("198.51.100.30", 44000), + ) + proxy.reaper(peer_id) + + messages = output.getvalue() + self.assertNotIn("(PROXY)(HDSTACK)", messages) + self.assertNotIn("backend:", messages) + self.assertIn("assigned to port:54000.", messages) + self.assertIn("Port:44000 Removed.", messages) + finally: + self._restore_modules(saved_modules) + def _import_proxy_module(self): saved_modules = self._install_proxy_stubs() try: diff --git a/tests/test_hdstack.py b/tests/test_hdstack.py new file mode 100644 index 0000000..f764f03 --- /dev/null +++ b/tests/test_hdstack.py @@ -0,0 +1,426 @@ +import configparser +import tempfile +import unittest +from pathlib import Path + +import config +import hdstack_runtime + + +ROOT = Path(__file__).resolve().parents[1] +DOCKER_CONFIGS = ROOT / 'docker-configs' + + +class HDStackRuntimeTests(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name) + self.main_config = self.root / 'freedmr.cfg' + self.bridge_config = self.root / 'freedmr-bridge.cfg' + self.bridge_rules = self.root / 'rules.py' + self.main_config.write_text(MAIN_CONFIG, encoding='utf-8') + self.bridge_config.write_text(BRIDGE_CONFIG, encoding='utf-8') + self.bridge_rules.write_text('BRIDGES = {}\n', encoding='utf-8') + + def generate(self, instances=2, base_id=23400): + return hdstack_runtime.generate_hdstack( + self.main_config, + self.bridge_config, + self.bridge_rules, + self.root / 'generated', + instances, + base_id, + ) + + def test_two_worker_layout_uses_base_id_plus_one_and_disjoint_ranges(self): + layout = self.generate() + aggregator = _read_config(layout.aggregator_config) + workers = [_read_config(path) for path in layout.worker_configs] + + self.assertEqual(layout.backend_ranges, ((54000, 54099), (54100, 54199))) + self.assertEqual(aggregator.getint('GLOBAL', 'SERVER_ID'), 23400) + self.assertEqual( + [worker.getint('GLOBAL', 'SERVER_ID') for worker in workers], + [23401, 23402], + ) + self.assertEqual( + [worker.getint('SYSTEM', 'PORT') for worker in workers], + [54000, 54100], + ) + self.assertEqual( + [worker.getint('SYSTEM', 'GENERATOR') for worker in workers], + [100, 100], + ) + + def test_five_worker_layout_is_derived_without_hard_coded_second_worker(self): + layout = self.generate(instances=5) + + self.assertEqual( + layout.backend_ranges, + ( + (54000, 54099), + (54100, 54199), + (54200, 54299), + (54300, 54399), + (54400, 54499), + ), + ) + self.assertEqual(len(layout.worker_configs), 5) + self.assertEqual(layout.worker_report_ports, (4322, 4323, 4324, 4325, 4326)) + + def test_aggregator_has_only_external_and_internal_openbridge_systems(self): + layout = self.generate() + aggregator = _read_config(layout.aggregator_config) + + self.assertFalse(aggregator.getboolean('REPORTS', 'REPORT')) + self.assertTrue(aggregator.getboolean('GLOBAL', 'GEN_STAT_BRIDGES')) + self.assertNotIn('SYSTEM', aggregator.sections()) + self.assertNotIn('ECHO', aggregator.sections()) + self.assertIn('EXTERNAL-FBP', aggregator.sections()) + self.assertEqual(aggregator.getint('EXTERNAL-FBP', 'PORT'), 62044) + protocol_sections = [ + section for section in aggregator.sections() + if aggregator.has_option(section, 'MODE') + ] + self.assertTrue(protocol_sections) + self.assertTrue(all(aggregator.get(section, 'MODE') == 'OPENBRIDGE' + for section in protocol_sections)) + + def test_workers_copy_system_policy_and_only_add_aggregator_link(self): + source = _read_config(self.main_config) + layout = self.generate() + workers = [_read_config(path) for path in layout.worker_configs] + + for number, worker in enumerate(workers, 1): + expected = dict(source.items('SYSTEM')) + expected['port'] = str(54000 + (number - 1) * 100) + self.assertEqual(dict(worker.items('SYSTEM')), expected) + self.assertNotIn('EXTERNAL-FBP', worker.sections()) + self.assertEqual( + [ + section for section in worker.sections() + if worker.has_option(section, 'MODE') + ], + ['SYSTEM', 'OBP-HDSTACK-AGGREGATOR'], + ) + self.assertEqual( + worker.getint('OBP-HDSTACK-AGGREGATOR', 'NETWORK_ID'), + 23400, + ) + self.assertFalse(worker.getboolean('GLOBAL', 'ENABLE_API')) + self.assertFalse(worker.getboolean('GLOBAL', 'DATA_GATEWAY')) + + def test_internal_obp_links_are_reciprocal_and_use_expected_ids(self): + layout = self.generate() + aggregator = _read_config(layout.aggregator_config) + workers = [_read_config(path) for path in layout.worker_configs] + + for number, worker in enumerate(workers, 1): + aggregator_link = 'OBP-HDSTACK-WORKER-{}'.format(number) + self.assertEqual( + ( + aggregator.getint(aggregator_link, 'PORT'), + aggregator.getint(aggregator_link, 'TARGET_PORT'), + aggregator.getint(aggregator_link, 'NETWORK_ID'), + ), + (7100 + number, 7000 + number, 23400 + number), + ) + self.assertEqual( + ( + worker.getint('OBP-HDSTACK-AGGREGATOR', 'PORT'), + worker.getint('OBP-HDSTACK-AGGREGATOR', 'TARGET_PORT'), + worker.getint('OBP-HDSTACK-AGGREGATOR', 'NETWORK_ID'), + ), + (7000 + number, 7100 + number, 23400), + ) + self.assertEqual(worker.getint('OBP-HDSTACK-AGGREGATOR', 'PROTO_VER'), 1) + self.assertFalse( + worker.getboolean('OBP-HDSTACK-AGGREGATOR', 'ENHANCED_OBP') + ) + + def test_reporting_mux_uses_workers_only_and_preserves_client_acl(self): + layout = self.generate() + supervisor = layout.supervisor_config.read_text(encoding='utf-8') + workers = [_read_config(path) for path in layout.worker_configs] + + self.assertIn('--listen-port 4321', supervisor) + self.assertIn('--client 127.0.0.1', supervisor) + self.assertIn('--client 192.0.2.20', supervisor) + self.assertIn('--backend 1,127.0.0.1,4322', supervisor) + self.assertIn('--backend 2,127.0.0.1,4323', supervisor) + self.assertNotIn('--backend 0,', supervisor) + self.assertNotIn('--backend 3,', supervisor) + self.assertEqual( + [worker.getint('REPORTS', 'REPORT_PORT') for worker in workers], + [4322, 4323], + ) + self.assertTrue(all( + worker.get('REPORTS', 'REPORT_CLIENTS') == '127.0.0.1' + for worker in workers + )) + + def test_supervisor_runs_aggregator_workers_bridge_proxy_mux_and_playback(self): + layout = self.generate() + supervisor = layout.supervisor_config.read_text(encoding='utf-8') + + for program in ( + 'hdstack-aggregator', + 'hdstack-worker-1', + 'hdstack-worker-2', + 'bridge', + 'playback', + 'proxy', + 'report-mux', + ): + self.assertIn('[program:{}]'.format(program), supervisor) + self.assertIn( + 'bridge.py -c {} -r {}'.format(self.bridge_config, self.bridge_rules), + supervisor, + ) + self.assertIn( + 'FDPROXY_BACKEND_PORT_RANGES="[[54000,54099],[54100,54199]]"', + supervisor, + ) + + def test_reporting_disabled_omits_mux_and_worker_reporting(self): + self.main_config.write_text( + MAIN_CONFIG.replace('REPORT: True', 'REPORT: False'), + encoding='utf-8', + ) + layout = self.generate() + supervisor = layout.supervisor_config.read_text(encoding='utf-8') + + self.assertNotIn('[program:report-mux]', supervisor) + for path in layout.worker_configs: + self.assertFalse(_read_config(path).getboolean('REPORTS', 'REPORT')) + + def test_generated_configs_are_accepted_by_current_freedmr_parser(self): + layout = self.generate() + for path in (layout.aggregator_config,) + layout.worker_configs: + parsed = config.build_config(str(path)) + self.assertIn('GLOBAL', parsed) + self.assertIn('SYSTEMS', parsed) + + def test_mutable_worker_state_files_are_distinct(self): + layout = self.generate(instances=3) + workers = [_read_config(path) for path in layout.worker_configs] + + self.assertEqual( + [worker.get('ALIASES', 'SUB_MAP_FILE') for worker in workers], + ['sub_map-hdstack-1.pkl', 'sub_map-hdstack-2.pkl', 'sub_map-hdstack-3.pkl'], + ) + self.assertEqual( + [worker.get('ALIASES', 'KEYS_FILE') for worker in workers], + ['keys-hdstack-1.json', 'keys-hdstack-2.json', 'keys-hdstack-3.json'], + ) + self.assertTrue(all( + not worker.getboolean('ALIASES', 'TRY_DOWNLOAD') for worker in workers + )) + + def test_generation_does_not_modify_supplied_configuration_or_rules(self): + before = ( + self.main_config.read_bytes(), + self.bridge_config.read_bytes(), + self.bridge_rules.read_bytes(), + ) + + self.generate() + + self.assertEqual( + before, + ( + self.main_config.read_bytes(), + self.bridge_config.read_bytes(), + self.bridge_rules.read_bytes(), + ), + ) + + def test_invalid_instance_count_base_id_and_system_template_are_rejected(self): + for instances in (1, 6): + with self.subTest(instances=instances): + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate(instances=instances) + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate(base_id=0xFFFFFFFF - 1) + + self.main_config.write_text( + MAIN_CONFIG.replace('MODE: MASTER', 'MODE: PEER'), + encoding='utf-8', + ) + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate() + + def test_missing_bridge_files_and_port_collisions_fail_before_generation(self): + self.bridge_rules.unlink() + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate() + self.bridge_rules.write_text('BRIDGES = {}\n', encoding='utf-8') + self.main_config.write_text( + MAIN_CONFIG.replace('PORT: 62044', 'PORT: 7001'), + encoding='utf-8', + ) + with self.assertRaises(hdstack_runtime.HDStackConfigurationError): + self.generate() + + +class HDStackContainerTests(unittest.TestCase): + def test_standard_entrypoint_preserves_single_mode_and_generates_two_to_five(self): + entrypoint = (DOCKER_CONFIGS / 'entrypoint-proxy').read_text(encoding='utf-8') + + self.assertIn('HDSTACK_COUNT="${HDSTACK:-1}"', entrypoint) + self.assertIn('[ "$HDSTACK_COUNT" = 1 ]', entrypoint) + self.assertIn('/etc/supervisor/conf.d/supervisord.conf', entrypoint) + self.assertIn('[ "$HDSTACK_COUNT" -ge 2 ]', entrypoint) + self.assertIn('[ "$HDSTACK_COUNT" -le 5 ]', entrypoint) + self.assertIn('--base-id "$HDSTACK_BASEID"', entrypoint) + self.assertIn('HDSTACK_BRIDGE_CONFIG', entrypoint) + self.assertIn('HDSTACK_BRIDGE_RULES', entrypoint) + + def test_standard_compose_documents_opt_in_without_changing_default(self): + compose = (DOCKER_CONFIGS / 'docker-compose.yml').read_text(encoding='utf-8') + + self.assertIn('#- HDSTACK=2', compose) + self.assertIn('#- HDSTACK_BASEID=23400', compose) + self.assertIn('#- HDSTACK_BRIDGE_CONFIG=', compose) + self.assertIn('#- HDSTACK_BRIDGE_RULES=', compose) + + def test_image_has_no_obsolete_static_hdstack_supervisor(self): + dockerfile = (DOCKER_CONFIGS / 'Dockerfile-ci').read_text(encoding='utf-8') + + self.assertNotIn('supervisord-hdstack.conf', dockerfile) + self.assertFalse((DOCKER_CONFIGS / 'supervisord-hdstack.conf').exists()) + self.assertFalse((DOCKER_CONFIGS / 'docker-compose-hdstack.yml').exists()) + + +def _read_config(path): + parser = configparser.ConfigParser(interpolation=None) + loaded = parser.read(path) + if not loaded: + raise AssertionError('could not read {}'.format(path)) + return parser + + +MAIN_CONFIG = """ +[GLOBAL] +SERVER_ID: 9999 +GEN_STAT_BRIDGES: False +ENABLE_API: True +DATA_GATEWAY: True +USE_ACL: True +REG_ACL: PERMIT:ALL +SUB_ACL: DENY:1 +TGID_TS1_ACL: PERMIT:ALL +TGID_TS2_ACL: PERMIT:ALL + +[REPORTS] +REPORT: True +REPORT_PORT: 4321 +REPORT_CLIENTS: 127.0.0.1,192.0.2.20 + +[LOGGER] +LOG_FILE: /dev/null +LOG_HANDLERS: console-timed +LOG_LEVEL: INFO +LOG_NAME: FreeDMR + +[ALIASES] +TRY_DOWNLOAD: True +PATH: ./json/ +SUB_MAP_FILE: sub_map.pkl +KEYS_FILE: keys.json + +[ALLSTAR] +ENABLED: False + +[EXTERNAL-FBP] +MODE: OPENBRIDGE +ENABLED: True +IP: 127.0.0.1 +PORT: 62044 +NETWORK_ID: 23499 +PASSPHRASE: external +TARGET_IP: 127.0.0.1 +TARGET_PORT: 62045 +USE_ACL: True +SUB_ACL: DENY:1 +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: True +ENHANCED_OBP: True +PROTO_VER: 5 + +[SYSTEM] +MODE: MASTER +ENABLED: True +REPEAT: True +MAX_PEERS: 1 +IP: 127.0.0.1 +PORT: 54000 +PASSPHRASE: +GROUP_HANGTIME: 5 +USE_ACL: True +REG_ACL: DENY:1 +SUB_ACL: DENY:1 +TGID_TS1_ACL: PERMIT:ALL +TGID_TS2_ACL: PERMIT:ALL +DEFAULT_UA_TIMER: 10 +SINGLE_MODE: True +VOICE_IDENT: True +DIAL_A_TG: True +DYNAMIC_TG_ROUTING: True +NETWORK_DIRECT_DIAL_SLOT: 1 +TS1_STATIC: +TS2_STATIC: 2350 +GENERATOR: 100 +ALLOW_UNREG_ID: True +PROXY_CONTROL: True +ANNOUNCEMENT_LANGUAGE: en_GB + +[ECHO] +MODE: PEER +ENABLED: True +PORT: 54916 +""".lstrip() + + +BRIDGE_CONFIG = """ +[GLOBAL] +SERVER_ID: 23500 +USE_ACL: True +REG_ACL: PERMIT:ALL +SUB_ACL: DENY:1 +TGID_TS1_ACL: PERMIT:ALL +TGID_TS2_ACL: PERMIT:ALL + +[REPORTS] +REPORT: False +REPORT_PORT: 4421 + +[LOGGER] + +[ALIASES] +TRY_DOWNLOAD: False + +[ALLSTAR] +ENABLED: False + +[TO-AGGREGATOR] +MODE: OPENBRIDGE +ENABLED: True +IP: 127.0.0.1 +PORT: 63001 +NETWORK_ID: 23400 +PASSPHRASE: bridge +TARGET_IP: 127.0.0.1 +TARGET_PORT: 63000 +USE_ACL: False +SUB_ACL: DENY:1 +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 +""".lstrip() + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_hdstack_udp.py b/tests/test_hdstack_udp.py new file mode 100644 index 0000000..14454bb --- /dev/null +++ b/tests/test_hdstack_udp.py @@ -0,0 +1,353 @@ +import os +import pickle +import re +import socket +import subprocess +import tempfile +import time +import unittest +from pathlib import Path + +import hdstack_runtime +from tests.harness.fdmr_monitor_compat import consume_config_message +from reporting_const import REPORT_OPCODES +from tests.harness.udp_blackbox import ( + DependencySandbox, + FreeDmrProcess, + HbpRepeater, + HotspotProxyProcess, + PacketSpec, + PROXY_RUNTIME_MODULES, + free_udp_port, + free_udp_port_range, + require_udp_integration_enabled, + write_hotspot_proxy_config, +) +from tests.test_hdstack import BRIDGE_CONFIG, MAIN_CONFIG + + +ROOT = Path(__file__).resolve().parents[1] + + +class HDStackUdpBlackBoxTests(unittest.TestCase): + def test_two_workers_route_through_aggregator_and_report_through_mux(self): + require_udp_integration_enabled() + if not _udp_ports_available((7001, 7002, 7101, 7102)): + self.skipTest('generated internal HDStack UDP ports are in use') + + with tempfile.TemporaryDirectory(prefix='freedmr-hdstack-udp-') as directory: + temporary = Path(directory) + hbp_base = free_udp_port_range(4) + report_base = _free_tcp_port_range(3) + bridge_port = free_udp_port() + bridge_peer_port = free_udp_port() + proxy_port = free_udp_port() + + main_config = temporary / 'freedmr.cfg' + bridge_config = temporary / 'freedmr-bridge.cfg' + bridge_rules = temporary / 'rules.py' + proxy_config = temporary / 'proxy.cfg' + bridge_link = _openbridge( + 'OBP-BRIDGE', + bridge_port, + bridge_peer_port, + 23500, + 'bridge', + ) + main_text = MAIN_CONFIG.replace( + '[EXTERNAL-FBP]\nMODE: OPENBRIDGE\nENABLED: True', + '[EXTERNAL-FBP]\nMODE: OPENBRIDGE\nENABLED: False', + ) + main_text = main_text.replace('TRY_DOWNLOAD: True', 'TRY_DOWNLOAD: False') + main_text = main_text.replace('GROUP_HANGTIME: 5', 'GROUP_HANGTIME: 0') + main_text = main_text.replace('PORT: 54000', 'PORT: {}'.format(hbp_base)) + main_text = main_text.replace('GENERATOR: 100', 'GENERATOR: 2') + main_text = main_text.replace('REPORT_PORT: 4321', 'REPORT_PORT: {}'.format(report_base)) + main_text = main_text.replace('\n[SYSTEM]', '\n{}\n[SYSTEM]'.format(bridge_link)) + main_config.write_text(main_text, encoding='utf-8') + + bridge_text = BRIDGE_CONFIG.replace('PORT: 63001', 'PORT: {}'.format(bridge_peer_port)) + bridge_text = bridge_text.replace( + 'TARGET_PORT: 63000', + 'TARGET_PORT: {}'.format(bridge_port), + ) + bridge_config.write_text(bridge_text, encoding='utf-8') + bridge_rules.write_text('BRIDGES = {}\n', encoding='utf-8') + + layout = hdstack_runtime.generate_hdstack( + main_config, + bridge_config, + bridge_rules, + temporary / 'generated', + 2, + 23400, + ) + write_hotspot_proxy_config( + proxy_config, + listen_port=proxy_port, + dest_port_start=hbp_base, + dest_port_end=hbp_base + 3, + ) + proxy_config.write_text( + proxy_config.read_text(encoding='utf-8').replace( + 'ClientInfo: False', + 'ClientInfo: True', + ), + encoding='utf-8', + ) + with proxy_config.open('a', encoding='utf-8') as config_file: + config_file.write( + 'BackendPortRanges: [[{},{}],[{},{}]]\n'.format( + hbp_base, + hbp_base + 1, + hbp_base + 2, + hbp_base + 3, + ) + ) + + dependencies = DependencySandbox( + ROOT, + extra_runtime_modules=PROXY_RUNTIME_MODULES, + ) + python = str(dependencies.resolve_python()) + processes = [] + logs = [] + repeaters = [] + try: + for config_path in (layout.aggregator_config,) + layout.worker_configs: + process = FreeDmrProcess(ROOT, config_path, python) + process.__enter__() + processes.append(process) + process.wait_for_start() + + bridge_process, bridge_log = _start_process( + [ + python, + 'bridge.py', + '-c', + str(bridge_config), + '-r', + str(bridge_rules), + '-l', + 'INFO', + ], + temporary / 'bridge.log', + ) + processes.append(bridge_process) + logs.append(bridge_log) + + mux_process, mux_log = _start_process( + [ + python, + 'report_mux.py', + '--listen-interface', + '127.0.0.1', + '--listen-port', + str(report_base), + '--client', + '127.0.0.1', + '--backend', + '1,127.0.0.1,{}'.format(report_base + 1), + '--backend', + '2,127.0.0.1,{}'.format(report_base + 2), + ], + temporary / 'mux.log', + ) + processes.append(mux_process) + logs.append(mux_log) + + proxy = HotspotProxyProcess(ROOT, proxy_config, python) + proxy.__enter__() + processes.append(proxy) + proxy.wait_for_start() + time.sleep(0.5) + _assert_processes_running(processes) + + first = HbpRepeater(proxy_port, 1001, bind_host='127.0.0.2') + second = HbpRepeater(proxy_port, 1002, bind_host='127.0.0.2') + repeaters.extend((first, second)) + first.login() + second.login() + first.send_dmr( + PacketSpec( + peer_id=1001, + rf_src=3120001, + dst_id=2350, + slot=2, + ) + ) + received = second.recv(timeout=3.0) + self.assertEqual(received.fields['dst_id'], (2350).to_bytes(3, 'big')) + + assignments = [ + int(port) + for port in re.findall(r'assigned to port:(\d+)', proxy.output()) + ] + self.assertGreaterEqual(len(assignments), 2) + self.assertIn(assignments[0], range(hbp_base, hbp_base + 2)) + self.assertIn(assignments[1], range(hbp_base + 2, hbp_base + 4)) + + report_message, combined = _request_report_config(report_base) + self.assertTrue(any(name.startswith('1:') for name in combined)) + self.assertTrue(any(name.startswith('2:') for name in combined)) + self.assertFalse(any(name.startswith('0:') for name in combined)) + + monitor_source = os.environ.get('FREEDMR_MONITOR_SOURCE') + if monitor_source: + connection_table = consume_config_message( + Path(monitor_source), + report_message, + ) + self.assertEqual( + set(connection_table['MASTERS']), + { + '1:SYSTEM-0', + '1:SYSTEM-1', + '2:SYSTEM-0', + '2:SYSTEM-1', + }, + ) + self.assertEqual( + set(connection_table['OPENBRIDGES']), + { + '1:OBP-HDSTACK-AGGREGATOR', + '2:OBP-HDSTACK-AGGREGATOR', + }, + ) + finally: + for repeater in repeaters: + repeater.close() + for process in reversed(processes): + if isinstance(process, FreeDmrProcess): + process.__exit__(None, None, None) + else: + _stop_process(process) + for log in logs: + log.close() + dependencies.cleanup() + + +def _openbridge(name, port, target_port, network_id, passphrase): + return """[{name}] +MODE: OPENBRIDGE +ENABLED: True +IP: 127.0.0.1 +PORT: {port} +NETWORK_ID: {network_id} +PASSPHRASE: {passphrase} +TARGET_IP: 127.0.0.1 +TARGET_PORT: {target_port} +USE_ACL: False +SUB_ACL: DENY:1 +TGID_ACL: PERMIT:ALL +RELAX_CHECKS: False +ENHANCED_OBP: True +PROTO_VER: 5 +""".format( + name=name, + port=port, + target_port=target_port, + network_id=network_id, + passphrase=passphrase, + ) + + +def _start_process(command, log_path): + log = log_path.open('w+') + process = subprocess.Popen( + command, + cwd=str(ROOT), + stdout=log, + stderr=subprocess.STDOUT, + text=True, + ) + return process, log + + +def _stop_process(process): + if process.poll() is not None: + return + process.terminate() + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait(timeout=5) + + +def _assert_processes_running(processes): + for process in processes: + proc = process.proc if isinstance(process, FreeDmrProcess) else process + if proc.poll() is not None: + raise AssertionError('HDStack process exited during startup') + + +def _udp_ports_available(ports): + sockets = [] + try: + for port in ports: + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + sock.bind(('127.0.0.1', port)) + sockets.append(sock) + return True + except OSError: + return False + finally: + for sock in sockets: + sock.close() + + +def _free_tcp_port_range(count): + for _attempt in range(100): + first = free_udp_port() + sockets = [] + try: + for port in range(first, first + count): + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.bind(('127.0.0.1', port)) + sockets.append(sock) + return first + except OSError: + pass + finally: + for sock in sockets: + sock.close() + raise RuntimeError('could not reserve a contiguous TCP port range') + + +def _request_report_config(port): + with socket.create_connection(('127.0.0.1', port), timeout=3.0) as report: + report.sendall(b'1:' + REPORT_OPCODES['CONFIG_REQ'] + b',') + deadline = time.monotonic() + 3.0 + while time.monotonic() < deadline: + message = _receive_netstring(report) + if message[:1] != REPORT_OPCODES['CONFIG_SND']: + continue + value = pickle.loads(message[1:]) + if any(name.startswith('1:') for name in value) and any( + name.startswith('2:') for name in value + ): + return message, value + raise AssertionError('report MUX did not provide both worker configurations') + + +def _receive_netstring(sock): + length = bytearray() + while True: + value = sock.recv(1) + if value == b':': + break + if not value: + raise AssertionError('report MUX closed during netstring length') + length.extend(value) + payload_length = int(length.decode('ascii')) + payload = b'' + while len(payload) < payload_length: + payload += sock.recv(payload_length - len(payload)) + if sock.recv(1) != b',': + raise AssertionError('invalid report MUX netstring terminator') + return payload + + +if __name__ == '__main__': + unittest.main() diff --git a/tests/test_report_mux.py b/tests/test_report_mux.py new file mode 100644 index 0000000..4ec347e --- /dev/null +++ b/tests/test_report_mux.py @@ -0,0 +1,309 @@ +import os +import argparse +import pickle +import unittest +from contextlib import redirect_stdout +from io import StringIO +from pathlib import Path + +import report_mux +from reporting_const import REPORT_OPCODES +from tests.harness.fdmr_monitor_compat import ( + consume_bridge_message, + consume_config_message, + consume_event_message, +) + + +class ReportMuxTests(unittest.TestCase): + def test_merges_config_and_disambiguates_colliding_system_names(self): + mux = report_mux.ReportMuxFactory() + client = _FakeReportClient() + mux.add_client(client) + + mux.ingest(1, _report_message('CONFIG_SND', {'SYSTEM-001': {'PORT': 54000}})) + mux.ingest(2, _report_message('CONFIG_SND', {'SYSTEM-001': {'PORT': 54100}})) + + combined = pickle.loads(client.messages[-1][1:]) + self.assertEqual( + combined, + { + '1:SYSTEM-001': {'PORT': 54000}, + '2:SYSTEM-001': {'PORT': 54100}, + }, + ) + + def test_merges_bridge_state_and_changes_only_system_field(self): + mux = report_mux.ReportMuxFactory() + client = _FakeReportClient() + mux.add_client(client) + backend_one = { + '2350': [ + { + 'SYSTEM': 'SYSTEM-001', + 'TS': 2, + 'TGID': 2350, + 'ACTIVE': True, + } + ] + } + backend_two = { + '2350': [ + { + 'SYSTEM': 'SYSTEM-001', + 'TS': 1, + 'TGID': 2350, + 'ACTIVE': False, + } + ] + } + + mux.ingest(1, _report_message('BRIDGE_SND', backend_one)) + mux.ingest(2, _report_message('BRIDGE_SND', backend_two)) + + combined = pickle.loads(client.messages[-1][1:]) + self.assertEqual( + combined['2350'], + [ + { + 'SYSTEM': '1:SYSTEM-001', + 'TS': 2, + 'TGID': 2350, + 'ACTIVE': True, + }, + { + 'SYSTEM': '2:SYSTEM-001', + 'TS': 1, + 'TGID': 2350, + 'ACTIVE': False, + }, + ], + ) + self.assertEqual(backend_one['2350'][0]['SYSTEM'], 'SYSTEM-001') + self.assertEqual(backend_two['2350'][0]['SYSTEM'], 'SYSTEM-001') + + def test_prefixes_only_bridge_event_system_field(self): + event = b'GROUP VOICE,END,RX,SYSTEM-001,1234,5678,9012,2,2350,4.25' + + rewritten = report_mux.prefix_bridge_event(3, event) + + self.assertEqual( + rewritten, + b'GROUP VOICE,END,RX,3:SYSTEM-001,1234,5678,9012,2,2350,4.25', + ) + + def test_malformed_source_does_not_disturb_healthy_backend(self): + mux = report_mux.ReportMuxFactory() + client = _FakeReportClient() + mux.add_client(client) + mux.ingest(1, _report_message('CONFIG_SND', {'SYSTEM-001': {'PORT': 54000}})) + healthy_message_count = len(client.messages) + + with redirect_stdout(StringIO()): + mux.ingest(2, REPORT_OPCODES['CONFIG_SND'] + b'not-a-pickle') + + self.assertEqual(len(client.messages), healthy_message_count) + combined = pickle.loads(client.messages[-1][1:]) + self.assertEqual(combined, {'1:SYSTEM-001': {'PORT': 54000}}) + + def test_disconnect_removes_only_that_backend_snapshot(self): + mux = report_mux.ReportMuxFactory() + client = _FakeReportClient() + mux.add_client(client) + mux.ingest(1, _report_message('CONFIG_SND', {'SYSTEM-001': {}})) + mux.ingest(2, _report_message('CONFIG_SND', {'SYSTEM-001': {}})) + + mux.backend_disconnected(1) + + combined = pickle.loads(client.messages[-1][1:]) + self.assertEqual(combined, {'2:SYSTEM-001': {}}) + + def test_new_client_receives_merged_cached_snapshots(self): + mux = report_mux.ReportMuxFactory() + mux.ingest(1, _report_message('CONFIG_SND', {'SYSTEM-001': {}})) + mux.ingest( + 1, + _report_message( + 'BRIDGE_SND', + {'2350': [{'SYSTEM': 'SYSTEM-001', 'TS': 2}]}, + ), + ) + client = _FakeReportClient() + + mux.add_client(client) + + self.assertEqual( + [message[:1] for message in client.messages], + [REPORT_OPCODES['CONFIG_SND'], REPORT_OPCODES['BRIDGE_SND']], + ) + + def test_failed_downstream_client_is_isolated(self): + mux = report_mux.ReportMuxFactory() + healthy = _FakeReportClient() + failed = _FakeReportClient(fail=True) + mux.clients.update((healthy, failed)) + + with redirect_stdout(StringIO()): + mux.ingest( + 1, + REPORT_OPCODES['BRDG_EVENT'] + + b'GROUP VOICE,START,RX,SYSTEM-001,1,2,3,1,4', + ) + + self.assertEqual( + healthy.messages, + [ + REPORT_OPCODES['BRDG_EVENT'] + + b'GROUP VOICE,START,RX,1:SYSTEM-001,1,2,3,1,4' + ], + ) + self.assertNotIn(failed, mux.clients) + + def test_backend_prefixes_are_stable_for_arbitrary_backend_count(self): + self.assertEqual( + [ + report_mux.prefixed_system_name(number, 'SYSTEM-001') + for number in (1, 2, 3) + ], + ['1:SYSTEM-001', '2:SYSTEM-001', '3:SYSTEM-001'], + ) + + def test_report_client_access_matches_configured_allow_list(self): + restricted = report_mux.ReportMuxFactory(('127.0.0.1',)) + self.assertIsInstance( + restricted.buildProtocol(_FakeAddress('127.0.0.1')), + report_mux.ReportMuxClient, + ) + with redirect_stdout(StringIO()): + self.assertIsNone(restricted.buildProtocol(_FakeAddress('192.0.2.1'))) + + unrestricted = report_mux.ReportMuxFactory(('*',)) + self.assertIsInstance( + unrestricted.buildProtocol(_FakeAddress('192.0.2.1')), + report_mux.ReportMuxClient, + ) + + def test_backend_parser_requires_deterministic_positive_number(self): + self.assertEqual( + report_mux.parse_backend('2,127.0.0.1,4323'), + (2, '127.0.0.1', 4323), + ) + with self.assertRaises(argparse.ArgumentTypeError): + report_mux.parse_backend('0,127.0.0.1,4323') + + +class PinnedFdmrMonitorCompatibilityTests(unittest.TestCase): + @classmethod + def setUpClass(cls): + source = os.environ.get('FREEDMR_MONITOR_SOURCE') + if not source: + raise unittest.SkipTest( + 'set FREEDMR_MONITOR_SOURCE to the patched pinned monitor.py' + ) + cls.monitor_source = Path(source) + + def test_mux_output_is_consumed_by_pinned_dashboard(self): + mux = report_mux.ReportMuxFactory() + client = _FakeReportClient() + mux.add_client(client) + + for backend_number in (1, 2): + mux.ingest( + backend_number, + _report_message( + 'CONFIG_SND', + { + 'SYSTEM-0': { + 'ENABLED': True, + 'MODE': 'MASTER', + 'REPEAT': True, + 'PEERS': {}, + }, + 'OBP-HDSTACK-AGGREGATOR': { + 'ENABLED': True, + 'MODE': 'OPENBRIDGE', + 'NETWORK_ID': (23400).to_bytes(4, 'big'), + 'TARGET_IP': '127.0.0.1', + 'TARGET_PORT': 7100 + backend_number, + }, + }, + ), + ) + + config_message = client.messages[-1] + connection_table = consume_config_message( + self.monitor_source, + config_message, + ) + self.assertEqual( + set(connection_table['MASTERS']), + {'1:SYSTEM-0', '2:SYSTEM-0'}, + ) + self.assertEqual( + set(connection_table['OPENBRIDGES']), + { + '1:OBP-HDSTACK-AGGREGATOR', + '2:OBP-HDSTACK-AGGREGATOR', + }, + ) + + bridge_entry = { + 'SYSTEM': 'SYSTEM-0', + 'TS': 2, + 'TGID': (2350).to_bytes(3, 'big'), + 'ACTIVE': True, + 'TIMEOUT': 600, + 'TO_TYPE': 'ON', + 'OFF': [], + 'ON': [(2350).to_bytes(3, 'big')], + 'RESET': [], + 'TIMER': 200.0, + } + mux.ingest(1, _report_message('BRIDGE_SND', {'2350': [bridge_entry]})) + mux.ingest(2, _report_message('BRIDGE_SND', {'2350': [bridge_entry]})) + + bridge_rows = consume_bridge_message( + self.monitor_source, + client.messages[-1], + ) + self.assertEqual( + [row['SYSTEM'] for row in bridge_rows['2350']], + ['1:SYSTEM-0', '2:SYSTEM-0'], + ) + + mux.ingest( + 1, + REPORT_OPCODES['BRDG_EVENT'] + + b'GROUP VOICE,START,RX,SYSTEM-0,1234,5678,9012,2,2350', + ) + event_state, dashboard_messages = consume_event_message( + self.monitor_source, + client.messages[-1], + ) + self.assertEqual(event_state['1234']['sys'], '1:SYSTEM-0') + self.assertTrue(dashboard_messages) + self.assertEqual(dashboard_messages[-1][1], 'log') + + +class _FakeAddress: + def __init__(self, host): + self.host = host + + +class _FakeReportClient: + def __init__(self, fail=False): + self.fail = fail + self.messages = [] + + def sendString(self, message): + if self.fail: + raise RuntimeError('closed') + self.messages.append(message) + + +def _report_message(opcode, value): + return REPORT_OPCODES[opcode] + pickle.dumps(value, protocol=2) + + +if __name__ == '__main__': + unittest.main()