From 3738f563a74fef4b6ab4f2460e10d8acc7868a14 Mon Sep 17 00:00:00 2001 From: Simon Date: Sat, 5 Sep 2026 02:59:45 +0000 Subject: [PATCH] Handle large HDStack reporting snapshots --- docs/reporting-relay.md | 7 ++++++- report_mux.py | 17 ++++++++++++++++ tests/test_report_mux.py | 42 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 65 insertions(+), 1 deletion(-) diff --git a/docs/reporting-relay.md b/docs/reporting-relay.md index aa47bd8..c40b5b7 100644 --- a/docs/reporting-relay.md +++ b/docs/reporting-relay.md @@ -47,4 +47,9 @@ 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. +not affect the radio proxy. Its backend connections remain open between +updates and accept bounded reporting frames up to 10 MiB so populated worker +snapshots do not trigger Twisted's smaller default netstring limit. A lost +backend connection is reported with its backend number and is re-established +independently. See [hdstack.md](hdstack.md) for the generated container +topology and configuration. diff --git a/report_mux.py b/report_mux.py index b860301..18974ea 100644 --- a/report_mux.py +++ b/report_mux.py @@ -30,6 +30,9 @@ 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') @@ -194,6 +197,8 @@ class ReportMuxFactory(Factory): class ReportMuxSource(NetstringReceiver): """One upstream FreeDMR reporting connection.""" + MAX_LENGTH = REPORT_FRAME_MAX_LENGTH + def __init__(self, factory): self._factory = factory @@ -215,10 +220,22 @@ class ReportMuxSourceFactory(ReconnectingClientFactory): 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) diff --git a/tests/test_report_mux.py b/tests/test_report_mux.py index 4ec347e..c49f630 100644 --- a/tests/test_report_mux.py +++ b/tests/test_report_mux.py @@ -8,6 +8,7 @@ from pathlib import Path import report_mux from reporting_const import REPORT_OPCODES +from twisted.internet.testing import StringTransport from tests.harness.fdmr_monitor_compat import ( consume_bridge_message, consume_config_message, @@ -191,6 +192,43 @@ class ReportMuxTests(unittest.TestCase): with self.assertRaises(argparse.ArgumentTypeError): report_mux.parse_backend('0,127.0.0.1,4323') + def test_source_accepts_snapshot_above_twisted_default_limit(self): + mux = report_mux.ReportMuxFactory() + downstream = _FakeReportClient() + mux.add_client(downstream) + source = report_mux.ReportMuxSourceFactory(1, mux).buildProtocol( + _FakeAddress('127.0.0.1') + ) + transport = StringTransport() + source.makeConnection(transport) + message = _report_message( + 'CONFIG_SND', + {'SYSTEM-001': {'PEERS': {'large': b'x' * 100_000}}}, + ) + self.assertGreater(len(message), 99_999) + + source.dataReceived(_netstring(message)) + + self.assertEqual(source.brokenPeer, 0) + combined = pickle.loads(downstream.messages[-1][1:]) + self.assertEqual( + combined, + {'1:SYSTEM-001': {'PEERS': {'large': b'x' * 100_000}}}, + ) + + def test_source_rejects_snapshot_above_reporting_limit(self): + mux = report_mux.ReportMuxFactory() + source = report_mux.ReportMuxSourceFactory(1, mux).buildProtocol( + _FakeAddress('127.0.0.1') + ) + source.makeConnection(StringTransport()) + + source.dataReceived( + str(report_mux.REPORT_FRAME_MAX_LENGTH + 1).encode('ascii') + b':' + ) + + self.assertEqual(source.brokenPeer, 1) + class PinnedFdmrMonitorCompatibilityTests(unittest.TestCase): @classmethod @@ -305,5 +343,9 @@ def _report_message(opcode, value): return REPORT_OPCODES[opcode] + pickle.dumps(value, protocol=2) +def _netstring(message): + return str(len(message)).encode('ascii') + b':' + message + b',' + + if __name__ == '__main__': unittest.main()