From 84ccbc8f5c2e9586bcd6e7f532cbc39b561ad19f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Tue, 9 Jun 2026 15:39:30 -0400 Subject: [PATCH] fix(talker-alias): passthrough from radio, dedupe relay and logs Buffer DMRA before VHEAD, wait for hotspot TA in both mode, relay passthrough DMRA once per stream, and dedupe DEBUG policy/embed lines. --- src/adn_server/application/bridge/lc_ta.py | 70 +++++++++++++++-- .../application/bridge_use_cases.py | 9 +-- .../application/talker_alias_use_cases.py | 75 ++++++++++++++----- .../twisted_adapters/udp_hbp.py | 27 ++++++- tests/talker_alias/test_bridge_inject.py | 3 +- tests/talker_alias/test_dmra_early_buffer.py | 65 ++++++++++++++++ tests/talker_alias/test_format.py | 11 ++- tests/talker_alias/test_relay_dedupe.py | 55 ++++++++++++++ 8 files changed, 281 insertions(+), 34 deletions(-) create mode 100644 tests/talker_alias/test_dmra_early_buffer.py create mode 100644 tests/talker_alias/test_relay_dedupe.py diff --git a/src/adn_server/application/bridge/lc_ta.py b/src/adn_server/application/bridge/lc_ta.py index d1537f9..0a756e8 100644 --- a/src/adn_server/application/bridge/lc_ta.py +++ b/src/adn_server/application/bridge/lc_ta.py @@ -37,6 +37,10 @@ class BridgeLcTaMixin: "OPENBRIDGE", ) + def _master_ta_wait_source(self, source_system: str) -> bool: + """``both`` mode waits for hotspot TA only on HBP MASTER sources (not OBP).""" + return self._config.get("SYSTEMS", {}).get(source_system, {}).get("MODE") == "MASTER" + def _cancel_both_ta_wait(self, source_system: str, stream_id: bytes) -> None: wait = self._both_ta_wait.pop(self._both_ta_key(source_system, stream_id), None) if wait and wait.get("timer") and getattr(wait["timer"], "cancel", None): @@ -109,6 +113,8 @@ class BridgeLcTaMixin: continue if st.get("TX_TA_ON"): continue + if st.get("_ta_embed_kind") == "passthrough" and not force_inject: + continue if st.get("TX_STREAM_ID") == stream_id and st.get("TX_RFS") == rf_src: target = st.get("_ta_target_system", source_system) self._init_talker_alias_embed( @@ -137,9 +143,10 @@ class BridgeLcTaMixin: blocks = self._get_stream_dmra_blocks(source_system, stream_id) if not blocks or not passthrough_complete(blocks): return - key = self._both_ta_key(source_system, stream_id) - wait = self._both_ta_wait.get(key) - if wait: + relay_key = self._both_ta_key(source_system, stream_id) + already_relayed = relay_key in self._passthrough_relayed + wait = self._both_ta_wait.get(relay_key) + if wait and not already_relayed: wait["peer"] = peer_id self._cancel_both_ta_wait(source_system, stream_id) for target_system in wait["targets"]: @@ -147,7 +154,9 @@ class BridgeLcTaMixin: source_system, target_system, rf_src, stream_id, peer_id, force=True, ) self._apply_both_ta_embed(source_system, rf_src, stream_id) - self._relay_passthrough_dmra(source_system, peer_id, rf_src, stream_id) + if not already_relayed: + self._relay_passthrough_dmra(source_system, peer_id, rf_src, stream_id) + self._passthrough_relayed.add(relay_key) def _relay_passthrough_dmra( self, @@ -179,9 +188,39 @@ class BridgeLcTaMixin: ) src_cfg = self._config.get("SYSTEMS", {}).get(source_system, {}) if src_cfg.get("MODE") == "MASTER" and src_cfg.get("REPEAT", True): - self._talker_alias.clear_stream(source_system, stream_id) self.send_talker_alias_local_repeat(source_system, peer_id, rf_src, stream_id) + def _dispatch_talker_alias_on_bridge_open( + self, + target_st: dict[str, Any], + source_system: str, + target_system: str, + rf_src: bytes, + stream_id: bytes, + source_peer: bytes, + ) -> None: + """Prepare TA on a new bridged HBP leg (VHEAD or first burst after hangtime).""" + settings = talker_alias_settings(self._config, source_system) + if not settings["enabled"]: + return + blocks = self._get_stream_dmra_blocks(source_system, stream_id) + have_source_ta = bool(blocks and passthrough_complete(blocks)) + if settings["mode"] == "both" and self._master_ta_wait_source(source_system) and not have_source_ta: + self._register_both_ta_wait( + source_system, target_system, rf_src, stream_id, source_peer, + ) + return + self._init_talker_alias_embed( + target_st, + source_system, + target_system, + rf_src, + stream_id, + ) + self._send_talker_alias_to_target( + source_system, target_system, rf_src, stream_id, source_peer, + ) + def _send_talker_alias_to_target( self, source_system: str, @@ -201,8 +240,12 @@ class BridgeLcTaMixin: return blocks = self._get_stream_dmra_blocks(source_system, stream_id) have_passthrough = bool(blocks and passthrough_complete(blocks)) - if have_passthrough and force: - self._talker_alias.clear_stream(target_system, stream_id) + if have_passthrough: + if force: + if not self._talker_alias.should_resend_passthrough_dmra(target_system, stream_id): + return + elif not self._talker_alias.should_send_on_vhead(target_system, stream_id): + return elif not self._talker_alias.should_send_on_vhead(target_system, stream_id): return if not have_passthrough: @@ -224,7 +267,11 @@ class BridgeLcTaMixin: except Exception as e: logger.warning("(ROUTER) send_dmra_to_system %s failed: %s", target_system, e) return - self._talker_alias.mark_dmra_sent(target_system, stream_id) + self._talker_alias.mark_dmra_sent( + target_system, + stream_id, + kind="passthrough" if have_passthrough else "inject", + ) sid = int_id(stream_id) if peer_count: logger.debug( @@ -311,6 +358,7 @@ class BridgeLcTaMixin: def clear_talker_alias_stream(self, system_name: str, stream_id: bytes) -> None: """Release per-stream TA dedupe state after VTERM.""" self._cancel_both_ta_wait(system_name, stream_id) + self._passthrough_relayed.discard(self._both_ta_key(system_name, stream_id)) self._talker_alias.clear_stream(system_name, stream_id) if not self._get_protocols: return @@ -375,6 +423,11 @@ class BridgeLcTaMixin: st["TX_TA_PHASE"] = 0 # First B1–B4 cycle carries group LC; TA on the next cycle. st["TX_TA_ON"] = False + blocks = self._get_stream_dmra_blocks(source_system, stream_id) + if blocks and passthrough_complete(blocks) and not force_inject: + st["_ta_embed_kind"] = "passthrough" + else: + st["_ta_embed_kind"] = "inject" def _rewrite_embed_lc( self, @@ -412,4 +465,5 @@ class BridgeLcTaMixin: st.pop("TX_TA_PHASE", None) st.pop("TX_TA_BLOCK_COUNT", None) st.pop("TX_TA_ON", None) + st.pop("_ta_embed_kind", None) diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index c19d184..6120f2a 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -89,6 +89,8 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi self._talker_alias = TalkerAliasUseCases(config, ta_emblc_encoder=ta_emblc_encoder) # (source_system, stream_id) -> {rf_src, peer, targets, timer} self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {} + # Passthrough DMRA/embed relay already applied for this source stream. + self._passthrough_relayed: set[tuple[str, bytes]] = set() def _send_bridge_event(self, event: str | bytes) -> bool: """Send BRDG_EVENT via ReportingUseCases when REPORT is enabled.""" @@ -562,12 +564,13 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi _ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc) _ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc) _ts_st["TX_EMB_LC"] = self._encode_emblc(dst_lc) - self._init_talker_alias_embed( + self._dispatch_talker_alias_on_bridge_open( _ts_st, system_name, entry["SYSTEM"], rf_src, stream_id, + peer_id, ) logger.info( "(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s", @@ -583,10 +586,6 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi int_id(entry_tgid_b), ) ) - # First successful forward to HBP (may not be VHEAD if earlier frames were hangtime-blocked). - self._send_talker_alias_to_target( - system_name, entry["SYSTEM"], rf_src, stream_id, peer_id, - ) _ts_st["TX_TIME"] = pkt_time _ts_st["TX_TYPE"] = dtype_vseq # Slot bit rewrite (legacy bridge.py 457-460 / 770-773) diff --git a/src/adn_server/application/talker_alias_use_cases.py b/src/adn_server/application/talker_alias_use_cases.py index 1b270e0..1943eb1 100644 --- a/src/adn_server/application/talker_alias_use_cases.py +++ b/src/adn_server/application/talker_alias_use_cases.py @@ -133,17 +133,60 @@ class TalkerAliasUseCases: self._config = config self._ta_emblc = ta_emblc_encoder self._sent_streams: set[tuple[str, bytes]] = set() - self._embed_logged: set[tuple[str, bytes]] = set() + self._sent_kind: dict[tuple[str, bytes], str] = {} + self._embed_logged: set[tuple[str, str, bytes]] = set() + self._policy_logged: set[tuple[str, str, bytes, str]] = set() + + def clear_dmra_sent(self, system_name: str, stream_id: bytes) -> None: + key = (system_name, stream_id) + self._sent_streams.discard(key) + self._sent_kind.pop(key, None) def clear_stream(self, system_name: str, stream_id: bytes) -> None: - self._sent_streams.discard((system_name, stream_id)) - self._embed_logged.discard((system_name, stream_id)) + self.clear_dmra_sent(system_name, stream_id) + self._embed_logged = {k for k in self._embed_logged if k[2] != stream_id} + self._policy_logged = {k for k in self._policy_logged if k[2] != stream_id} def should_send_on_vhead(self, target_system: str, stream_id: bytes) -> bool: return (target_system, stream_id) not in self._sent_streams - def mark_dmra_sent(self, target_system: str, stream_id: bytes) -> None: - self._sent_streams.add((target_system, stream_id)) + def should_resend_passthrough_dmra(self, target_system: str, stream_id: bytes) -> bool: + """True when passthrough DMRA should replace a prior inject on this leg.""" + key = (target_system, stream_id) + if key not in self._sent_streams: + return True + return self._sent_kind.get(key) == "inject" + + def mark_dmra_sent(self, target_system: str, stream_id: bytes, *, kind: str) -> None: + key = (target_system, stream_id) + self._sent_streams.add(key) + self._sent_kind[key] = kind + + def _log_policy_once( + self, + source_system: str, + target_system: str, + stream_id: bytes, + kind: str, + text: str, + *, + suffix: str = "", + via: str, + ) -> None: + log_key = (source_system, target_system, stream_id, kind) + if log_key in self._policy_logged: + return + self._policy_logged.add(log_key) + if kind == "passthrough": + logger.debug( + "(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s", + source_system, text, via, target_system, int_id(stream_id), + ) + else: + logger.debug( + "(%s) *TALKER ALIAS* inject '%s'%s via %s -> %s stream %s", + source_system, text, suffix, via, target_system, int_id(stream_id), + ) def packets_for_stream( self, @@ -168,16 +211,14 @@ class TalkerAliasUseCases: if not have_passthrough: return None text = decode_ta_from_blocks(blocks) - logger.debug( - "(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s", - source_system, text, via, target, int_id(stream_id), + self._log_policy_once( + source_system, target, stream_id, "passthrough", text, via=via, ) return passthrough_packets_from_blocks(rf_src, blocks) if mode == "inject": text = format_talker_alias_text(self._config, rf_src) - logger.debug( - "(%s) *TALKER ALIAS* inject '%s' via %s -> %s stream %s", - source_system, text, via, target, int_id(stream_id), + self._log_policy_once( + source_system, target, stream_id, "inject", text, via=via, ) return build_dmra_packets(rf_src, text, settings["text_formats"][0]) # both: prefer the source's own TA. If a valid MMDVM DMRA buffer arrived, @@ -185,17 +226,15 @@ class TalkerAliasUseCases: # through unchanged in the DMRD voice, so do NOT inject a template here. if have_passthrough: text = decode_ta_from_blocks(blocks) - logger.debug( - "(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s", - source_system, text, via, target, int_id(stream_id), + self._log_policy_once( + source_system, target, stream_id, "passthrough", text, via=via, ) return passthrough_packets_from_blocks(rf_src, blocks) # Legacy resolve_ta (both): inject template when no DMRA buffer yet (VHEAD). text = format_talker_alias_text(self._config, rf_src) suffix = " (no source TA yet)" if fallback_inject else "" - logger.debug( - "(%s) *TALKER ALIAS* inject '%s'%s via %s -> %s stream %s", - source_system, text, suffix, via, target, int_id(stream_id), + self._log_policy_once( + source_system, target, stream_id, "inject", text, suffix=suffix, via=via, ) return build_dmra_packets(rf_src, text, settings["text_formats"][0]) @@ -225,7 +264,7 @@ class TalkerAliasUseCases: mode = settings["mode"] target = target_system or source_system via = "repeat" if source_system == target else "bridge" - log_key = (source_system, stream_id) + log_key = (source_system, target, stream_id) def _log(text: str, suffix: str) -> None: if log_key in self._embed_logged: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 8db9f0a..d3cbf7b 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -318,6 +318,30 @@ class HBPProtocol(DatagramProtocol): if not self._ta_buffer_enabled() or not stream_id: return self._dmra_rf_stream[(peer_id, rf_src)] = stream_id + self._promote_dmra_provisional_key(rf_src, stream_id) + + def _promote_dmra_provisional_key(self, provisional: bytes, stream_id: bytes) -> None: + """Move TA blocks buffered under ``rf_src`` (pre-VHEAD) to the real stream id.""" + if provisional == stream_id: + return + entry = self._dmra_by_stream.pop(provisional, None) + if not entry: + return + target = self._dmra_by_stream.setdefault( + stream_id, + { + "blocks": {}, + "rf_src": entry.get("rf_src", b""), + "peer": entry.get("peer", b""), + "last": entry.get("last", time.time()), + }, + ) + target["blocks"].update(entry.get("blocks", {})) + target["last"] = max(float(target.get("last", 0)), float(entry.get("last", 0))) + if entry.get("rf_src"): + target["rf_src"] = entry["rf_src"] + if entry.get("peer"): + target["peer"] = entry["peer"] def store_dmra_packet(self, peer_id: bytes, data: bytes) -> None: """Buffer one DMRA block from a hotspot (MASTER receive path).""" @@ -327,7 +351,8 @@ class HBPProtocol(DatagramProtocol): rf_src, block_id, payload = parsed stream_id = self._dmra_rf_stream.get((peer_id, rf_src)) if not stream_id: - return + # Legacy hblink: DMRA may arrive before the first DMRD; key by rf_src until then. + stream_id = rf_src now = time.time() entry = self._dmra_by_stream.setdefault( stream_id, diff --git a/tests/talker_alias/test_bridge_inject.py b/tests/talker_alias/test_bridge_inject.py index 74080ca..18cbd1b 100644 --- a/tests/talker_alias/test_bridge_inject.py +++ b/tests/talker_alias/test_bridge_inject.py @@ -83,7 +83,8 @@ def test_both_mode_without_ta_still_rewrites_group_embedded_lc() -> None: ) ts_st = scenario.protocols["MASTER-B"].STATUS[2] - assert ts_st.get("TX_TA_EMB") is not None # both mode injects template at VHEAD + # both + MASTER source: TA overlay deferred until source TA or 2s fallback (see docs). + assert ts_st.get("TX_TA_EMB") is None bursts = [p for p in scenario.capture.for_system("MASTER-B") if p.fields["dtype_vseq"] == 1] assert bursts, "voice burst B should be forwarded to MASTER-B" bits = bitarray(endian="big") diff --git a/tests/talker_alias/test_dmra_early_buffer.py b/tests/talker_alias/test_dmra_early_buffer.py new file mode 100644 index 0000000..c61ea90 --- /dev/null +++ b/tests/talker_alias/test_dmra_early_buffer.py @@ -0,0 +1,65 @@ +"""DMRA may arrive before DMRD VHEAD — legacy buffers under rf_src until stream note.""" + +from __future__ import annotations + +from adn_server.domain import bytes_3, bytes_4 +from adn_server.domain.talker_alias import build_dmra_packet, decode_ta_from_blocks +from adn_server.infrastructure.hbp_constants import DMRA + +from tests.support.hbp_repeat_stack import build_hbp_repeat_stack +from tests.harness.deterministic import DeterministicScenario, PacketSpec +from tests.talker_alias.test_mmdvm_wire import mmdvm_wire_blocks + + +_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 test_dmra_before_vhead_passthrough_on_repeat() -> None: + """Regression: early DMRA was dropped when stream_id mapping did not exist yet.""" + stack = build_hbp_repeat_stack(talker_alias=True) + stack.config["GLOBAL"]["TALKER_ALIAS_MODE"] = "passthrough" + stack.register_peer(_PEER_TX, _ADDR_TX) + stack.register_peer(_PEER_RX, _ADDR_RX) + + rf_src = bytes_3(3120001) + text = "CE5RPY Radio" + for block_id, payload in mmdvm_wire_blocks(text).items(): + stack.inject(build_dmra_packet(rf_src, block_id, payload), _ADDR_TX) + + base = PacketSpec(peer_id=730039210, rf_src=3120001, dst_id=7304, slot=2, stream_id=0xA1B2C3D4) + stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX) + + assert stack.dmra_capture, "passthrough DMRA should be sent on VHEAD" + packets, exclude = stack.dmra_capture[0] + assert exclude == _PEER_TX + assert all(p[:4] == DMRA for p in packets) + rebuilt = {i: packets[i][8:15] for i in range(len(packets))} + assert decode_ta_from_blocks(rebuilt) == text + + dmra_to_rx = [p for p, _ in stack.transport.sent if p[:4] == DMRA and _ == _ADDR_RX] + assert dmra_to_rx, "listening peer should receive passthrough DMRA" + + +def test_promote_provisional_dmra_buffer() -> None: + stack = build_hbp_repeat_stack(talker_alias=True) + hbp = stack.hbp + peer = _PEER_TX + rf_src = bytes_3(3120001) + stream_id = bytes_4(0x12345678) + blocks = mmdvm_wire_blocks("TA early") + + for block_id, payload in blocks.items(): + hbp.store_dmra_packet(peer, build_dmra_packet(rf_src, block_id, payload)) + + assert hbp.get_dmra_blocks(rf_src) is not None + assert hbp.get_dmra_blocks(stream_id) is None + + hbp.note_dmrd_stream(peer, rf_src, stream_id) + + promoted = hbp.get_dmra_blocks(stream_id) + assert promoted is not None + assert decode_ta_from_blocks(promoted) == "TA early" + assert hbp.get_dmra_blocks(rf_src) is None diff --git a/tests/talker_alias/test_format.py b/tests/talker_alias/test_format.py index a29dc6d..a55e30d 100644 --- a/tests/talker_alias/test_format.py +++ b/tests/talker_alias/test_format.py @@ -21,9 +21,18 @@ def test_talker_alias_clear_stream_allows_resend() -> None: key = ("MASTER-B", stream_id) assert ta.should_send_on_vhead("MASTER-B", stream_id) is True - ta.mark_dmra_sent("MASTER-B", stream_id) + ta.mark_dmra_sent("MASTER-B", stream_id, kind="inject") assert ta.should_send_on_vhead("MASTER-B", stream_id) is False ta.clear_stream("MASTER-B", stream_id) assert key not in ta._sent_streams assert ta.should_send_on_vhead("MASTER-B", stream_id) is True + + +def test_clear_dmra_sent_preserves_embed_log_dedupe() -> None: + ta = make_talker_alias_use_cases(talker_alias_config()) + stream_id = PacketSpec(stream_id=0x93939393).data()[16:20] + ta._embed_logged.add(("MASTER-A", "MASTER-B", stream_id)) + + ta.clear_dmra_sent("MASTER-B", stream_id) + assert ("MASTER-A", "MASTER-B", stream_id) in ta._embed_logged diff --git a/tests/talker_alias/test_relay_dedupe.py b/tests/talker_alias/test_relay_dedupe.py new file mode 100644 index 0000000..49c3d3e --- /dev/null +++ b/tests/talker_alias/test_relay_dedupe.py @@ -0,0 +1,55 @@ +"""Talker Alias DMRA relay and log deduplication.""" + +from __future__ import annotations + +from adn_server.application.bridge_use_cases import BridgeUseCases +from adn_server.domain import bytes_3, bytes_4 +from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter +from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder +from adn_server.domain.dmr.bptc import encode_emblc + +from tests.harness.scenarios import make_talker_alias_use_cases, talker_alias_config +from tests.talker_alias.test_mmdvm_wire import mmdvm_wire_blocks + + +def test_should_resend_passthrough_only_after_inject() -> None: + ta = make_talker_alias_use_cases(talker_alias_config()) + stream_id = bytes_4(0xA1B2C3D4) + + assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is True + ta.mark_dmra_sent("SYSTEM", stream_id, kind="inject") + assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is True + ta.mark_dmra_sent("SYSTEM", stream_id, kind="passthrough") + assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is False + + +def test_on_dmra_fragment_stored_relays_passthrough_once() -> None: + """Second completion callback must not re-send passthrough DMRA on the same stream.""" + config = talker_alias_config() + config["GLOBAL"]["TALKER_ALIAS_MODE"] = "both" + config["SYSTEMS"]["MASTER-A"]["REPEAT"] = True + blocks = mmdvm_wire_blocks("CE5RPY") + sent: list[str] = [] + + def _send_dmra(target_system: str, packets: list[bytes], exclude_peer: bytes | None = None) -> int: + sent.append(target_system) + return 1 + + bridge = BridgeUseCases( + InMemoryBridgeRouter(), + config, + send_dmra_to_system=_send_dmra, + get_dmra_blocks=lambda _sys, _sid: blocks, + encode_emblc=encode_emblc, + ta_emblc_encoder=default_ta_emblc_encoder, + ) + peer = bytes_4(1001) + rf_src = bytes_3(3120001) + stream_id = bytes_4(0xB1B2B3B4) + + bridge.on_dmra_fragment_stored("MASTER-A", peer, rf_src, stream_id) + first_count = len(sent) + assert first_count == 1 + + bridge.on_dmra_fragment_stored("MASTER-A", peer, rf_src, stream_id) + assert len(sent) == first_count