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/tests/test_report_mux.py

352 lines
11 KiB

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

Powered by TurnKey Linux.