From 89f6e178078ff0eb41d4dbbb17c5dc2964043a64 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Thu, 18 Jun 2026 16:20:33 -0400 Subject: [PATCH] fix: echo TG 9990 monitor parity and UA session exclusions MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Exclude 9990–9999 from SINGLE/UA locks, remap inject-only echo TX events, resolve bridge TX peer_id for monitor display, and fix dynamic TG restore callback. --- .../application/dynamic_tg_use_cases.py | 2 +- .../application/report/monitor_topology.py | 17 +++++-- src/adn_server/application/routing/helpers.py | 10 +++- .../application/routing_use_cases.py | 20 ++++++-- tests/application/test_dynamic_tg_persist.py | 46 +++++++++++++++++++ tests/application/test_monitor_topology.py | 7 +++ .../application/test_peer_single_downlink.py | 21 +++++++++ 7 files changed, 113 insertions(+), 10 deletions(-) diff --git a/src/adn_server/application/dynamic_tg_use_cases.py b/src/adn_server/application/dynamic_tg_use_cases.py index 324aac4..62213f4 100644 --- a/src/adn_server/application/dynamic_tg_use_cases.py +++ b/src/adn_server/application/dynamic_tg_use_cases.py @@ -113,7 +113,7 @@ class DynamicTgUseCases: ] tgids = restore_peer_ua_entries_to_memory(sys_cfg, peer_id, active, now=now) if self._on_restored is not None: - self._on_restored(peer_id, system_name, sys_cfg, active, now) + self._on_restored(peer_id, system_name, sys_cfg, active, now=now) if tgids: logger.info( "(DYNAMIC_TG) Restored %s TG(s) for peer %s on %s: %s", diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py index 35a9beb..7f8e437 100644 --- a/src/adn_server/application/report/monitor_topology.py +++ b/src/adn_server/application/report/monitor_topology.py @@ -31,7 +31,7 @@ from __future__ import annotations import copy from typing import Any -from adn_server.application.routing.helpers import peer_should_receive_group_voice +from adn_server.application.routing.helpers import is_special_tg, peer_should_receive_group_voice from adn_server.application.proxy.deployment import is_proxy_inject_only, proxy_target_system from adn_server.domain.value_objects import bytes_4, int_id @@ -246,7 +246,12 @@ def _peers_receiving_tgid( def _echo_tx_target_peer(parts: list[str], peers: dict[Any, Any]) -> bytes | None: - """Echo/static downlink: field 5 is 9990 and field 6 resolves one hotspot.""" + """Echo/service downlink TX: field 5 may be 9990 or hotspot id; field 8 is 9990–9999.""" + tgid_slot = _voice_event_tgid_slot(parts) + if tgid_slot is not None: + tgid, _ = tgid_slot + if is_special_tg(str(tgid)): + return _peer_key_from_voice_csv(parts, peers) if len(parts) <= 5: return None try: @@ -320,13 +325,19 @@ def remap_inject_proxy_voice_events( trx = parts[2].strip() if len(parts) > 2 else "" if trx == "TX": + tgid_slot = _voice_event_tgid_slot(parts) echo_peer = _echo_tx_target_peer(parts, peers) if echo_peer is not None: slot = slot_map.get(echo_peer) if slot is not None: + tx_parts = list(parts) + if tgid_slot is not None: + tgid, _ = tgid_slot + if is_special_tg(str(tgid)): + tx_parts[5] = str(tgid) return [ _remap_voice_event_to_slot( - parts, target=target, slot=slot, peer_key=echo_peer + tx_parts, target=target, slot=slot, peer_key=echo_peer ) ] try: diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 0f87294..ec24fd0 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -104,8 +104,14 @@ def tg4000_reset_on_vhead(int_dst_id: int, frame_type: int, dtype_vseq: int) -> def is_ua_session_tgid(tgid: int) -> bool: - """True when a keyed TG may be stored as a user-activated dynamic session.""" - return int(tgid) > 0 and int(tgid) != 4000 + """True when a keyed TG may be stored as a user-activated dynamic session. + + Excludes TG 4000 (reset command) and service/echo 9990–9999 (no SINGLE lock). + """ + t = int(tgid) + if t <= 0 or t == 4000: + return False + return not is_special_tg(str(t)) def obp_target_bcsq_quenches_stream( diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 906bc4c..14455f6 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -417,6 +417,11 @@ class RoutingUseCases( bridge_match_slot=bridge_match_slot, dst_int=dst_int, ) + _tx_report_peer = int_id(peer_id) + if not source_is_obp: + _tx_report_peer = int_id( + resolve_voice_peer_id(peer_id, rf_src, system_name, systems_cfg) + ) forwarded = [] _leg_iter: list[tuple[str, dict[str, Any]]] = [ ( @@ -496,7 +501,7 @@ class RoutingUseCases( ) self._send_routing_event( "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid) + entry["SYSTEM"], int_id(stream_id), _tx_report_peer, int_id(rf_src), entry.get("TS", 1), int_id(target_tgid) ) ) if "EMB_LC" not in _target_status[stream_id]: @@ -533,7 +538,7 @@ class RoutingUseCases( call_duration = pkt_time - _target_status[stream_id].get("START", pkt_time) self._send_routing_event( "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration + entry["SYSTEM"], int_id(stream_id), _tx_report_peer, int_id(rf_src), entry.get("TS", 1), int_id(target_tgid), call_duration ) ) elif dtype_vseq in (1, 2, 3, 4): @@ -624,7 +629,7 @@ class RoutingUseCases( "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( entry["SYSTEM"], int_id(stream_id), - int_id(peer_id), + _tx_report_peer, int_id(rf_src), entry_ts, int_id(entry_tgid_b), @@ -648,11 +653,18 @@ class RoutingUseCases( dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197] call_duration = pkt_time - _ts_st.get("TX_START", pkt_time) _end_peer = _ts_st.get("TX_PEER", peer_id) + _end_report_peer = int_id(_end_peer) + if not source_is_obp: + _end_report_peer = int_id( + resolve_voice_peer_id( + _end_peer, rf_src, system_name, systems_cfg + ) + ) self._send_routing_event( "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( entry["SYSTEM"], int_id(stream_id), - int_id(_end_peer), + _end_report_peer, int_id(rf_src), entry_ts, int_id(entry_tgid_b), diff --git a/tests/application/test_dynamic_tg_persist.py b/tests/application/test_dynamic_tg_persist.py index 09e191e..989dfce 100644 --- a/tests/application/test_dynamic_tg_persist.py +++ b/tests/application/test_dynamic_tg_persist.py @@ -26,6 +26,7 @@ from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases from adn_server.application.routing.helpers import register_peer_ua_session from adn_server.domain.dynamic_tg import DynamicTgEntry from tests.application.test_peer_single_downlink import _peer_id, _sys_cfg +from twisted.internet import defer class _CaptureStore: @@ -66,3 +67,48 @@ def test_persist_after_register_stores_absolute_expires_from_memory() -> None: entry = store.replaced[0] assert entry.expires_at == now + 60.0 * 60.0 assert entry.single_mode is True + + +def test_restore_peer_passes_now_as_keyword_to_on_restored() -> None: + """sync_restored_dynamic_tgs declares ``now`` keyword-only; restore must match.""" + import time + + peer_id = _peer_id() + sys_cfg = _sys_cfg() + now = time.time() + entry = DynamicTgEntry( + int_id=730039101, + system_name="MASTER-A", + slot=2, + tgid=7304, + single_mode=True, + expires_at=now + 3600.0, + updated_at=now, + ) + seen: dict[str, float] = {} + + def on_restored( + _peer_id: bytes, + _system_name: str, + _sys_cfg: dict, + entries: list[DynamicTgEntry], + *, + now: float, + ) -> None: + seen["now"] = now + seen["count"] = float(len(entries)) + + class _Store: + def load_peer(self, int_id: int, system_name: str): + return defer.succeed([entry]) + + def purge_expired(self, now: float) -> None: + pass + + uc = DynamicTgUseCases(_Store(), on_restored=on_restored) + out: list[list[int]] = [] + d = uc.restore_peer(peer_id, "MASTER-A", sys_cfg) + d.addCallback(lambda tgids: out.append(tgids)) + assert out == [[7304]] + assert seen["count"] == 1.0 + assert isinstance(seen.get("now"), float) diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py index 3770174..74061b1 100644 --- a/tests/application/test_monitor_topology.py +++ b/tests/application/test_monitor_topology.py @@ -142,6 +142,13 @@ def test_expand_inject_proxy_emits_all_virtual_masters() -> None: {"parts": {3: "SYSTEM-4", 5: "9990"}}, id="tx_echo_keeps_9990_for_hotspot_rx", ), + pytest.param( + "GROUP VOICE,START,TX,SYSTEM,4100887026,730039101,730039101,2,9990", + [(730039101,)], + {730039101: 4}, + {"parts": {3: "SYSTEM-4", 5: "9990"}}, + id="tx_echo_peer_id_in_field5_dst_9990", + ), pytest.param( "GROUP VOICE,START,RX,SYSTEM,4100887026,73003,7300392,2,9990", [(730039101,)], diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index 7de9885..6d24376 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -388,6 +388,27 @@ def test_register_peer_ua_session_ignores_4000() -> None: assert peer_single_exclusive_tgid(peer_single, 2, sys_cfg, peer_id=peer_id, now=now) is None +def test_register_peer_ua_session_ignores_echo_9990() -> None: + """Echo TG 9990 must not create SINGLE=1 exclusive listen lock.""" + peer_id = _peer_id() + sys_cfg = _sys_cfg() + now = 1_000_000.0 + peer = {"OPTIONS": b"TS2=730444;SINGLE=1;TIMER=5;"} + register_peer_ua_session(peer, peer_id, 2, 730444, sys_cfg, now=now) + register_peer_ua_session(peer, peer_id, 2, 9990, sys_cfg, now=now) + assert peer_single_exclusive_tgid(peer, 2, sys_cfg, peer_id=peer_id, now=now) == 730444 + + +def test_register_peer_ua_session_ignores_service_999x() -> None: + """On-demand / service 9991–9999 are not UA sessions (same class as echo 9990).""" + peer_id = _peer_id() + sys_cfg = _sys_cfg() + now = 1_000_000.0 + peer = {"OPTIONS": b"TS2=730444;SINGLE=1;TIMER=5;"} + register_peer_ua_session(peer, peer_id, 2, 9999, sys_cfg, now=now) + assert peer_single_exclusive_tgid(peer, 2, sys_cfg, peer_id=peer_id, now=now) is None + + def test_new_tx_replaces_single_session_tg() -> None: peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"} sys_cfg = _sys_cfg()