Handle large HDStack reporting snapshots

master
Simon 3 weeks ago
parent d7f119af7f
commit 3738f563a7

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.