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 twisted.internet.testing import StringTransport 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') 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 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) def _netstring(message): return str(len(message)).encode('ascii') + b':' + message + b',' if __name__ == '__main__': unittest.main()