You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
FreeDMR/report_mux.py

300 lines
11 KiB

###############################################################################
# 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
REPORT_FRAME_MAX_LENGTH = 10 * 1024 * 1024
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."""
MAX_LENGTH = REPORT_FRAME_MAX_LENGTH
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):
print(
'(REPORT-MUX) backend {} reporting connection lost: {}'.format(
self.backend_number,
reason.getErrorMessage(),
)
)
self.mux.backend_disconnected(self.backend_number)
ReconnectingClientFactory.clientConnectionLost(self, connector, reason)
def clientConnectionFailed(self, connector, reason):
print(
'(REPORT-MUX) backend {} reporting connection failed: {}'.format(
self.backend_number,
reason.getErrorMessage(),
)
)
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()

Powered by TurnKey Linux.