import os import pickle import re import socket import subprocess import tempfile import time import unittest from pathlib import Path import hdstack_runtime from tests.harness.fdmr_monitor_compat import consume_config_message from reporting_const import REPORT_OPCODES from tests.harness.udp_blackbox import ( HBPF_DATA_SYNC, DependencySandbox, FreeDmrProcess, HbpRepeater, HotspotProxyProcess, PacketSpec, PROXY_RUNTIME_MODULES, free_udp_port, free_udp_port_range, require_udp_integration_enabled, write_hotspot_proxy_config, ) from tests.test_hdstack import BRIDGE_CONFIG, MAIN_CONFIG ROOT = Path(__file__).resolve().parents[1] class HDStackUdpBlackBoxTests(unittest.TestCase): def test_two_workers_route_through_aggregator_and_report_through_mux(self): require_udp_integration_enabled() generated_ports = ( 7001, 7002, 7101, 7102, 7200, 7201, 54915, 54916, 54917, 54918, ) if not _udp_ports_available(generated_ports): self.skipTest('generated internal HDStack UDP ports are in use') with tempfile.TemporaryDirectory(prefix='freedmr-hdstack-udp-') as directory: temporary = Path(directory) hbp_base = free_udp_port_range(4) report_base = _free_tcp_port_range(3) proxy_port = free_udp_port() main_config = temporary / 'freedmr.cfg' bridge_config = temporary / 'freedmr-bridge.cfg' bridge_rules = temporary / 'rules.py' proxy_config = temporary / 'proxy.cfg' main_text = MAIN_CONFIG.replace( '[EXTERNAL-FBP]\nMODE: OPENBRIDGE\nENABLED: True', '[EXTERNAL-FBP]\nMODE: OPENBRIDGE\nENABLED: False', ) main_text = main_text.replace('TRY_DOWNLOAD: True', 'TRY_DOWNLOAD: False') main_text = main_text.replace('GROUP_HANGTIME: 5', 'GROUP_HANGTIME: 0') main_text = main_text.replace('PORT: 54000', 'PORT: {}'.format(hbp_base)) main_text = main_text.replace('GENERATOR: 100', 'GENERATOR: 2') main_text = main_text.replace('REPORT_PORT: 4321', 'REPORT_PORT: {}'.format(report_base)) main_config.write_text(main_text, encoding='utf-8') bridge_config.write_text(BRIDGE_CONFIG, encoding='utf-8') bridge_rules.write_text('BRIDGES = {}\n', encoding='utf-8') layout = hdstack_runtime.generate_hdstack( main_config, bridge_config, bridge_rules, temporary / 'generated', 2, 23400, ) write_hotspot_proxy_config( proxy_config, listen_port=proxy_port, dest_port_start=hbp_base, dest_port_end=hbp_base + 3, ) proxy_config.write_text( proxy_config.read_text(encoding='utf-8').replace( 'ClientInfo: False', 'ClientInfo: True', ), encoding='utf-8', ) with proxy_config.open('a', encoding='utf-8') as config_file: config_file.write( 'BackendPortRanges: [[{},{}],[{},{}]]\n'.format( hbp_base, hbp_base + 1, hbp_base + 2, hbp_base + 3, ) ) dependencies = DependencySandbox( ROOT, extra_runtime_modules=PROXY_RUNTIME_MODULES, ) python = str(dependencies.resolve_python()) processes = [] logs = [] repeaters = [] try: for number, config_path in enumerate(layout.loro_configs, 1): loro_process, loro_log = _start_process( [ python, 'playback.py', '-c', str(config_path), '-l', 'INFO', ], temporary / 'loro-{}.log'.format(number), ) processes.append(loro_process) logs.append(loro_log) for config_path in (layout.aggregator_config,) + layout.worker_configs: process = FreeDmrProcess(ROOT, config_path, python) process.__enter__() processes.append(process) process.wait_for_start() bridge_process, bridge_log = _start_process( [ python, 'bridge.py', '-c', str(layout.bridge_config), '-r', str(bridge_rules), '-l', 'INFO', ], temporary / 'bridge.log', ) processes.append(bridge_process) logs.append(bridge_log) mux_process, mux_log = _start_process( [ python, 'report_mux.py', '--listen-interface', '127.0.0.1', '--listen-port', str(report_base), '--client', '127.0.0.1', '--backend', '1,127.0.0.1,{}'.format(report_base + 1), '--backend', '2,127.0.0.1,{}'.format(report_base + 2), ], temporary / 'mux.log', ) processes.append(mux_process) logs.append(mux_log) proxy = HotspotProxyProcess(ROOT, proxy_config, python) proxy.__enter__() processes.append(proxy) proxy.wait_for_start() time.sleep(0.5) _assert_processes_running(processes) # A first BCKA can precede the reciprocal listener. Allow the # next scheduled exchange to open both FBP send gates. time.sleep(10.5) first_peer_id = 312000101 second_peer_id = 312000201 first = HbpRepeater( proxy_port, first_peer_id, bind_host='127.0.0.2', ) second = HbpRepeater( proxy_port, second_peer_id, bind_host='127.0.0.2', ) repeaters.extend((first, second)) first.login() second.login() first.send_dmr( PacketSpec( peer_id=first_peer_id, rf_src=3120001, dst_id=3120002, slot=2, stream_id=0x11223344, call_type='unit', frame_type=HBPF_DATA_SYNC, dtype_vseq=6, ) ) try: received = second.recv(timeout=3.0) except TimeoutError as err: details = '\n'.join( 'process {}:\n{}'.format(number, process.output()) for number, process in enumerate(processes[:3]) ) raise AssertionError(details) from err self.assertEqual( received.fields['dst_id'], (3120002).to_bytes(3, 'big'), ) self.assertEqual(received.fields['call_type'], 'unit') first.send_dmr( PacketSpec( peer_id=first_peer_id, rf_src=3120001, dst_id=2350, slot=2, ) ) received = second.recv(timeout=3.0) self.assertEqual(received.fields['dst_id'], (2350).to_bytes(3, 'big')) assignments = [ int(port) for port in re.findall(r'assigned to port:(\d+)', proxy.output()) ] self.assertGreaterEqual(len(assignments), 2) self.assertIn(assignments[0], range(hbp_base, hbp_base + 2)) self.assertIn(assignments[1], range(hbp_base + 2, hbp_base + 4)) report_message, combined = _request_report_config(report_base) self.assertTrue(any(name.startswith('1:') for name in combined)) self.assertTrue(any(name.startswith('2:') for name in combined)) self.assertFalse(any(name.startswith('0:') for name in combined)) monitor_source = os.environ.get('FREEDMR_MONITOR_SOURCE') if monitor_source: connection_table = consume_config_message( Path(monitor_source), report_message, ) self.assertEqual( set(connection_table['MASTERS']), { '1:SYSTEM-0', '1:SYSTEM-1', '2:SYSTEM-0', '2:SYSTEM-1', }, ) self.assertEqual( set(connection_table['OPENBRIDGES']), { '1:OBP-HDSTACK-AGGREGATOR', '2:OBP-HDSTACK-AGGREGATOR', }, ) finally: for repeater in repeaters: repeater.close() for process in reversed(processes): if isinstance(process, FreeDmrProcess): process.__exit__(None, None, None) else: _stop_process(process) for log in logs: log.close() dependencies.cleanup() def _start_process(command, log_path): log = log_path.open('w+') process = subprocess.Popen( command, cwd=str(ROOT), stdout=log, stderr=subprocess.STDOUT, text=True, ) process.test_log_path = log_path return process, log def _stop_process(process): if process.poll() is not None: return process.terminate() try: process.wait(timeout=5) except subprocess.TimeoutExpired: process.kill() process.wait(timeout=5) def _assert_processes_running(processes): for number, process in enumerate(processes, 1): proc = process.proc if isinstance(process, FreeDmrProcess) else process if proc.poll() is not None: if isinstance(process, FreeDmrProcess): output = process.output() else: log_path = process.test_log_path try: output = log_path.read_text(encoding='utf-8') except OSError: output = ''.format(log_path) raise AssertionError( 'HDStack process {} exited during startup:\n{}'.format( number, output, ) ) def _udp_ports_available(ports): sockets = [] try: for port in ports: sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock.bind(('127.0.0.1', port)) sockets.append(sock) return True except OSError: return False finally: for sock in sockets: sock.close() def _free_tcp_port_range(count): for _attempt in range(100): first = free_udp_port() sockets = [] try: for port in range(first, first + count): sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) sock.bind(('127.0.0.1', port)) sockets.append(sock) return first except OSError: pass finally: for sock in sockets: sock.close() raise RuntimeError('could not reserve a contiguous TCP port range') def _request_report_config(port): with socket.create_connection(('127.0.0.1', port), timeout=3.0) as report: report.sendall(b'1:' + REPORT_OPCODES['CONFIG_REQ'] + b',') deadline = time.monotonic() + 3.0 while time.monotonic() < deadline: message = _receive_netstring(report) if message[:1] != REPORT_OPCODES['CONFIG_SND']: continue value = pickle.loads(message[1:]) if any(name.startswith('1:') for name in value) and any( name.startswith('2:') for name in value ): return message, value raise AssertionError('report MUX did not provide both worker configurations') def _receive_netstring(sock): length = bytearray() while True: value = sock.recv(1) if value == b':': break if not value: raise AssertionError('report MUX closed during netstring length') length.extend(value) payload_length = int(length.decode('ascii')) payload = b'' while len(payload) < payload_length: payload += sock.recv(payload_length - len(payload)) if sock.recv(1) != b',': raise AssertionError('invalid report MUX netstring terminator') return payload if __name__ == '__main__': unittest.main()