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_hdstack_udp.py

395 lines
14 KiB

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 = '<could not read {}>'.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()

Powered by TurnKey Linux.