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