From 11970e0552203fa1e013121c8997f805504898b4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 24 Jul 2026 23:28:33 -0400 Subject: [PATCH] fix: target private calls to the known hotspot instead of broadcasting Use SUB_MAP's stored peer id for delivery (repeat, unit-data, and pvt_call_received), add Talker Alias support for private calls, and report the receiving hotspot to the monitor. --- .../application/report/monitor_topology.py | 33 +++++ src/adn_server/application/report/payloads.py | 26 +++- .../application/routing/hbp_forward.py | 21 ++- src/adn_server/application/routing/lc_ta.py | 14 +- .../application/routing_use_cases.py | 81 +++++++++-- .../infrastructure/bootstrap/peer_server.py | 6 + .../twisted_adapters/udp_hbp.py | 106 +++++++++++--- tests/application/test_monitor_topology.py | 59 ++++++++ tests/application/test_report_payloads.py | 57 ++++++++ tests/harness/deterministic.py | 5 + .../test_hbp_private_call_targeting.py | 134 +++++++++++++++++ .../test_hbp_repeat_private_call.py | 135 ++++++++++++++++++ tests/routing/test_unit_data_routing.py | 67 ++++++++- 13 files changed, 704 insertions(+), 40 deletions(-) create mode 100644 tests/infrastructure/test_hbp_private_call_targeting.py create mode 100644 tests/infrastructure/test_hbp_repeat_private_call.py diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py index f76e6d6..833103b 100644 --- a/src/adn_server/application/report/monitor_topology.py +++ b/src/adn_server/application/report/monitor_topology.py @@ -221,6 +221,26 @@ def _voice_event_stream_id(parts: list[str]) -> bytes | None: return None +def _private_voice_tx_dest_peer_id(parts: list[str]) -> int | None: + """Destination hotspot peer id trailing a PRIVATE VOICE TX event, if present. + + routing_use_cases._pvt_call_received appends this (SUB_MAP's known peer for the + destination) after the legacy 9 START / 10 END fields. Private call destinations + are subscriber ids, never a real static TG -- _peers_receiving_tgid's group-voice + subscription fan-out can never resolve one, so this field is the only way to know + which single hotspot slot the event belongs to. + """ + if len(parts) < 2 or parts[0].strip() != "PRIVATE VOICE" or parts[2].strip() != "TX": + return None + base_len = 10 if parts[1].strip() == "END" else 9 + if len(parts) <= base_len: + return None + try: + return int(parts[-1].strip()) + except ValueError: + return None + + def _peers_receiving_tgid( connected: list[tuple[Any, dict[str, Any]]], *, @@ -374,6 +394,19 @@ def remap_inject_proxy_voice_events( stream_id = _voice_event_stream_id(parts) if trx == "TX": + dest_peer_id = _private_voice_tx_dest_peer_id(parts) + if dest_peer_id is not None: + peer_key = bytes_4(dest_peer_id) + if peer_key not in _connected_peer_keys(peers): + return [] + slot = slot_map.get(peer_key) + if slot is None: + return [] + return [ + _remap_voice_event_to_slot( + parts, target=target, slot=slot, peer_key=peer_key + ) + ] tgid_slot = _voice_event_tgid_slot(parts) echo_peer = _echo_tx_target_peer( parts, peers, status=downlink_ctx.status if downlink_ctx else None diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index 5611df6..4de9252 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -52,6 +52,17 @@ _CSV_FAMILIES = { "GROUP VOICE": "GROUP", "PRIVATE VOICE": "PRIVATE", "UNIT DATA": "UNIT", + # routing_use_cases._dtype_labels emits one of these instead of the generic + # "UNIT DATA" label for dtype_vseq in (3, 6, 7, 8) — private-call CSBK setup + # signaling and SMS/GPS (ARS/LRRP) header + VCSBK blocks. Without these, + # parse_bridge_event_csv() returns None and the event is silently dropped + # ("voice_event not emitted (unmapped CSV)"), so a private call whose CSBK + # handshake never escalates to voice (e.g. destination unreachable) never + # shows up in the monitor at all. + "UNIT CSBK": "UNIT", + "UNIT DATA HEADER": "UNIT", + "UNIT VCSBK 1/2 DATA BLOCK": "UNIT", + "UNIT VCSBK 3/4 DATA BLOCK": "UNIT", } @@ -515,7 +526,20 @@ def parse_bridge_event_csv(event: str, *, ts: float | None = None) -> dict[str, voice["duration_s"] = None elif phase != "END": voice["duration_s"] = None - voice["is_announcement"] = len(parts) > 9 and parts[-1] == "1" + if call_family == "PRIVATE" and direction == "TX": + # routing_use_cases._pvt_call_received appends the destination hotspot's own + # peer id (from SUB_MAP) after the legacy fields -- there's no static TG a + # private call's subscriber-id destination could ever match, so this is the + # only way the monitor can know which single hotspot is receiving. + voice["is_announcement"] = False + base_len = 10 if phase == "END" else 9 + if len(parts) > base_len: + try: + voice["dest_peer_id"] = int(parts[-1]) + except ValueError: + pass + else: + voice["is_announcement"] = len(parts) > 9 and parts[-1] == "1" return voice diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 6d36e12..049afb3 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -353,11 +353,28 @@ class HbpForwardMixin: rf_src: bytes, stream_id: bytes, peer_id: bytes, + d_peer_id: bytes | None = None, ) -> None: - """Legacy sendDataToHBP: forward a unit-data packet to an HBP (MASTER/PEER) target.""" + """Legacy sendDataToHBP: forward a unit-data packet to an HBP (MASTER/PEER) target. + + ``d_peer_id`` (the exact hotspot the destination was last heard on, from SUB_MAP + or the 6/7-digit peer-ID match) lets this go directly to that one peer via + ``send_peer`` instead of ``send_system``/``send_peers`` broadcasting the private + DATA to every other hotspot connected to ``d_system`` — the same cross-talk this + server already avoids for private voice REPEAT (``_pvt_repeat_targets``). Falls + back to the system-wide broadcast when the target peer isn't known or connected + (unchanged legacy parity for that case).""" _tmp_data = b"".join([data[:15], tmp_bits.to_bytes(1, "big"), data[16:20], dmrpkt]) try: - self._send_to_system(d_system, _tmp_data) + if d_peer_id is not None: + protocols = self._get_protocols() if self._get_protocols else {} + proto = protocols.get(d_system) + if proto is not None and d_peer_id in getattr(proto, "_peers", {}): + proto.send_peer(d_peer_id, _tmp_data) + else: + self._send_to_system(d_system, _tmp_data) + else: + self._send_to_system(d_system, _tmp_data) except Exception as exc: logger.warning("(%s) send_data_to_hbp %s failed: %s", source_system, d_system, exc) return diff --git a/src/adn_server/application/routing/lc_ta.py b/src/adn_server/application/routing/lc_ta.py index 428fe4d..fd2a9e3 100644 --- a/src/adn_server/application/routing/lc_ta.py +++ b/src/adn_server/application/routing/lc_ta.py @@ -49,7 +49,7 @@ from typing import Any from bitarray import bitarray from ...domain import int_id -from ...domain.dmr.const import LC_OPT +from ...domain.dmr.const import LC_OPT_G, LC_OPT_U from ...domain.talker_alias import DMRA_BLOCK_COUNT from ..talker_alias_use_cases import passthrough_complete, talker_alias_settings from .helpers import EMB_LC_SLICE @@ -417,11 +417,14 @@ class LcTaMixin: dst_id: bytes, slot: int, stream_id: bytes, + call_type: str = "group", ) -> None: """REPEAT on VHEAD: standalone DMRA plus embedded TA state for downlink DMRD. WPSD/MMDVMHost ignores standalone DMRA UDP; it only displays Talker Alias decoded - from embedded LC inside repeated voice bursts (B–E). + from embedded LC inside repeated voice bursts (B–E). Private (unit) calls use the + same passthrough/inject policy as group calls -- only the LC FLCO opt byte differs + (LC_OPT_U vs LC_OPT_G), since the destination is a subscriber ID, not a talkgroup. """ settings = talker_alias_settings(self._config, system_name) if settings["enabled"] and self._get_protocols: @@ -433,7 +436,8 @@ class LcTaMixin: st = {} status[slot] = st if st.get("REP_STREAM_ID") != stream_id: - dst_lc = LC_OPT + dst_id + rf_src + lc_opt = LC_OPT_U if call_type == "unit" else LC_OPT_G + dst_lc = lc_opt + dst_id + rf_src st["REP_STREAM_ID"] = stream_id st["REP_EMB_LC"] = self._encode_emblc(dst_lc) self._init_talker_alias_embed( @@ -563,8 +567,8 @@ class LcTaMixin: ) -> None: """Replace embedded LC on voice bursts B–E (legacy bridge.py parity). - Group-call superframes always carry the **destination** group LC (``emb_key``), - re-encoded for the rewritten TGID — this is required for the voice to be accepted + Superframes always carry the **destination** LC (``emb_key`` — group LC re-encoded + for the rewritten TGID, or unit LC for private calls) — this is required for the voice to be accepted by the receiving MMDVM (a mismatched embedded LC causes packet loss). When a Talker Alias is available (``TX_TA_EMB``: injected template or the source TA re-encoded from its DMRA/voice blocks) it is overlaid on alternate superframes. diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index ab8d5bc..5106a76 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -1161,7 +1161,9 @@ class RoutingUseCases( # SUB_MAP lookup (legacy ~2297-2312 / ~3099-3114) sub_map = self._config.get("_SUB_MAP", {}) if dst_id in sub_map: - _d_system, _d_slot, _d_time = sub_map[dst_id] + _sub_map_entry = sub_map[dst_id] + _d_system, _d_slot, _d_time = _sub_map_entry[:3] + _d_peer_id = _sub_map_entry[3] if len(_sub_map_entry) > 3 else None _d_proto = protocols.get(_d_system) if _d_proto: _dst_slot = getattr(_d_proto, "STATUS", {}).get(_d_slot, {}) @@ -1170,7 +1172,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _tmp_bits = _bits ^ (1 << 7) if slot != _d_slot else _bits - self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_d_peer_id) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( system_name, _d_system, _d_slot, _int_dst_id, _dst_slot, _d_sys_cfg, pkt_time, @@ -1197,7 +1199,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits - self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( system_name, _d_system, _d_slot, _int_dst_id, _dst_slot, _d_sys_cfg, pkt_time, @@ -1212,7 +1214,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits - self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id) + self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( system_name, _d_system, _d_slot, _int_dst_id, _dst_slot, _d_sys_cfg, pkt_time, @@ -1241,7 +1243,7 @@ class RoutingUseCases( _bits = data[15] if len(data) > 15 else 0 sub_map = self._config.get("_SUB_MAP", {}) if sub_map is not None: - sub_map[rf_src] = (system_name, slot, pkt_time) + sub_map[rf_src] = (system_name, slot, pkt_time, peer_id) systems_cfg = self._config.get("SYSTEMS", {}) protocols = self._get_protocols() if self._get_protocols else {} source_proto = protocols.get(system_name) @@ -1260,14 +1262,31 @@ class RoutingUseCases( ) return slot_st["RX_START"] = pkt_time + self._pvt_same_system_dst_peer_id = None if dst_id in sub_map: - if sub_map[dst_id][0] != system_name: - self._pvt_targets = [sub_map[dst_id][0]] + _dst_entry = sub_map[dst_id] + if _dst_entry[0] != system_name: + self._pvt_targets = [_dst_entry[0]] + # Exact hotspot the destination was last heard on, if SUB_MAP has it + # (4th element) -- lets delivery target that one peer on the target + # system instead of send_to_system's broadcast-to-every-peer default. + self._pvt_target_peer_ids = { + _dst_entry[0]: (_dst_entry[3] if len(_dst_entry) > 3 else None) + } else: self._pvt_targets = [] + self._pvt_target_peer_ids = {} logger.error("PRIVATE call to a subscriber on the same system, send nothing") + # Delivery is already handled by the MASTER's own local REPEAT + # (_pvt_repeat_targets in udp_hbp.py) -- this branch only exists so the + # monitor also learns who is *receiving*, which nothing else reports. + if len(_dst_entry) > 3 and _dst_entry[3] is not None: + _dest_peers = getattr(source_proto, "_peers", {}) if source_proto else {} + if _dst_entry[3] in _dest_peers: + self._pvt_same_system_dst_peer_id = _dst_entry[3] else: self._pvt_targets = [] + self._pvt_target_peer_ids = {} logger.info( "(%s) *PRIVATE CALL START* STREAM ID: %s SUB: %s PEER: %s DST: %s, TS: %s, FORWARD: %s", system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, self._pvt_targets, @@ -1280,12 +1299,22 @@ class RoutingUseCases( system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id) ) ) + if self._pvt_same_system_dst_peer_id is not None: + self._send_routing_event( + "PRIVATE VOICE,START,TX,{},{},{},{},{},{},{}".format( + system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), + int_id(self._pvt_same_system_dst_peer_id), + ) + ) for _target in getattr(self, "_pvt_targets", []): target_proto = protocols.get(_target) if not target_proto: continue _target_status = getattr(target_proto, "STATUS", {}) _target_system = systems_cfg.get(_target, {}) + _target_peer_id = None + if _target_system.get("MODE") != "OPENBRIDGE": + _target_peer_id = getattr(self, "_pvt_target_peer_ids", {}).get(_target) if _target_system.get("MODE") == "OPENBRIDGE": if _target_system.get("ENHANCED_OBP") and "_bcka" in _target_system and _target_system["_bcka"] < pkt_time - 60: continue @@ -1345,16 +1374,27 @@ class RoutingUseCases( ts_st["TX_PEER"] = peer_id logger.info("(%s) PRIVATE call bridged to HBP System: %s TS: %s, DST: %s", system_name, _target, slot, int_id(dst_id)) if not _unit_data: - self._send_routing_event( - "PRIVATE VOICE,START,TX,{},{},{},{},{},{}".format( - _target, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), - ).encode("utf-8", "ignore") - ) + if _target_peer_id is not None and _target_peer_id in getattr(target_proto, "_peers", {}): + self._send_routing_event( + "PRIVATE VOICE,START,TX,{},{},{},{},{},{},{}".format( + _target, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), + int_id(_target_peer_id), + ) + ) + else: + self._send_routing_event( + "PRIVATE VOICE,START,TX,{},{},{},{},{},{}".format( + _target, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), + ) + ) ts_st["TX_TIME"] = pkt_time ts_st["TX_TYPE"] = dtype_vseq send_data = data try: - self._send_to_system(_target, send_data) + if _target_peer_id is not None and target_proto is not None and _target_peer_id in getattr(target_proto, "_peers", {}): + target_proto.send_peer(_target_peer_id, send_data) + else: + self._send_to_system(_target, send_data) except Exception as e: logger.warning("(ROUTER) send_to_system %s failed: %s", _target, e) if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM: @@ -1375,6 +1415,21 @@ class RoutingUseCases( call_duration, ) ) + _same_system_dst_peer_id = getattr(self, "_pvt_same_system_dst_peer_id", None) + if _same_system_dst_peer_id is not None: + self._send_routing_event( + "PRIVATE VOICE,END,TX,{},{},{},{},{},{},{:.2f},{}".format( + system_name, + int_id(stream_id), + int_id(peer_id), + int_id(rf_src), + slot, + int_id(dst_id), + call_duration, + int_id(_same_system_dst_peer_id), + ) + ) + self._pvt_same_system_dst_peer_id = None if slot_st: if _unit_data: # Keep stream continuity for multi-frame unit data without marking the diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 7717d3b..60b5625 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -687,6 +687,12 @@ def run_peer_server( def sighup_reload_config(_sig, _frame): """Reload adn-server.yaml (SYSTEMS, GLOBAL); keeps active streams on unchanged listeners.""" logger.info("(CONFIG-RELOAD) SIGHUP received, scheduling reload") + if aliases_cfg.get("SUB_MAP_FILE"): + try: + sub_map_store.save(sub_map_path, sub_map) + logger.info("(SUBSCRIBER) Writing SUB_MAP to disk (SIGHUP)") + except Exception as e: + logger.warning("(SUBSCRIBER) Cannot write SUB_MAP to file: %s", e) reactor.callLater(0, _do_config_reload) signal.signal(signal.SIGTERM, sig_handler) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index c78b9b2..d9e23ce 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -207,7 +207,7 @@ class HBPProtocol(DatagramProtocol): on_deactivate_dynamic_relays: Callable[[str], None] | None = None, on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None, on_talker_alias_local_repeat: Callable[[str, bytes, bytes, bytes], None] | None = None, - on_talker_alias_repeat_prepare: Callable[[str, bytes, bytes, bytes, int, bytes], None] | None = None, + on_talker_alias_repeat_prepare: Callable[[str, bytes, bytes, bytes, int, bytes, str], None] | None = None, on_talker_alias_repeat_burst: Callable[[str, int, bytes, int, bytes], bytes] | None = None, on_talker_alias_stream_end: Callable[[str, bytes], None] | None = None, on_dmra_fragment_stored: Callable[[str, bytes, bytes, bytes], None] | None = None, @@ -423,6 +423,57 @@ class HBPProtocol(DatagramProtocol): index = self._ensure_downlink_index() return index.candidates(slot, tgid, connected_count=connected) + def _pvt_repeat_targets(self, dst_id: bytes) -> tuple[bytes, ...] | None: + """Which peer(s) get a private-call REPEAT, using _SUB_MAP's last-heard peer. + + Returns ``None`` when we've never heard this destination transmit — + the caller should fall back to broadcasting to every peer (legacy + parity: hblink.py's master_datagramReceived repeats unconditionally + when the destination is unknown). Returns ``()`` when we know the + destination isn't reachable from here at all (heard on a different + system, or that specific peer isn't currently connected) — refusing + to blast the call to peers it can't possibly be for. Returns a + 1-tuple with the exact peer_id when it's still connected here. + + The report_slot / "SYSTEM-N" monitor display name is NOT used for + this — it's cosmetic and can be reassigned to a different peer across + refreshes (self-service peers without a stable report_slot fall back + to a sorted-by-id allocation recomputed each time). The raw peer_id + is what's actually stable. + """ + sub_map = self._CONFIG.get("_SUB_MAP") + if not sub_map: + return None + entry = sub_map.get(dst_id) + if entry is None: + return None + d_system = entry[0] + d_peer_id = entry[3] if len(entry) > 3 else None + if d_system != self._system or d_peer_id is None: + return () + if d_peer_id in self._peers: + return (d_peer_id,) + return () + + def _purge_sub_map_for_peer(self, peer_id: bytes) -> None: + """Drop every SUB_MAP entry last heard on ``peer_id``. + + Called when a hotspot finishes (re)connecting (RPTC -> CONNECTION + "YES"). A subscriber's last-known hotspot is only trustworthy between + that hotspot's connects: while it was offline the radio may have + moved elsewhere, and there's no way to tell without seeing it + transmit again. So on reconnect we forget where any subscriber was + last heard on this peer, rather than risk a stale entry silently + misdirecting a private call. The subscriber simply can't receive a + private call again until it transmits and gets re-learned. + """ + sub_map = self._CONFIG.get("_SUB_MAP") + if not sub_map: + return + stale = [rf_src for rf_src, entry in sub_map.items() if len(entry) > 3 and entry[3] == peer_id] + for rf_src in stale: + del sub_map[rf_src] + def _peer_mesh_config(self) -> PeerMeshConfig: _global = self._CONFIG.get("GLOBAL", {}) _sid = _global.get("SERVER_ID", b"\x00\x00\x00\x00") @@ -579,6 +630,15 @@ class HBPProtocol(DatagramProtocol): def _peer_should_receive_dmrd(self, peer_id: bytes, packet: bytes) -> bool: if peer_id not in self._peers: return False + if len(packet) >= 16 and (packet[15] & 0x40): + # Private (unit) call: parse_dmrd_route_fields()/parse_dmrd_burst_fields() + # return None for these by design (they're group/vcsbk-only helpers), which + # would otherwise fall into the "parsed is None" branch below and, with more + # than one connected peer, incorrectly block delivery to everyone. Legacy + # repeats private calls unconditionally to every peer of the same master + # (hblink.py master_datagramReceived) — no per-peer OPTIONS/TG filtering + # applies, since the destination radio (not the hotspot) decides relevance. + return True parsed = parse_dmrd_route_fields(packet) if parsed is None: if self._inject_multi_peer_options_filter(): @@ -1190,10 +1250,12 @@ class HBPProtocol(DatagramProtocol): logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY SYSTEM TS2 ACL", self._system, int_id(_stream_id), int_id(_dst_id)) self._laststrid[_slot] = _stream_id return - # SUB_MAP update (legacy routerHBP.dmrd_received) + # SUB_MAP update (legacy routerHBP.dmrd_received). 4th element + # (peer_id) is new — lets same-system private-call repeat target + # the exact hotspot instead of broadcasting to every peer. sub_map = self._CONFIG.get("_SUB_MAP") if sub_map is not None: - sub_map[_rf_src] = (self._system, _slot, pkt_time) + sub_map[_rf_src] = (self._system, _slot, pkt_time, _peer_id) self.note_dmrd_stream(_peer_id, _rf_src, _stream_id) if ( _call_type in ("group", "vcsbk") @@ -1274,26 +1336,23 @@ class HBPProtocol(DatagramProtocol): if ( _repeat_ok and self._config.get("REPEAT", True) - and _call_type in ("group", "vcsbk") + and _call_type in ("group", "vcsbk", "unit") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD ): if self._on_talker_alias_repeat_prepare: self._on_talker_alias_repeat_prepare( - self._system, _peer_id, _rf_src, _dst_id, _slot, _stream_id, + self._system, _peer_id, _rf_src, _dst_id, _slot, _stream_id, _call_type, ) elif self._on_talker_alias_local_repeat: self._on_talker_alias_local_repeat( self._system, _peer_id, _rf_src, _stream_id, ) - if ( - _repeat_ok - and self._config.get("REPEAT", True) - and _call_type in ("group", "vcsbk") - ): + if _repeat_ok and self._config.get("REPEAT", True): _repeat_tail = _data[15:] if ( - _dtype_vseq in (1, 2, 3, 4) + _call_type in ("group", "vcsbk", "unit") + and _dtype_vseq in (1, 2, 3, 4) and len(_data) >= 53 and self._on_talker_alias_repeat_burst ): @@ -1302,9 +1361,18 @@ class HBPProtocol(DatagramProtocol): ) _repeat_tail = b"".join([_data[15:20], _dmrpkt_out, _data[53:]]) _repeat_pkt = b"".join([_data[:11], _peer_id, _repeat_tail]) - for _peer in self._iter_downlink_peers(_repeat_pkt): - if _peer != _peer_id: - self.send_peer(_peer, _repeat_pkt) + if _call_type == "unit": + _pvt_targets = self._pvt_repeat_targets(_dst_id) + else: + _pvt_targets = None + if _pvt_targets is None: + for _peer in self._iter_downlink_peers(_repeat_pkt): + if _peer != _peer_id: + self.send_peer(_peer, _repeat_pkt) + else: + for _peer in _pvt_targets: + if _peer != _peer_id: + self.send_peer(_peer, _repeat_pkt) # TG 4000: reset after REPEAT so peers see the packet (legacy order) if self._handle_tg4000_packet( _peer_id, _slot, _int_dst_id, _call_type, _frame_type, _dtype_vseq, @@ -1388,7 +1456,7 @@ class HBPProtocol(DatagramProtocol): self._on_in_band_signalling(self._system, _slot, _dst_id, pkt_time) if ( _accepted - and _call_type in ("group", "vcsbk") + and _call_type in ("group", "vcsbk", "unit") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM and _slot in self.STATUS @@ -1555,6 +1623,7 @@ class HBPProtocol(DatagramProtocol): _rptc_field_str(self._peers[_peer_id]["SOFTWARE_ID"]), _rptc_field_str(self._peers[_peer_id]["DESCRIPTION"]), ) + self._purge_sub_map_for_peer(_peer_id) self._refresh_connected_peer_count() self._mark_downlink_index_dirty() self._config_push_throttle.note_peer_connected() @@ -1748,10 +1817,11 @@ class HBPProtocol(DatagramProtocol): logger.info("(%s) CALL DROPPED WITH STREAM ID %s ON TGID %s BY SYSTEM TS2 ACL", self._system, int_id(_stream_id), int_id(_dst_id)) self._laststrid[_slot] = _stream_id return - # SUB_MAP update (legacy routerHBP.dmrd_received) + # SUB_MAP update (legacy routerHBP.dmrd_received). 4th element + # (peer_id) is new — see the MASTER-mode write site for why. sub_map = self._CONFIG.get("_SUB_MAP") if sub_map is not None: - sub_map[_rf_src] = (self._system, _slot, pkt_time) + sub_map[_rf_src] = (self._system, _slot, pkt_time, _peer_id) # TG 4000: reset after ACL/SUB_MAP (legacy order — routerHBP.dmrd_received) if self._handle_tg4000_packet( _peer_id, _slot, _int_dst_id, _call_type, _frame_type, _dtype_vseq, @@ -2389,7 +2459,7 @@ def HBPProtocolFactory( on_deactivate_dynamic_relays: Callable[[str], None] | None = None, on_obp_bcsq_received: Callable[[str, bytes, bytes], None] | None = None, on_talker_alias_local_repeat: Callable[[str, bytes, bytes, bytes], None] | None = None, - on_talker_alias_repeat_prepare: Callable[[str, bytes, bytes, bytes, int, bytes], None] | None = None, + on_talker_alias_repeat_prepare: Callable[[str, bytes, bytes, bytes, int, bytes, str], None] | None = None, on_talker_alias_repeat_burst: Callable[[str, int, bytes, int, bytes], bytes] | None = None, on_talker_alias_stream_end: Callable[[str, bytes], None] | None = None, on_dmra_fragment_stored: Callable[[str, bytes, bytes, bytes], None] | None = None, diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py index 9dc696f..2e3c462 100644 --- a/tests/application/test_monitor_topology.py +++ b/tests/application/test_monitor_topology.py @@ -303,6 +303,65 @@ def test_remap_voice_event_passes_through_non_proxy_systems() -> None: assert remap_inject_proxy_voice_event(raw, {}, {}) == raw +def test_private_voice_tx_remaps_to_dest_peer_id_not_tg_subscription() -> None: + """Regression: a PRIVATE VOICE TX event was silently dropped -- the generic TX + path treats field 8 as a talkgroup and fans out via TG-subscription eligibility + (_peers_receiving_tgid), which a subscriber-id destination never matches. The + trailing dest_peer_id field (routing_use_cases._pvt_call_received) must route it + directly to that one hotspot's SYSTEM-N row instead.""" + peers = { + bytes_4(730039110): _peer(), + bytes_4(730039101): _peer(), + } + peer_slots = {bytes_4(730039110): 1, bytes_4(730039101): 8} + config = _proxy_config(peers) + raw = "PRIVATE VOICE,START,TX,SYSTEM,1610544978,730039110,7300391,2,7300392,730039101" + events = remap_inject_proxy_voice_events(raw, config, config["SYSTEMS"], peer_slots) + assert len(events) == 1 + parts = events[0].split(",") + assert parts[3] == "SYSTEM-8" + assert parts[2] == "TX" + assert parts[-1] == "730039101" + + +def test_private_voice_tx_dropped_when_dest_peer_not_connected() -> None: + peers = {bytes_4(730039110): _peer()} + peer_slots = {bytes_4(730039110): 1} + config = _proxy_config(peers) + raw = "PRIVATE VOICE,START,TX,SYSTEM,1610544978,730039110,7300391,2,7300392,730039101" + events = remap_inject_proxy_voice_events(raw, config, config["SYSTEMS"], peer_slots) + assert events == [] + + +def test_private_voice_tx_without_dest_peer_id_falls_back_to_legacy_tg_fanout() -> None: + """Old-server-format event (no trailing peer id) -- must not crash, and since a + subscriber id is never a static TG, no companion TX is fabricated.""" + peers = { + bytes_4(730039110): _peer(), + bytes_4(730039101): _peer(), + } + peer_slots = {bytes_4(730039110): 1, bytes_4(730039101): 8} + config = _proxy_config(peers) + raw = "PRIVATE VOICE,START,TX,SYSTEM,1610544978,730039110,7300391,2,7300392" + events = remap_inject_proxy_voice_events(raw, config, config["SYSTEMS"], peer_slots) + assert events == [] + + +def test_private_voice_end_tx_remaps_to_dest_peer_id_with_duration_field() -> None: + peers = { + bytes_4(730039110): _peer(), + bytes_4(730039101): _peer(), + } + peer_slots = {bytes_4(730039110): 1, bytes_4(730039101): 8} + config = _proxy_config(peers) + raw = "PRIVATE VOICE,END,TX,SYSTEM,1610544978,730039110,7300391,2,7300392,5.87,730039101" + events = remap_inject_proxy_voice_events(raw, config, config["SYSTEMS"], peer_slots) + assert len(events) == 1 + parts = events[0].split(",") + assert parts[3] == "SYSTEM-8" + assert parts[-1] == "730039101" + + def test_hangtime_blocks_monitor_tx_fanout_to_blocked_peer() -> None: """Companion TX / OBP fan-out must not light peers blocked by GROUP_HANGTIME.""" import time diff --git a/tests/application/test_report_payloads.py b/tests/application/test_report_payloads.py index 45d8c51..be44e34 100644 --- a/tests/application/test_report_payloads.py +++ b/tests/application/test_report_payloads.py @@ -253,6 +253,63 @@ def test_parse_group_voice_start_matches_example(validator: jsonschema.Draft2020 validator.validate(doc) +@pytest.mark.parametrize( + "label", + [ + "UNIT CSBK", + "UNIT DATA HEADER", + "UNIT VCSBK 1/2 DATA BLOCK", + "UNIT VCSBK 3/4 DATA BLOCK", + ], +) +def test_parse_unit_dtype_labels_map_to_unit_family(label: str) -> None: + """Regression: routing_use_cases._dtype_labels emits these instead of the + generic "UNIT DATA" for dtype_vseq in (3, 6, 7, 8) — private-call CSBK + setup and SMS/GPS header/VCSBK blocks. Before these were added to + _CSV_FAMILIES, parse_bridge_event_csv() returned None and the event was + silently dropped, so a private call whose CSBK handshake never escalates + to voice never reached the monitor's Last Heard log at all.""" + csv = f"{label},DATA,RX,SYSTEM-8,656846929,730039210,7300392,2,714001" + doc = parse_bridge_event_csv(csv) + assert doc is not None + assert doc["call_family"] == "UNIT" + assert doc["phase"] == "DATA" + + +def test_parse_bridge_event_csv_unknown_family_still_none() -> None: + csv = "SOMETHING ELSE,DATA,RX,SYSTEM-8,656846929,730039210,7300392,2,714001" + assert parse_bridge_event_csv(csv) is None + + +def test_parse_private_voice_tx_captures_trailing_dest_peer_id() -> None: + """Regression: the JSON voice_event schema had no field for the destination + hotspot's peer id (routing_use_cases._pvt_call_received appends it after the + legacy CSV fields) -- it was silently dropped in this CSV->JSON conversion, + so the monitor could never know which single hotspot was receiving a private + call, and always reconstructed a bare "0" in its place (report_mapper. + voice_event_to_csv_parts's is_announcement default).""" + csv = "PRIVATE VOICE,START,TX,SYSTEM,1610544978,730039110,7300391,2,7300392,730039101" + doc = parse_bridge_event_csv(csv) + assert doc is not None + assert doc["dest_peer_id"] == 730039101 + assert doc["is_announcement"] is False + + +def test_parse_private_voice_end_tx_captures_dest_peer_id_after_duration() -> None: + csv = "PRIVATE VOICE,END,TX,SYSTEM,1610544978,730039110,7300391,2,7300392,5.87,730039101" + doc = parse_bridge_event_csv(csv) + assert doc is not None + assert doc["dest_peer_id"] == 730039101 + assert doc["duration_s"] == pytest.approx(5.87) + + +def test_parse_private_voice_rx_has_no_dest_peer_id() -> None: + csv = "PRIVATE VOICE,START,RX,SYSTEM,1610544978,730039110,7300391,2,7300392" + doc = parse_bridge_event_csv(csv) + assert doc is not None + assert "dest_peer_id" not in doc + + def test_routing_delta_matches_example(validator: jsonschema.Draft202012Validator) -> None: with (_EXAMPLES_DIR / "routing_table.json").open(encoding="utf-8") as fh: previous = json.load(fh) diff --git a/tests/harness/deterministic.py b/tests/harness/deterministic.py index fd3b46b..2634775 100644 --- a/tests/harness/deterministic.py +++ b/tests/harness/deterministic.py @@ -207,6 +207,11 @@ class FakeHbpProtocol: def __init__(self, name: str) -> None: self.name = name self.STATUS: dict[int, dict[str, Any]] = {1: {}, 2: {}} + self._peers: dict[bytes, Any] = {} + self.sent_to_peer: list[tuple[bytes, bytes]] = [] + + def send_peer(self, peer_id: bytes, packet: bytes) -> None: + self.sent_to_peer.append((peer_id, packet)) class FakeReportSender: diff --git a/tests/infrastructure/test_hbp_private_call_targeting.py b/tests/infrastructure/test_hbp_private_call_targeting.py new file mode 100644 index 0000000..3f9f951 --- /dev/null +++ b/tests/infrastructure/test_hbp_private_call_targeting.py @@ -0,0 +1,134 @@ +# ADN DMR Peer Server - tests infrastructure hbp private call targeting +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""SUB_MAP-based precise private-call targeting on a MASTER with 3+ peers. + +Legacy (hblink.py master_datagramReceived) broadcasts every private call to +every connected peer unconditionally -- the receiving hotspot's own ACL/sub +list is the only filter, so a private call between two unrelated users still +reaches every other hotspot's air interface. _pvt_repeat_targets narrows this: +once SUB_MAP has recorded which peer_id last carried a given subscriber's +traffic, a private call to that subscriber is delivered only to that peer, +never broadcast. Unknown destinations still fall back to broadcast (legacy +parity) -- see test_hbp_repeat_private_call.py.""" + +from __future__ import annotations + +from tests.harness.deterministic import DeterministicScenario, PacketSpec +from tests.support.hbp_repeat_stack import build_hbp_repeat_stack + +from adn_server.domain import bytes_3, bytes_4 +from adn_server.infrastructure.hbp_constants import RPTC, RPTK, RPTL +from adn_server.infrastructure.twisted_adapters.udp_hbp import _calc_hash, _get_passphrase_bytes + +_PEER_TX = bytes_4(730039210) +_PEER_RX = bytes_4(730039101) +_PEER_OTHER = bytes_4(730039199) +_ADDR_TX = ("10.0.0.1", 62001) +_ADDR_RX = ("10.0.0.2", 62002) +_ADDR_OTHER = ("10.0.0.3", 62003) + +_DST_SUB = 7304011 + + +def _private_spec() -> PacketSpec: + return PacketSpec( + peer_id=730039210, + rf_src=7300392, + dst_id=_DST_SUB, + slot=2, + call_type="unit", + stream_id=0xA1B2C3D4, + payload=b"\x00" * 33, + ) + + +def _fire_private_call(stack, base: PacketSpec) -> None: + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + stack.transport.clear() + stack.inject_spec( + DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX, + ) + + +def test_private_call_delivered_only_to_last_known_peer_not_broadcast() -> None: + stack = build_hbp_repeat_stack(talker_alias=True) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;") + stack.register_peer(_PEER_OTHER, _ADDR_OTHER, options="TS2=7304;") + stack.config["_SUB_MAP"] = { + bytes_3(_DST_SUB): (stack.system_name, 2, 1_700_000_000.0, _PEER_RX), + } + + _fire_private_call(stack, _private_spec()) + + assert stack.transport.for_addr(_ADDR_RX), "known peer must still receive the call" + assert not stack.transport.for_addr(_ADDR_OTHER), "uninvolved peer must not see this private call" + + +def test_private_call_dropped_when_known_peer_no_longer_connected() -> None: + stack = build_hbp_repeat_stack(talker_alias=True) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_OTHER, _ADDR_OTHER, options="TS2=7304;") + # _PEER_RX was last heard here, but is no longer registered/connected. + stack.config["_SUB_MAP"] = { + bytes_3(_DST_SUB): (stack.system_name, 2, 1_700_000_000.0, _PEER_RX), + } + + _fire_private_call(stack, _private_spec()) + + assert not stack.transport.for_addr(_ADDR_OTHER), "must not blast to peers that can't be the destination" + + +def test_private_call_not_delivered_locally_when_known_on_different_system() -> None: + stack = build_hbp_repeat_stack(talker_alias=True) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_OTHER, _ADDR_OTHER, options="TS2=7304;") + stack.config["_SUB_MAP"] = { + bytes_3(_DST_SUB): ("MASTER-B", 2, 1_700_000_000.0, _PEER_RX), + } + + _fire_private_call(stack, _private_spec()) + + assert not stack.transport.for_addr(_ADDR_OTHER), "destination known elsewhere must not repeat locally" + + +def test_sub_map_entries_for_peer_purged_on_reconnect() -> None: + """Once a hotspot finishes (re)logging in, any SUB_MAP entry pointing at it is + stale -- the radio may have moved to a different hotspot while this one was + offline and there's no way to tell without hearing it transmit again.""" + stack = build_hbp_repeat_stack(talker_alias=True) + stack.hbp._config["PASSPHRASE"] = b"test-passphrase" + stack.config["_SUB_MAP"] = { + bytes_3(_DST_SUB): (stack.system_name, 2, 1_700_000_000.0, _PEER_RX), + bytes_3(1234567): (stack.system_name, 1, 1_700_000_000.0, _PEER_OTHER), + } + passphrase = _get_passphrase_bytes(stack.hbp._config) + + stack.hbp.datagramReceived(RPTL + _PEER_RX, _ADDR_RX) + salt = bytes_4(stack.hbp._peers[_PEER_RX]["SALT"]) + stack.hbp.datagramReceived(RPTK + _PEER_RX + _calc_hash(salt, passphrase), _ADDR_RX) + stack.hbp.datagramReceived( + RPTC + _PEER_RX + b"CE5RPY " + b"\x00" * 85 + b"4", _ADDR_RX, + ) + + assert stack.hbp._peers[_PEER_RX]["CONNECTION"] == "YES" + assert bytes_3(_DST_SUB) not in stack.config["_SUB_MAP"], "entry pointing at the reconnected peer must be gone" + assert bytes_3(1234567) in stack.config["_SUB_MAP"], "entries for other peers must be untouched" diff --git a/tests/infrastructure/test_hbp_repeat_private_call.py b/tests/infrastructure/test_hbp_repeat_private_call.py new file mode 100644 index 0000000..a6ed2c4 --- /dev/null +++ b/tests/infrastructure/test_hbp_repeat_private_call.py @@ -0,0 +1,135 @@ +# ADN DMR Peer Server - tests infrastructure hbp repeat private call +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""REPEAT path through real HBPProtocol for private (unit-to-unit) calls. + +Regression: the raw intra-MASTER REPEAT loop was gated on +`call_type in ("group", "vcsbk")`, so two peers registered under the same +MASTER system never saw each other's private calls at all — legacy +(hblink.py master_datagramReceived) repeats every call_type unconditionally. +This test exercises the real infrastructure layer (not the routing-only +harness in tests/harness/deterministic.py, which calls dmrd_received() +directly and never touches this REPEAT loop).""" + +from __future__ import annotations + +from tests.harness.deterministic import DeterministicScenario, PacketSpec +from tests.support.hbp_repeat_stack import build_hbp_repeat_stack + +from adn_server.domain import bytes_4 + +_PEER_TX = bytes_4(730039210) +_PEER_RX = bytes_4(730039101) +_ADDR_TX = ("10.0.0.1", 62001) +_ADDR_RX = ("10.0.0.2", 62002) + + +def _private_spec() -> PacketSpec: + return PacketSpec( + peer_id=730039210, + rf_src=7300392, + dst_id=7304011, # 7-digit subscriber ID, not a talkgroup + slot=2, + call_type="unit", + stream_id=0xA1B2C3D4, + payload=b"\x00" * 33, + ) + + +def test_private_call_is_repeated_to_other_peer_on_same_master() -> None: + stack = build_hbp_repeat_stack(talker_alias=True) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;") + base = _private_spec() + + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + stack.transport.clear() + stack.inject_spec( + DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX, + ) + + downlink = stack.transport.for_addr(_ADDR_RX) + assert downlink, "private call must be repeated to the other peer on the same MASTER" + assert downlink[0][11:15] == _PEER_RX + + +def test_private_call_reaches_peer_under_inject_only_proxy_with_multiple_peers() -> None: + """Regression: parse_dmrd_burst_fields() returns None for private (unit) calls by + design (it's a group/vcsbk-only helper) — _peer_should_receive_dmrd's fallback for + "parsed is None" under an inject-only multi-peer proxy is `connected_count <= 1`, + which blocks delivery entirely once a 2nd hotspot connects. Real deployments run + with PROXY.TARGET_SYSTEM set and many peers, so this was silently dropping every + private call — reproduced here since build_hbp_repeat_stack's default config has + no PROXY section (the other tests in this file don't exercise this branch at all).""" + stack = build_hbp_repeat_stack(talker_alias=True) + stack.config["PROXY"] = {"TARGET_SYSTEM": stack.system_name} + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;") + base = _private_spec() + + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + stack.transport.clear() + stack.inject_spec( + DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX, + ) + + downlink = stack.transport.for_addr(_ADDR_RX) + assert downlink, "private call must reach the other peer even under inject-only multi-peer proxy" + assert downlink[0][11:15] == _PEER_RX + + +def test_private_call_repeat_leaves_burst_payload_unchanged_when_talker_alias_disabled() -> None: + stack = build_hbp_repeat_stack(talker_alias=False) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;") + base = _private_spec() + + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + stack.transport.clear() + uplink = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1) + stack.inject_spec(uplink, _ADDR_TX) + + downlink = stack.transport.for_addr(_ADDR_RX) + assert len(downlink) == 1 + assert downlink[0][20:53] == uplink.payload + assert stack.hbp.STATUS.get(2, {}).get("TX_TA_EMB") is None + + +def test_private_call_repeat_embeds_unit_lc_when_talker_alias_enabled() -> None: + """Same rules as group calls: inject the configured TA when the source doesn't + supply its own -- only the LC opt byte differs (LC_OPT_U, subscriber destination + instead of a talkgroup).""" + stack = build_hbp_repeat_stack(talker_alias=True) + stack.register_peer(_PEER_TX, _ADDR_TX, options="TS2=7304;") + stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;") + base = _private_spec() + uplink = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1) + + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + stack.transport.clear() + stack.inject_spec(uplink, _ADDR_TX) + + downlink = stack.transport.for_addr(_ADDR_RX) + assert len(downlink) == 1 + slot_st = stack.hbp.STATUS[2] + assert slot_st.get("TX_TA_EMB") is not None + assert slot_st.get("REP_EMB_LC") is not None + # AMBE voice bits (outside the embedded LC slice) still pass through untouched. + assert downlink[0][20:53][:14] == uplink.payload[:14] diff --git a/tests/routing/test_unit_data_routing.py b/tests/routing/test_unit_data_routing.py index 059cd47..286f1d1 100644 --- a/tests/routing/test_unit_data_routing.py +++ b/tests/routing/test_unit_data_routing.py @@ -34,7 +34,7 @@ from tests.harness.deterministic import ( ) from tests.routing.unit_data_helpers import idle_hbp_slot -from adn_server.domain import bytes_3 +from adn_server.domain import bytes_3, bytes_4 from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM DAPRS_GATEWAY_ID = 900999 @@ -311,3 +311,68 @@ def test_unit_data_pvt_call_does_not_emit_private_voice_monitor_events() -> None assert scenario.report_factory is not None private_voice = [e for e in scenario.report_factory.events if e.startswith("PRIVATE VOICE")] assert private_voice == [] + + +@pytest.mark.behavior +def test_unit_data_sub_map_with_peer_id_targets_only_that_peer() -> None: + """Regression: cross-talk -- unit DATA to a known subscriber must reach only the + one hotspot it was last heard on, not every peer of the destination system + (send_to_system/send_peers broadcasts to all peers when no peer is targeted).""" + dst_sub = 712345 + dst_peer = bytes_4(730039101) + config = minimal_config(("MASTER-A", "MASTER-B")) + config["_SUB_MAP"] = {bytes_3(dst_sub): ("MASTER-B", 2, 1000.0, dst_peer)} + scenario = DeterministicScenario(config=config) + scenario.protocols["MASTER-B"].STATUS[2] = idle_hbp_slot() + scenario.protocols["MASTER-B"]._peers[dst_peer] = {} + base = PacketSpec(call_type="unit", dst_id=dst_sub, stream_id=0x52525253, slot=2) + + scenario.inject_unit("MASTER-A", DeterministicScenario.unit_data_header_spec(base)) + + assert scenario.protocols["MASTER-B"].sent_to_peer, "must deliver directly to the known peer" + assert scenario.protocols["MASTER-B"].sent_to_peer[0][0] == dst_peer + assert_not_forwarded(scenario, "MASTER-B") + + +@pytest.mark.behavior +def test_unit_data_sub_map_peer_not_connected_falls_back_to_broadcast() -> None: + """When the known peer isn't (or is no longer) connected, fall back to the old + broadcast-to-system behavior rather than silently dropping the packet.""" + dst_sub = 712345 + dst_peer = bytes_4(730039101) + config = minimal_config(("MASTER-A", "MASTER-B")) + config["_SUB_MAP"] = {bytes_3(dst_sub): ("MASTER-B", 2, 1000.0, dst_peer)} + scenario = DeterministicScenario(config=config) + scenario.protocols["MASTER-B"].STATUS[2] = idle_hbp_slot() + # dst_peer intentionally NOT registered in scenario.protocols["MASTER-B"]._peers. + base = PacketSpec(call_type="unit", dst_id=dst_sub, stream_id=0x52525254, slot=2) + + scenario.inject_unit("MASTER-A", DeterministicScenario.unit_data_header_spec(base)) + + assert not scenario.protocols["MASTER-B"].sent_to_peer + assert_forwarded(scenario, "MASTER-B", count=1, call_type="unit") + + +@pytest.mark.behavior +def test_pvt_call_sub_map_with_peer_id_targets_only_that_peer() -> None: + """Same cross-talk fix for the 7-digit pvt_call_received cross-system path (private + voice and ARS/LRRP downlink use this, independent of the unit-data SUB_MAP branch).""" + dst_peer = bytes_4(730039101) + config = minimal_config(("D-APRS", "SYSTEM")) + config["_SUB_MAP"] = {bytes_3(HOTSPOT_SUB_ID): ("SYSTEM", 2, 1000.0, dst_peer)} + scenario = DeterministicScenario(config=config) + scenario.protocols["SYSTEM"].STATUS[2] = idle_hbp_slot() + scenario.protocols["SYSTEM"]._peers[dst_peer] = {} + base = PacketSpec( + call_type="unit", + rf_src=DAPRS_GATEWAY_ID, + dst_id=HOTSPOT_SUB_ID, + stream_id=0x33416425, + slot=2, + ) + + scenario.inject_unit("D-APRS", DeterministicScenario.unit_data_header_spec(base)) + + assert scenario.protocols["SYSTEM"].sent_to_peer, "must deliver directly to the known peer" + assert scenario.protocols["SYSTEM"].sent_to_peer[0][0] == dst_peer + assert_not_forwarded(scenario, "SYSTEM")