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.
300 lines
11 KiB
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()
|