From 86983597b2f23f4dbcc3b0bc7dc1b9b11cd7dc28 Mon Sep 17 00:00:00 2001 From: ce5rpy <169016246+ce5rpy@users.noreply.github.com> Date: Sat, 4 Jul 2026 15:49:19 -0400 Subject: [PATCH] fix: echo TG 9990 isolation for multi-peer inject proxy (#32) Echo (TG 9990-9999) must be point-to-point: only the originating hotspot receives the playback. Two bugs in the inject-only proxy broke this: Data plane (udp_hbp.py): _peer_should_receive_dmrd had a fuzzy-match fallback (peer_id // 100 == rf_src) for special TGs that leaked the first VHEAD to every hotspot sharing the user's DMR-id prefix. Removed the fallback so echo delivers only to the exact RX_PEER. Report plane (monitor_topology.py): _echo_tx_target_peer resolved the monitor chip peer via fuzzy rf_src matching, which picked the wrong hotspot when a user has several radios sharing a base id (e.g. rf_src 7140023 matched peer 714002301 instead of the real originator 714000103). Now resolves from STATUS[slot].RX_PEER when available. --- .../application/report/monitor_topology.py | 30 ++++- .../twisted_adapters/udp_hbp.py | 7 +- tests/application/test_monitor_topology.py | 58 +++++++++ .../test_echo_special_tg_downlink.py | 123 ++++++++++++++++++ 4 files changed, 211 insertions(+), 7 deletions(-) create mode 100644 tests/infrastructure/test_echo_special_tg_downlink.py diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py index 007a39b..590ea8e 100644 --- a/src/adn_server/application/report/monitor_topology.py +++ b/src/adn_server/application/report/monitor_topology.py @@ -267,12 +267,32 @@ def _peers_receiving_tgid( return out -def _echo_tx_target_peer(parts: list[str], peers: dict[Any, Any]) -> bytes | None: - """Echo/service downlink TX: field 5 may be 9990 or hotspot id; field 8 is 9990–9999.""" +def _echo_tx_target_peer( + parts: list[str], + peers: dict[Any, Any], + *, + status: dict[int, dict[str, Any]] | None = None, +) -> bytes | None: + """Echo/service downlink TX: field 5 may be 9990 or hotspot id; field 8 is 9990–9999. + + For special TGs (9990–9999) the audio is point-to-point: the downlink goes only to + the peer that originated the call. The runtime ``STATUS[slot]["RX_PEER"]`` holds that + exact peer id, so we prefer it. Fuzzy matching on ``rf_src`` (field 6) resolves the + wrong hotspot when a user has several radios sharing a DMR-id prefix (e.g. + rf_src 7140023 matches peer 714002301 via ``// 100`` instead of the real 714000103). + """ tgid_slot = _voice_event_tgid_slot(parts) if tgid_slot is not None: - tgid, _ = tgid_slot + tgid, voice_slot = tgid_slot if is_special_tg(str(tgid)): + if status is not None: + slot_st = status.get(voice_slot, {}) + rx_peer = slot_st.get("RX_PEER", b"") + if rx_peer and rx_peer != b"\x00\x00\x00\x00": + rx_b = bytes_4(int_id(rx_peer)) + connected = _connected_peer_keys(peers) + if rx_b in connected: + return rx_b return _peer_key_from_voice_csv(parts, peers) if len(parts) <= 5: return None @@ -350,7 +370,9 @@ def remap_inject_proxy_voice_events( if trx == "TX": tgid_slot = _voice_event_tgid_slot(parts) - echo_peer = _echo_tx_target_peer(parts, peers) + echo_peer = _echo_tx_target_peer( + parts, peers, status=downlink_ctx.status if downlink_ctx else None + ) if echo_peer is not None: slot = slot_map.get(echo_peer) if slot is not None: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 04d2e7b..c4c9be6 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -569,15 +569,16 @@ class HBPProtocol(DatagramProtocol): if rx_peer and rx_peer != b"\x00\x00\x00\x00": return bytes_4(int_id(peer_id)) == bytes_4(int_id(rx_peer)) return self._cached_connected_peer_count() == 1 - # Parrot / echo (9990–9999): not in per-hotspot OPTIONS; deliver to last RX peer on slot. + # Parrot / echo (9990–9999): point-to-point echo — deliver ONLY to the exact peer + # that originated the call (RX_PEER), never to other hotspots of the same user. + # Legacy parity: each MASTER had a single peer, so echo naturally returned only + # to the caller. Multi-peer inject-only proxy must enforce the same explicitly. if is_special_tg(str(tgid)): slot_st = self.STATUS.get(slot, {}) if int_id(slot_st.get("RX_TGID", b"\x00\x00\x00")) == tgid: rx_peer = slot_st.get("RX_PEER", b"") if rx_peer and rx_peer != b"\x00\x00\x00\x00": return bytes_4(int_id(peer_id)) == bytes_4(int_id(rx_peer)) - if len(packet) >= 8 and peer_matches_rf_source(peer_id, packet[5:8], self._peers): - return True return self._cached_connected_peer_count() == 1 if len(packet) >= 8 and peer_matches_rf_source(peer_id, packet[5:8], self._peers): ctx = self._downlink_ctx() diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py index 268800e..ae060cd 100644 --- a/tests/application/test_monitor_topology.py +++ b/tests/application/test_monitor_topology.py @@ -372,3 +372,61 @@ def test_non_proxy_systems_pass_through_unchanged() -> None: "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, } assert expand_inject_proxy_systems({"PROXY": {"TARGET_SYSTEM": "SYSTEM"}}, systems) is systems + + +def test_echo_tx_uses_rx_peer_when_rf_src_fuzzy_matches_sibling() -> None: + """Echo TX must resolve to RX_PEER from STATUS, not fuzzy-match rf_src. + + Reproduces the monitor display bug: user 7140023 has two hotspots, + 714000103 (base 7140001) and 714002301 (base 7140023). The echo playback + arrives with peer_id=9990, rf_src=7140023. Fuzzy matching (// 100) resolves + 714002301 instead of the true originator 714000103 held in STATUS.RX_PEER. + """ + from adn_server.application.routing.downlink import DownlinkContext + + peer_origin = bytes_4(714000103) # hotspot that transmitted to 9990 + peer_sibling = bytes_4(714002301) # same user, base == rf_src + peers = {peer_origin: _peer(), peer_sibling: _peer()} + config = _proxy_config(peers) + peer_slots = {peer_origin: 1, peer_sibling: 5} + status = { + 2: { + "RX_PEER": peer_origin, + "RX_TGID": bytes_4(9990)[1:], + } + } + ctx = DownlinkContext( + config=config, + system_name="SYSTEM", + sys_cfg=config["SYSTEMS"]["SYSTEM"], + peers=peers, + status=status, + connected_count=2, + ) + # Echo playback TX: peer_id=9990, rf_src=7140023 (user DMR ID), dst=9990 slot 2 + raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7140023,2,9990" + events = remap_inject_proxy_voice_events( + raw, config, config["SYSTEMS"], peer_slots, downlink_ctx=ctx, + ) + assert len(events) == 1 + parts = events[0].split(",") + # Must remap to SYSTEM-1 (peer_origin slot), NOT SYSTEM-5 (peer_sibling) + assert parts[3] == "SYSTEM-1" + assert parts[5] == "9990" # echo id preserved for TX chip display + + +def test_echo_tx_falls_back_without_downlink_ctx() -> None: + """Without STATUS (no downlink_ctx), echo TX cannot resolve when rf_src + does not fuzzy-match any connected peer — event passes through unchanged.""" + peer = bytes_4(714000103) + peers = {peer: _peer()} + config = _proxy_config(peers) + peer_slots = {peer: 1} + raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7140023,2,9990" + events = remap_inject_proxy_voice_events( + raw, config, config["SYSTEMS"], peer_slots, + ) + # rf_src 7140023 does not match peer 714000103 via // 100 or prefix; + # legacy leaves the event unchanged when the hotspot cannot be resolved. + assert len(events) == 1 + assert events[0] == raw diff --git a/tests/infrastructure/test_echo_special_tg_downlink.py b/tests/infrastructure/test_echo_special_tg_downlink.py new file mode 100644 index 0000000..987cc53 --- /dev/null +++ b/tests/infrastructure/test_echo_special_tg_downlink.py @@ -0,0 +1,123 @@ +# ADN DMR Peer Server - echo / special TG downlink filter non-regression +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +# Regression test: TG 9990 (echo) and special TGs 9990-9999 must only be +# delivered to the exact peer that originated the call, never replicated +# to other hotspots of the same user (same DMR ID base). + +from __future__ import annotations + +from adn_server.application.routing.helpers import synthetic_group_dmrd_route_packet +from adn_server.domain import bytes_3, bytes_4 +from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config +from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol + + +def _master_config() -> dict: + config = { + "GLOBAL": {"USE_ACL": False}, + "SYSTEMS": { + "TEST": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 8, + "GROUP_HANGTIME": 0, + } + }, + } + ensure_system_runtime_config(config) + return config + + +def _make_echo_packet(rf_src_base: int, tgid: int = 9990) -> bytes: + """Build a DMRD group packet for TG 9990 with a given rf_src (3-byte base).""" + pkt = synthetic_group_dmrd_route_packet(2, tgid) + return pkt[:5] + bytes_3(rf_src_base) + pkt[8:] + + +def test_echo_9990_only_to_originating_peer() -> None: + """Echo playback must deliver ONLY to the RX_PEER that originated the call.""" + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_a = bytes_4(730039101) + peer_b = bytes_4(730039210) + proto._peers = { + peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)}, + peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)}, + } + # Simulate: peer_a transmitted to TG 9990 on slot 2 (RX_PEER set by dmrd_received). + proto.STATUS[2] = { + "RX_PEER": peer_a, + "RX_TGID": bytes_3(9990), + "RX_STREAM_ID": b"\x00\x00\x00\x01", + } + rf_src_base = 7300392 # both peers share this base ID (same user) + echo_pkt = _make_echo_packet(rf_src_base, 9990) + assert proto._peer_should_receive_dmrd(peer_a, echo_pkt) + assert not proto._peer_should_receive_dmrd(peer_b, echo_pkt) + + +def test_echo_9990_no_fuzzy_match_when_no_rx_peer() -> None: + """When RX_PEER is not set, echo must NOT use fuzzy rf_src matching. + + With multiple peers connected, no peer should receive the echo if the + originating RX_PEER is unknown (single-peer fallback is the only exception). + """ + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_a = bytes_4(730039101) + peer_b = bytes_4(730039210) + proto._peers = { + peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)}, + peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)}, + } + # RX_PEER not set (stale / different TG on slot) + proto.STATUS[2] = { + "RX_PEER": b"\x00\x00\x00\x00", + "RX_TGID": bytes_3(0), + "RX_STREAM_ID": b"\x00\x00\x00\x00", + } + rf_src_base = 7300392 + echo_pkt = _make_echo_packet(rf_src_base, 9990) + assert not proto._peer_should_receive_dmrd(peer_a, echo_pkt) + assert not proto._peer_should_receive_dmrd(peer_b, echo_pkt) + + +def test_echo_9990_single_peer_fallback() -> None: + """When only one peer is connected, echo is delivered to it (legacy parity).""" + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_a = bytes_4(730039101) + proto._peers = { + peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)}, + } + proto.STATUS[2] = { + "RX_PEER": b"\x00\x00\x00\x00", + "RX_TGID": bytes_3(0), + "RX_STREAM_ID": b"\x00\x00\x00\x00", + } + rf_src_base = 7300392 + echo_pkt = _make_echo_packet(rf_src_base, 9990) + assert proto._peer_should_receive_dmrd(peer_a, echo_pkt) + + +def test_special_tg_9991_same_isolation() -> None: + """All special TGs 9990-9999 share the same point-to-point isolation.""" + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_a = bytes_4(730039101) + peer_b = bytes_4(730039210) + proto._peers = { + peer_a: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62031)}, + peer_b: {"CONNECTION": "YES", "OPTIONS": b"", "SOCKADDR": ("127.0.0.1", 62032)}, + } + proto.STATUS[2] = { + "RX_PEER": peer_b, + "RX_TGID": bytes_3(9991), + "RX_STREAM_ID": b"\x00\x00\x00\x01", + } + rf_src_base = 7300392 + pkt = _make_echo_packet(rf_src_base, 9991) + assert proto._peer_should_receive_dmrd(peer_b, pkt) + assert not proto._peer_should_receive_dmrd(peer_a, pkt)