From 9fa3893a9e565f74709e1a8e1d81175b1fb53d90 Mon Sep 17 00:00:00 2001 From: Simon Date: Fri, 10 Jul 2026 13:00:14 +0000 Subject: [PATCH] Harden hotspot proxy and add UDP coverage --- docs/testing.md | 8 + hotspot_proxy_v2.py | 72 ++++- tests/harness/udp_blackbox.py | 170 +++++++++- tests/test_auxiliary_tools.py | 501 ++++++++++++++++++++++++++++- tests/test_udp_blackbox_harness.py | 57 ++++ 5 files changed, 774 insertions(+), 34 deletions(-) diff --git a/docs/testing.md b/docs/testing.md index 2f04a8a..bd72919 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -19,6 +19,14 @@ If FreeDMR runtime dependencies are not installed, tests that import The deterministic suite includes static TG routing and packet rewrite coverage. It verifies cross-slot TS1-to-TS2 routing changes only the expected slot bit while preserving packet identity fields and bytes outside that header bit. +The UDP black-box suite also includes hotspot proxy scenarios that start +`hotspot_proxy_v2.py` in front of generated HBP masters and verify proxied +login, ping, routed `DMRD` traffic, multiple destination-port allocation and +port-exhaustion behaviour. That path needs the normal runtime requirements, +including `Pyro5`. Auxiliary proxy tests cover malformed packets, `PRBL`, +RPTL flood blacklisting, source-port migration, timeout reaping, full default +100-port allocation and the existing assumption that packets from the configured +master IP are master-side traffic. API controller coverage verifies the experimental HTTP/JSON API performs only small in-memory control-plane operations, returns a clear no-options response, diff --git a/hotspot_proxy_v2.py b/hotspot_proxy_v2.py index c72e076..0fbf1bd 100644 --- a/hotspot_proxy_v2.py +++ b/hotspot_proxy_v2.py @@ -35,6 +35,11 @@ __license__ = 'GNU GPLv3' __maintainer__ = 'Simon Adlem G7RZU' __email__ = 'simon@gb7fr.org.uk' +DEFAULT_LISTEN_PORT = 62031 +DEFAULT_DESTPORT_START = 54000 +DEFAULT_DESTPORT_COUNT = 100 +DEFAULT_DESTPORT_END = DEFAULT_DESTPORT_START + DEFAULT_DESTPORT_COUNT - 1 + def bool_from_env(value): return str(value).strip().lower() in ('1', 'true', 'yes', 'on') @@ -101,12 +106,21 @@ class Proxy(DatagramProtocol): self.IPBlackList = IPBlackList self.destPortStart = DestportStart self.destPortEnd = DestPortEnd - self.numPorts = DestPortEnd - DestportStart + self.numPorts = DestPortEnd - DestportStart + 1 self.privHelper = privHelper self.rptlTrack = rptlTrack + def packet_too_short(self, data, minimum, command): + if len(data) >= minimum: + return False + if self.debug: + print('(PROXY)Ignoring short {} packet: length {}, minimum {}'.format(command.decode('ascii', errors='ignore'), len(data), minimum)) + return True + def reaper(self,_peer_id): + if _peer_id not in self.peerTrack: + return if self.debug: print("dead",_peer_id) if self.clientinfo and _peer_id != b'\xff\xff\xff\xff': @@ -163,32 +177,46 @@ class Proxy(DatagramProtocol): _command = data[:4] if _command == PRBL: + if self.packet_too_short(data, 9, PRBL): + return _peer_id = data[4:8] - _bltime = data[8:].decode('UTF-8') - _bltime = float(_bltime) - try: - self.IPBlackList[self.peerTrack[_peer_id]['shost']] = _bltime - except KeyError: + try: + _bltime = float(data[8:].decode('UTF-8')) + except (UnicodeDecodeError, ValueError): + return + _peer = self.peerTrack.get(_peer_id) + if not _peer: return + self.IPBlackList[_peer['shost']] = _bltime if self.clientinfo: - print('(PROXY)Add to blacklist: host {}. Expire time {}'.format(self.peerTrack[_peer_id]['shost'],_bltime)) + print('(PROXY)Add to blacklist: host {}. Expire time {}'.format(_peer['shost'],_bltime)) if self.privHelper: - print('(PROXY)Ask priv_helper to add to iptables: host {}, port {}.'.format(self.peerTrack[_peer_id]['shost'],self.ListenPort)) - reactor.callInThread(self.privHelper.addBL,self.ListenPort,self.peerTrack[_peer_id]['shost']) + print('(PROXY)Ask priv_helper to add to iptables: host {}, port {}.'.format(_peer['shost'],self.ListenPort)) + reactor.callInThread(self.privHelper.addBL,self.ListenPort,_peer['shost']) return if _command == DMRD: + if self.packet_too_short(data, 53, DMRD): + return _peer_id = data[11:15] elif _command == RPTA: + if self.packet_too_short(data, 10, RPTA): + return if data[6:10] in self.peerTrack: _peer_id = data[6:10] else: - _peer_id = self.connTrack[port] + _peer_id = self.connTrack.get(port, False) elif _command == MSTN: + if self.packet_too_short(data, 10, MSTN): + return _peer_id = data[6:10] elif _command == MSTP: + if self.packet_too_short(data, 11, MSTP): + return _peer_id = data[7:11] elif _command == MSTC: + if self.packet_too_short(data, 9, MSTC): + return _peer_id = data[5:9] if self.debug: @@ -207,10 +235,16 @@ class Proxy(DatagramProtocol): _command = data[:4] if _command == DMRD: # DMRData -- encapsulated DMR data frame + if self.packet_too_short(data, 53, DMRD): + return _peer_id = data[11:15] elif _command == DMRA: # DMRAlias -- Talker Alias information + if self.packet_too_short(data, 8, DMRA): + return _peer_id = data[4:8] elif _command == RPTL: # RPTLogin -- a repeater wants to login + if self.packet_too_short(data, 8, RPTL): + return _peer_id = data[4:8] #if we have seen more than 20 RPTL packets from this IP since the RPTL tracking table was reset (every 60 secs) @@ -234,15 +268,25 @@ class Proxy(DatagramProtocol): return elif _command == RPTK: # Repeater has answered our login challenge + if self.packet_too_short(data, 8, RPTK): + return _peer_id = data[4:8] elif _command == RPTC: # Repeater is sending it's configuraiton OR disconnecting if data[:5] == RPTCL: # Disconnect command + if self.packet_too_short(data, 9, RPTCL): + return _peer_id = data[5:9] else: + if self.packet_too_short(data, 8, RPTC): + return _peer_id = data[4:8] # Configure Command elif _command == RPTO: # options + if self.packet_too_short(data, 8, RPTO): + return _peer_id = data[4:8] elif _command == RPTP: # RPTPing -- peer is pinging us + if self.packet_too_short(data, 11, RPTP): + return _peer_id = data[7:11] else: return @@ -340,11 +384,11 @@ if __name__ == '__main__': #*** CONFIG HERE *** Master = "127.0.0.1" - ListenPort = 62031 + ListenPort = DEFAULT_LISTEN_PORT #'' = all IPv4, '::' = all IPv4 and IPv6 (Dual Stack) ListenIP = '' - DestportStart = 54000 - DestPortEnd = 54100 + DestportStart = DEFAULT_DESTPORT_START + DestPortEnd = DEFAULT_DESTPORT_END Timeout = 30 Stats = False Debug = False @@ -417,7 +461,7 @@ if __name__ == '__main__': if CONNTRACK[port]: count = count+1 - totalPorts = DestPortEnd - DestportStart + totalPorts = DestPortEnd - DestportStart + 1 freePorts = totalPorts - count print("{} ports out of {} in use ({} free)".format(count,totalPorts,freePorts)) diff --git a/tests/harness/udp_blackbox.py b/tests/harness/udp_blackbox.py index 6d03f37..3d788df 100644 --- a/tests/harness/udp_blackbox.py +++ b/tests/harness/udp_blackbox.py @@ -48,6 +48,7 @@ BCVE = b"BCVE" FBP_VERSION = 5 FBP_PASSPHRASE = b"test-passphrase".ljust(20, b"\x00")[:20] REQUIRED_RUNTIME_MODULES = ("bitarray", "twisted", "setproctitle") +PROXY_RUNTIME_MODULES = ("Pyro5",) @dataclass(frozen=True) @@ -246,12 +247,18 @@ def _venv_python(venv_dir: Path) -> Path: return venv_dir / "bin" / "python" -def python_has_runtime_deps(python_executable: str | Path) -> bool: - return _runtime_import_check(python_executable).returncode == 0 +def python_has_runtime_deps( + python_executable: str | Path, + modules: tuple[str, ...] = REQUIRED_RUNTIME_MODULES, +) -> bool: + return _runtime_import_check(python_executable, modules).returncode == 0 -def _runtime_import_check(python_executable: str | Path) -> subprocess.CompletedProcess: - imports = "; ".join(f"import {module}" for module in REQUIRED_RUNTIME_MODULES) +def _runtime_import_check( + python_executable: str | Path, + modules: tuple[str, ...] = REQUIRED_RUNTIME_MODULES, +) -> subprocess.CompletedProcess: + imports = "; ".join(f"import {module}" for module in modules) return subprocess.run( [str(python_executable), "-c", imports], stdout=subprocess.PIPE, @@ -263,8 +270,13 @@ def _runtime_import_check(python_executable: str | Path) -> subprocess.Completed class DependencySandbox: """Resolve or create a Python runtime that can start bridge_master.py.""" - def __init__(self, repo_root: Path) -> None: + def __init__( + self, + repo_root: Path, + extra_runtime_modules: tuple[str, ...] = (), + ) -> None: self.repo_root = repo_root + self.runtime_modules = REQUIRED_RUNTIME_MODULES + extra_runtime_modules self._tempdir: tempfile.TemporaryDirectory | None = None def cleanup(self) -> None: @@ -275,13 +287,13 @@ class DependencySandbox: def resolve_python(self) -> str: explicit_python = os.environ.get("FREEDMR_UDP_PYTHON") if explicit_python: - if not python_has_runtime_deps(explicit_python): + if not python_has_runtime_deps(explicit_python, self.runtime_modules): raise unittest.SkipTest( "FREEDMR_UDP_PYTHON does not have FreeDMR runtime dependencies" ) return explicit_python - if python_has_runtime_deps(sys.executable): + if python_has_runtime_deps(sys.executable, self.runtime_modules): return sys.executable if os.environ.get("FREEDMR_UDP_BOOTSTRAP_VENV") != "1": @@ -304,13 +316,13 @@ class DependencySandbox: if not python_executable.exists(): venv.EnvBuilder(with_pip=True).create(venv_dir) - if not python_has_runtime_deps(python_executable): + if not python_has_runtime_deps(python_executable, self.runtime_modules): requirements = self.repo_root / "requirements.txt" subprocess.check_call( [str(python_executable), "-m", "pip", "install", "-r", str(requirements)] ) - import_check = _runtime_import_check(python_executable) + import_check = _runtime_import_check(python_executable, self.runtime_modules) if import_check.returncode != 0: raise RuntimeError( "test venv was created but FreeDMR dependencies still fail to import:\n" @@ -326,6 +338,28 @@ def free_udp_port() -> int: return sock.getsockname()[1] +def free_udp_port_range(count: int) -> int: + if count < 1: + raise ValueError("count must be at least 1") + + for _attempt in range(100): + start = random.randint(20000, 60000 - count) + sockets = [] + try: + for port in range(start, start + count): + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + sock.bind(("127.0.0.1", port)) + sockets.append(sock) + return start + except OSError: + pass + finally: + for sock in sockets: + sock.close() + + raise RuntimeError(f"could not reserve a contiguous UDP port range of {count}") + + def parse_udp_dmr_fields(packet: bytes) -> dict[str, object]: if packet[:4] != DMRE: return parse_dmr_fields(packet) @@ -589,6 +623,32 @@ SERVER: 127.0.0.1 PORT: 5038 NODE: 0 {''.join(systems)} + """, + encoding="utf-8", + ) + + +def write_hotspot_proxy_config( + path: Path, + *, + listen_port: int, + dest_port_start: int, + dest_port_end: int, + master: str = "127.0.0.1", +) -> None: + path.write_text( + f"""[PROXY] +Master: {master} +ListenPort: {listen_port} +ListenIP: 127.0.0.1 +DestportStart: {dest_port_start} +DestPortEnd: {dest_port_end} +Timeout: 30 +Stats: False +Debug: False +ClientInfo: False +BlackList: [] +IPBlackList: {{}} """, encoding="utf-8", ) @@ -606,12 +666,19 @@ class UdpCapture: class HbpRepeater: - def __init__(self, master_port: int, radio_id: int, timeout: float = 2.0) -> None: - self.master = ("127.0.0.1", master_port) + def __init__( + self, + master_port: int, + radio_id: int, + timeout: float = 2.0, + bind_host: str = "127.0.0.1", + master_host: str = "127.0.0.1", + ) -> None: + self.master = (master_host, master_port) self.radio_id = bytes_4(radio_id) self.timeout = timeout self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - self.sock.bind(("127.0.0.1", 0)) + self.sock.bind((bind_host, 0)) self.sock.settimeout(timeout) self.captures: list[UdpCapture] = [] @@ -676,7 +743,7 @@ class HbpRepeater: raise AssertionError(f"expected RPTACK config, got {configured.packet!r}") def ping(self) -> None: - self.send(RPTPING + b"\x00\x00\x00" + self.radio_id) + self.send(RPTPING + self.radio_id) pong = self.recv() if not pong.packet.startswith(MSTPONG): raise AssertionError(f"expected MSTPONG, got {pong.packet!r}") @@ -963,6 +1030,32 @@ class FreeDmrProcess: raise TimeoutError("FreeDMR did not reach startup wait window") +class HotspotProxyProcess(FreeDmrProcess): + def __enter__(self): + self.proc = subprocess.Popen( + [self.python_executable, "hotspot_proxy_v2.py", "-c", str(self.config_path)], + cwd=str(self.repo_root), + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + ) + self._reader = threading.Thread(target=self._read_output, daemon=True) + self._reader.start() + return self + + def wait_for_start(self, timeout: float = 8.0) -> None: + assert self.proc is not None + deadline = time.monotonic() + timeout + ready_at = time.monotonic() + 0.5 + while time.monotonic() < deadline: + if self.proc.poll() is not None: + raise RuntimeError("hotspot_proxy_v2 exited before startup:\n" + self.output()) + if time.monotonic() >= ready_at: + return + time.sleep(0.05) + raise TimeoutError("hotspot_proxy_v2 did not reach startup wait window") + + class UdpBlackBoxScenario: def __init__( self, @@ -980,17 +1073,42 @@ class UdpBlackBoxScenario: master_extra_config: str = "", fbp_systems: dict[str, int] | None = None, fbp_proto_versions: dict[str, int] | None = None, + enable_proxy: bool = False, + proxy_dest_count: int = 1, ) -> None: self.repo_root = repo_root or Path(__file__).resolve().parents[2] self.tempdir = tempfile.TemporaryDirectory(prefix="freedmr-udp-test-") self.config_path = Path(self.tempdir.name) / "freedmr-test.cfg" - self.system_ports = {"MASTER-A": free_udp_port(), "MASTER-B": free_udp_port()} + self.proxy_config_path = Path(self.tempdir.name) / "hotspot-proxy-test.cfg" + self.enable_proxy = enable_proxy + self.proxy_dest_ports = [] + self.proxy_port = free_udp_port() if enable_proxy else None + if enable_proxy: + proxy_dest_start = free_udp_port_range(proxy_dest_count) + self.proxy_dest_ports = list( + range(proxy_dest_start, proxy_dest_start + proxy_dest_count) + ) + master_b_port = free_udp_port() + while master_b_port in self.proxy_dest_ports: + master_b_port = free_udp_port() + self.system_ports = { + "MASTER-A": self.proxy_dest_ports[0], + "MASTER-B": master_b_port, + } + for index, port in enumerate(self.proxy_dest_ports[1:], start=1): + self.system_ports[f"PROXY-{index}"] = port + else: + self.system_ports = {"MASTER-A": free_udp_port(), "MASTER-B": free_udp_port()} self.fbp_network_ids = fbp_systems or {} self.fbp_proto_versions = fbp_proto_versions or {} self.fbp_system_ports = {name: free_udp_port() for name in self.fbp_network_ids} self.fbp_peers: dict[str, FbpPeer] = {} self.process: FreeDmrProcess | None = None - self.deps = DependencySandbox(self.repo_root) + self.proxy_process: HotspotProxyProcess | None = None + self.deps = DependencySandbox( + self.repo_root, + extra_runtime_modules=PROXY_RUNTIME_MODULES if enable_proxy else (), + ) self.config_options = { "global_use_acl": global_use_acl, "global_sub_acl": global_sub_acl, @@ -1024,9 +1142,26 @@ class UdpBlackBoxScenario: self.process = FreeDmrProcess(self.repo_root, self.config_path, python_executable) self.process.__enter__() self.process.wait_for_start() + if self.enable_proxy: + assert self.proxy_port is not None + write_hotspot_proxy_config( + self.proxy_config_path, + listen_port=self.proxy_port, + dest_port_start=self.proxy_dest_ports[0], + dest_port_end=self.proxy_dest_ports[-1], + ) + self.proxy_process = HotspotProxyProcess( + self.repo_root, + self.proxy_config_path, + python_executable, + ) + self.proxy_process.__enter__() + self.proxy_process.wait_for_start() return self def __exit__(self, exc_type, exc, tb) -> None: + if self.proxy_process is not None: + self.proxy_process.__exit__(exc_type, exc, tb) if self.process is not None: self.process.__exit__(exc_type, exc, tb) for peer in self.fbp_peers.values(): @@ -1037,5 +1172,10 @@ class UdpBlackBoxScenario: def repeater(self, system_name: str, radio_id: int) -> HbpRepeater: return HbpRepeater(self.system_ports[system_name], radio_id) + def proxied_repeater(self, radio_id: int) -> HbpRepeater: + if self.proxy_port is None: + raise RuntimeError("UdpBlackBoxScenario was not created with enable_proxy=True") + return HbpRepeater(self.proxy_port, radio_id, bind_host="127.0.0.2") + def fbp_peer(self, system_name: str) -> FbpPeer: return self.fbp_peers[system_name] diff --git a/tests/test_auxiliary_tools.py b/tests/test_auxiliary_tools.py index 210e2fa..6411dc1 100644 --- a/tests/test_auxiliary_tools.py +++ b/tests/test_auxiliary_tools.py @@ -77,11 +77,8 @@ class AuxiliaryToolTests(unittest.TestCase): self.assertTrue(fake_db.cursor_obj.closed) def test_proxy_environment_bool_parser(self): - saved_modules = self._install_proxy_stubs() + hotspot_proxy_v2, saved_modules = self._import_proxy_module() try: - import hotspot_proxy_v2 - hotspot_proxy_v2 = importlib.reload(hotspot_proxy_v2) - self.assertTrue(hotspot_proxy_v2.bool_from_env("1")) self.assertTrue(hotspot_proxy_v2.bool_from_env("true")) self.assertTrue(hotspot_proxy_v2.bool_from_env("yes")) @@ -91,6 +88,422 @@ class AuxiliaryToolTests(unittest.TestCase): finally: self._restore_modules(saved_modules) + def test_proxy_default_destination_range_matches_generated_masters(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + default_ports = range( + hotspot_proxy_v2.DEFAULT_DESTPORT_START, + hotspot_proxy_v2.DEFAULT_DESTPORT_END + 1, + ) + + self.assertEqual(hotspot_proxy_v2.DEFAULT_DESTPORT_COUNT, 100) + self.assertEqual(len(default_ports), hotspot_proxy_v2.DEFAULT_DESTPORT_COUNT) + finally: + self._restore_modules(saved_modules) + + def test_proxy_routes_login_and_master_dmrd(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + fake_reactor = _FakeReactor() + hotspot_proxy_v2.reactor = fake_reactor + transport = _FakeTransport() + peer_id = b"\x00\x00\x03\xe9" + conn_track = {54000: False} + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"RPTL" + peer_id, ("198.51.100.10", 40000)) + proxy.datagramReceived(_dmrd_packet(peer_id), ("127.0.0.1", 54000)) + + self.assertEqual(conn_track[54000], peer_id) + self.assertEqual(peer_track[peer_id]["shost"], "198.51.100.10") + self.assertEqual(transport.writes[0], (b"RPTL" + peer_id, ("127.0.0.1", 54000))) + self.assertEqual(transport.writes[-1], (_dmrd_packet(peer_id), ("198.51.100.10", 40000))) + finally: + self._restore_modules(saved_modules) + + def test_proxy_ignores_short_client_dmrd(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + conn_track = {54000: False} + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"DMRD", ("198.51.100.10", 40000)) + + self.assertEqual(conn_track[54000], False) + self.assertEqual(peer_track, {}) + self.assertEqual(transport.writes, []) + finally: + self._restore_modules(saved_modules) + + def test_proxy_ignores_unknown_master_rptack_port(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: False}, + {}, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"RPTACK" + b"\x00\x00\x03\xe9", ("127.0.0.1", 59999)) + + self.assertEqual(transport.writes, []) + finally: + self._restore_modules(saved_modules) + + def test_proxy_ignores_new_client_when_all_destination_ports_are_in_use(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + existing_peer = b"\x00\x00\x03\xe9" + new_peer = b"\x00\x00\x03\xea" + conn_track = {54000: existing_peer} + peer_track = { + existing_peer: { + "dport": 54000, + "sport": 40000, + "shost": "198.51.100.10", + "timer": _FakeTimer(), + } + } + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"RPTL" + new_peer, ("198.51.100.11", 40001)) + + self.assertEqual(conn_track, {54000: existing_peer}) + self.assertNotIn(new_peer, peer_track) + self.assertEqual(transport.writes, []) + finally: + self._restore_modules(saved_modules) + + def test_proxy_allocates_each_default_destination_port_once_before_exhaustion(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + conn_track = { + port: False + for port in range( + hotspot_proxy_v2.DEFAULT_DESTPORT_START, + hotspot_proxy_v2.DEFAULT_DESTPORT_END + 1, + ) + } + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + hotspot_proxy_v2.DEFAULT_LISTEN_PORT, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + hotspot_proxy_v2.DEFAULT_DESTPORT_START, + hotspot_proxy_v2.DEFAULT_DESTPORT_END, + None, + {}, + ) + proxy.transport = transport + + peer_ids = [] + for index in range(hotspot_proxy_v2.DEFAULT_DESTPORT_COUNT): + peer_id = (1000 + index).to_bytes(4, "big") + peer_ids.append(peer_id) + proxy.datagramReceived( + b"RPTL" + peer_id, + ("198.51.100.{}".format(index + 1), 40000 + index), + ) + + write_count_before_exhaustion = len(transport.writes) + proxy.datagramReceived(b"RPTL" + (9999).to_bytes(4, "big"), ("198.51.100.200", 40100)) + + self.assertEqual(len(peer_track), hotspot_proxy_v2.DEFAULT_DESTPORT_COUNT) + self.assertEqual(set(conn_track.values()), set(peer_ids)) + self.assertNotIn(False, conn_track.values()) + self.assertEqual(len(transport.writes), write_count_before_exhaustion) + finally: + self._restore_modules(saved_modules) + + def test_proxy_updates_client_source_port_for_existing_peer(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + peer_id = b"\x00\x00\x03\xe9" + conn_track = {54000: False} + peer_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"RPTL" + peer_id, ("198.51.100.10", 40000)) + proxy.datagramReceived(b"RPTPING" + peer_id, ("198.51.100.10", 40001)) + proxy.datagramReceived(b"MSTPONG" + peer_id, ("127.0.0.1", 54000)) + + self.assertEqual(peer_track[peer_id]["sport"], 40001) + self.assertEqual(transport.writes[-1], (b"MSTPONG" + peer_id, ("198.51.100.10", 40001))) + finally: + self._restore_modules(saved_modules) + + def test_proxy_reaper_releases_port_and_is_idempotent(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + transport = _FakeTransport() + peer_id = b"\x00\x00\x03\xe9" + conn_track = {54000: peer_id} + peer_track = { + peer_id: { + "dport": 54000, + "sport": 40000, + "shost": "198.51.100.10", + "timer": _FakeTimer(), + } + } + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + conn_track, + peer_track, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.reaper(peer_id) + write_count = len(transport.writes) + proxy.reaper(peer_id) + + self.assertEqual(conn_track[54000], False) + self.assertNotIn(peer_id, peer_track) + self.assertEqual(transport.writes[0], (b"RPTCL" + peer_id, ("127.0.0.1", 54000))) + self.assertEqual(transport.writes[1:], [(b"MSTCL", ("198.51.100.10", 40000))] * 3) + self.assertEqual(len(transport.writes), write_count) + finally: + self._restore_modules(saved_modules) + + def test_proxy_blacklists_known_peer_from_proxy_control_packet(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + peer_id = b"\x00\x00\x03\xe9" + ip_blacklist = {} + peer_track = { + peer_id: { + "dport": 54000, + "sport": 40000, + "shost": "198.51.100.10", + "timer": _FakeTimer(), + } + } + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: peer_id}, + peer_track, + [], + ip_blacklist, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"PRBL" + peer_id + b"12345.0", ("127.0.0.1", 54000)) + + self.assertEqual(ip_blacklist["198.51.100.10"], 12345.0) + finally: + self._restore_modules(saved_modules) + + def test_proxy_rptl_flood_blacklists_source_ip(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + peer_id = b"\x00\x00\x03\xe9" + ip_blacklist = {} + rptl_track = {} + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: False}, + {}, + [], + ip_blacklist, + 30, + False, + False, + 54000, + 54000, + None, + rptl_track, + ) + proxy.transport = transport + + for index in range(21): + proxy.datagramReceived(b"RPTL" + peer_id, ("198.51.100.10", 40000 + index)) + write_count = len(transport.writes) + proxy.datagramReceived(b"RPTPING" + peer_id, ("198.51.100.10", 40022)) + + self.assertIn("198.51.100.10", ip_blacklist) + self.assertNotIn("198.51.100.10", rptl_track) + self.assertEqual(len(transport.writes), write_count) + finally: + self._restore_modules(saved_modules) + + def test_proxy_treats_packets_from_master_ip_as_master_side(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: False}, + {}, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"RPTL" + b"\x00\x00\x03\xe9", ("127.0.0.1", 40000)) + + self.assertEqual(proxy.peerTrack, {}) + self.assertEqual(transport.writes, []) + finally: + self._restore_modules(saved_modules) + + def test_proxy_ignores_malformed_proxy_blacklist_packet(self): + hotspot_proxy_v2, saved_modules = self._import_proxy_module() + try: + hotspot_proxy_v2.reactor = _FakeReactor() + transport = _FakeTransport() + proxy = hotspot_proxy_v2.Proxy( + "127.0.0.1", + 62031, + {54000: False}, + {}, + [], + {}, + 30, + False, + False, + 54000, + 54000, + None, + {}, + ) + proxy.transport = transport + + proxy.datagramReceived(b"PRBL" + b"\x00\x00\x03\xe9" + b"not-a-float", ("127.0.0.1", 54000)) + + self.assertEqual(transport.writes, []) + finally: + self._restore_modules(saved_modules) + + def _import_proxy_module(self): + saved_modules = self._install_proxy_stubs() + try: + import hotspot_proxy_v2 + return importlib.reload(hotspot_proxy_v2), saved_modules + except Exception: + self._restore_modules(saved_modules) + raise + def _install_mysql_stub(self): mysql_module = types.ModuleType("mysql") connector_module = types.ModuleType("mysql.connector") @@ -108,7 +521,15 @@ class AuxiliaryToolTests(unittest.TestCase): sys.modules["mysql.connector"] = connector_module def _install_proxy_stubs(self): - stubbed = ["Pyro5", "Pyro5.api"] + stubbed = [ + "Pyro5", + "Pyro5.api", + "setproctitle", + "twisted", + "twisted.internet", + "twisted.internet.protocol", + "twisted.internet.task", + ] saved_modules = {name: sys.modules.get(name) for name in stubbed + ["hotspot_proxy_v2"]} pyro5_module = types.ModuleType("Pyro5") @@ -117,6 +538,26 @@ class AuxiliaryToolTests(unittest.TestCase): pyro5_module.api = pyro5_api_module sys.modules["Pyro5"] = pyro5_module sys.modules["Pyro5.api"] = pyro5_api_module + + setproctitle_module = types.ModuleType("setproctitle") + setproctitle_module.setproctitle = lambda title: None + sys.modules["setproctitle"] = setproctitle_module + + twisted_module = types.ModuleType("twisted") + twisted_internet_module = types.ModuleType("twisted.internet") + twisted_protocol_module = types.ModuleType("twisted.internet.protocol") + twisted_task_module = types.ModuleType("twisted.internet.task") + twisted_protocol_module.DatagramProtocol = object + twisted_task_module.LoopingCall = object + twisted_internet_module.protocol = twisted_protocol_module + twisted_internet_module.reactor = _FakeReactor() + twisted_internet_module.task = twisted_task_module + twisted_module.internet = twisted_internet_module + sys.modules["twisted"] = twisted_module + sys.modules["twisted.internet"] = twisted_internet_module + sys.modules["twisted.internet.protocol"] = twisted_protocol_module + sys.modules["twisted.internet.task"] = twisted_task_module + sys.modules.pop("hotspot_proxy_v2", None) return saved_modules @@ -155,5 +596,55 @@ class _FakeDB: self.committed = True +class _FakeTimer: + def __init__(self): + self.resets = [] + self.cancelled = False + + def reset(self, value=None): + self.resets.append(value) + + def cancel(self): + self.cancelled = True + + +class _FakeReactor: + def __init__(self): + self.timers = [] + self.thread_calls = [] + + def callLater(self, timeout, func, *args): + timer = _FakeTimer() + self.timers.append((timeout, func, args, timer)) + return timer + + def callInThread(self, func, *args): + self.thread_calls.append((func, args)) + func(*args) + + +class _FakeTransport: + def __init__(self): + self.writes = [] + + def write(self, data, addr): + self.writes.append((data, addr)) + + +def _dmrd_packet(peer_id): + return b"".join( + [ + b"DMRD", + b"\x01", + b"\x00\x00\x01", + b"\x00\x00\x02", + peer_id, + b"\x80", + b"\x01\x02\x03\x04", + b"\x55" * 33, + ] + ) + + if __name__ == "__main__": unittest.main() diff --git a/tests/test_udp_blackbox_harness.py b/tests/test_udp_blackbox_harness.py index 48cb0e8..01b65db 100644 --- a/tests/test_udp_blackbox_harness.py +++ b/tests/test_udp_blackbox_harness.py @@ -79,6 +79,63 @@ class UdpBlackBoxHarnessTest(unittest.TestCase): self.assertEqual(captured.fields["slot"], 2) self.assertEqual(captured.fields["stream_id"], bytes_4(0x01020304)) + def test_hotspot_proxy_routes_hbp_repeater_to_master(self): + require_udp_integration_enabled() + + with UdpBlackBoxScenario(enable_proxy=True) as scenario: + proxied = scenario.proxied_repeater(1001) + master_b = scenario.repeater("MASTER-B", 1002) + try: + proxied.login() + proxied.ping() + master_b.login() + + proxied.send_dmr(PacketSpec(peer_id=1001, rf_src=3120001, dst_id=91, slot=2)) + captured = master_b.recv(timeout=2.0) + finally: + proxied.close() + master_b.close() + + self.assertEqual(captured.packet[:4], b"DMRD") + self.assertEqual(captured.fields["rf_src"], bytes_3(3120001)) + self.assertEqual(captured.fields["dst_id"], bytes_3(91)) + self.assertEqual(captured.fields["slot"], 2) + self.assertEqual(captured.fields["stream_id"], bytes_4(0x01020304)) + + def test_hotspot_proxy_allocates_multiple_ports_and_refuses_exhausted_client(self): + require_udp_integration_enabled() + + with UdpBlackBoxScenario(enable_proxy=True, proxy_dest_count=2) as scenario: + proxied_a = scenario.proxied_repeater(1001) + proxied_b = scenario.proxied_repeater(1002) + refused = scenario.proxied_repeater(1003) + master_b = scenario.repeater("MASTER-B", 1004) + try: + proxied_a.login() + proxied_a.ping() + proxied_b.login() + proxied_b.ping() + master_b.login() + + proxied_a.send_dmr( + PacketSpec(peer_id=1001, rf_src=3120001, dst_id=91, slot=2) + ) + captured = master_b.recv(timeout=2.0) + + refused.send(b"RPTL" + bytes_4(1003)) + with self.assertRaises((TimeoutError, socket.timeout)): + refused.recv(timeout=0.5) + finally: + proxied_a.close() + proxied_b.close() + refused.close() + master_b.close() + + self.assertEqual(captured.packet[:4], b"DMRD") + self.assertEqual(captured.fields["peer_id"], bytes_4(1004)) + self.assertEqual(captured.fields["rf_src"], bytes_3(3120001)) + self.assertEqual(captured.fields["dst_id"], bytes_3(91)) + def test_startup_accepts_slot_and_route_specific_timer_config(self): require_udp_integration_enabled()