Add configurable HDStack backend scaling

master
Simon 4 weeks ago
parent 26c22834c8
commit d355654d8a

@ -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).

@ -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:

@ -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

@ -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.

@ -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.

@ -0,0 +1,563 @@
###############################################################################
# Copyright (C) 2026 Simon Adlem, G7RZU <g7rzu@gb7fr.org.uk>
#
# 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())

@ -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))

@ -0,0 +1,282 @@
###############################################################################
# Copyright (C) 2026 Simon Adlem, G7RZU <g7rzu@gb7fr.org.uk>
#
# 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()

@ -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

@ -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:

@ -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()

@ -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()

@ -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()
Loading…
Cancel
Save

Powered by TurnKey Linux.