From 0e61d5d473d028fffba6b3971f816c5a5aee7f58 Mon Sep 17 00:00:00 2001 From: yo Date: Fri, 25 Sep 2026 00:09:17 +0200 Subject: [PATCH] fix(plugins): plugin voice records RX state like a hotspot frame MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Found on a live master (2131): a plugin beacon on TG 213 reached the OpenBridges with 2 frames out of 115. PluginIngress marked the ingress MASTER slot as TX of the stream, while its RX fields still held the last radio's call. From the second frame on, group_voice_tg_ingress_collision saw TG 213 active on that slot (our own TX mark) with a stream and a source that were not ours (the stale RX ones), took it for a QSO and suppressed the uplink: only the MASTER's own hotspots heard the rest. A hotspot frame doesn't have this: udp_hbp records the slot's RX state (RX_STREAM_ID, RX_LC, RX_RFS, RX_TGID, RX_TYPE, RX_TIME…) once routing accepts it, so the next frame is the same stream, not a new one. PluginIngress now does exactly that instead of the TX mark. The slot still reads busy for routed voice while the stream plays (its RX leg), a radio or another stream still makes the next frame fail, and the terminator or the 1 s silence sweep sets RX_TYPE back to VTERM. Regression test: stale RX state on the slot, a whole beacon forwarded, header and embedded LC both the beacon's (the LC itself was never affected: routing builds each leg's LC from the target TG and rf_src). Co-Authored-By: Claude Opus 5.5 --- .../plugins/application/ingress.py | 60 ++++++++++++++----- tests/routing/test_plugin_send_routing.py | 42 +++++++++++-- 2 files changed, 84 insertions(+), 18 deletions(-) diff --git a/src/adn_server/application/plugins/application/ingress.py b/src/adn_server/application/plugins/application/ingress.py index 2d9e02c..5fa0d34 100644 --- a/src/adn_server/application/plugins/application/ingress.py +++ b/src/adn_server/application/plugins/application/ingress.py @@ -26,7 +26,10 @@ import logging import time from typing import Any, Callable -from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM +from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, int_id +from ....domain.dmr import decode +from ....domain.dmr.const import LC_OPT +from ....domain.hbp_protocol import STREAM_TO from ....domain.mesh_engine import server_id_bytes from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd from ...routing.helpers import master_dynamic_tg_slots, slot_voice_held_by_other_stream @@ -44,9 +47,11 @@ class PluginIngress: Group voice is routed like a scheduled announcement: through the bridges (OpenBridge legs included, as for any local ingress), and out to the hotspots of the MASTER it - enters on. While a plugin's stream plays it holds that MASTER slot (TX_TYPE=VHEAD, - TX_STREAM_ID, TX_RFS), so routed voice finds it busy; a radio or another stream on the - slot makes the frame fail, which tells the plugin to stop. The terminator frees it. + enters on. Each accepted frame records the slot's RX state exactly as ``udp_hbp`` does + for a hotspot's frame (RX_STREAM_ID, RX_LC, RX_RFS, RX_TGID, RX_TYPE, RX_TIME…): routing + then knows the stream it is continuing and the LC it carries, and routed voice finds + the slot busy. A radio or another stream on the slot makes the frame fail, which tells + the plugin to stop. The terminator frees the slot. """ def __init__( @@ -131,18 +136,12 @@ class PluginIngress: slot = getattr(proto, "STATUS", {}).get(header.slot) if proto is not None else None if slot is None: return False - if slot_voice_held_by_other_stream(slot, header.stream_id, now) or ( - self._route(master, pkt, now, server_id, plugin) is not True - ): + if _slot_taken(slot, header.stream_id, now) or self._route(master, pkt, now, server_id, plugin) is not True: self._release(slot, header.stream_id) self._off_air(header.stream_id, now) return False is_term = header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VTERM - slot["TX_TYPE"] = HBPF_SLT_VTERM if is_term else HBPF_SLT_VHEAD - slot["TX_STREAM_ID"] = header.stream_id - slot["TX_RFS"] = header.rf_src - slot["TX_TGID"] = header.dst_id - slot["TX_TIME"] = now + _record_rx(slot, header, pkt, server_id, now) self._send_local(master, pkt) if header.stream_id not in self._on_air: self._on_air_start(master, header, now) @@ -191,8 +190,41 @@ class PluginIngress: @staticmethod def _release(slot: dict[str, Any], stream_id: bytes) -> None: - if slot.get("TX_STREAM_ID") == stream_id: - slot["TX_TYPE"] = HBPF_SLT_VTERM + if slot.get("RX_STREAM_ID") == stream_id: + slot["RX_TYPE"] = HBPF_SLT_VTERM def _server_id(self) -> bytes: return server_id_bytes(self._config.get("GLOBAL", {}).get("SERVER_ID"))[:4] + + +def _slot_taken(slot: dict[str, Any], stream_id: bytes, now: float) -> bool: + """Live voice on the slot that is not this stream: a radio, or another stream.""" + for leg in ("RX", "TX"): + leg_type = slot.get(f"{leg}_TYPE") + if leg_type is None or leg_type == HBPF_SLT_VTERM: + continue + if slot.get(f"{leg}_STREAM_ID") == stream_id: + continue + if now - float(slot.get(f"{leg}_TIME", 0) or 0) < STREAM_TO: + return True + return False + + +def _record_rx(slot: dict[str, Any], header: Any, pkt: bytes, peer_id: bytes, now: float) -> None: + """The slot RX state ``udp_hbp`` records after routing accepts a hotspot's frame.""" + if header.stream_id != slot.get("RX_STREAM_ID"): + slot["RX_START"] = now + lc = LC_OPT + header.dst_id + header.rf_src + if header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VHEAD: + try: + lc = decode.voice_head_term(pkt[20:53])["LC"] + except Exception: + pass + slot["RX_LC"] = lc + slot["RX_PEER"] = peer_id + slot["RX_SEQ"] = header.seq + slot["RX_RFS"] = header.rf_src + slot["RX_TYPE"] = header.dtype_vseq + slot["RX_TGID"] = header.dst_id + slot["RX_TIME"] = now + slot["RX_STREAM_ID"] = header.stream_id diff --git a/tests/routing/test_plugin_send_routing.py b/tests/routing/test_plugin_send_routing.py index cf68e3a..425a8a6 100644 --- a/tests/routing/test_plugin_send_routing.py +++ b/tests/routing/test_plugin_send_routing.py @@ -35,6 +35,7 @@ from tests.routing.unit_data_helpers import idle_hbp_slot from adn_server.application.plugins.application.bus import PluginBus from adn_server.application.plugins.application.data_bridge import DataPluginBridge from adn_server.application.plugins.application.ingress import PluginIngress +from adn_server.application.routing.helpers import hbp_slot_blocks_group_voice from adn_server.application.plugins.domain.events import UnitDataFrame from adn_server.application.routing.announcement_ptt_inject import inject_plugin_dmrd from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_3, bytes_4 @@ -170,9 +171,11 @@ def test_the_master_slot_is_held_while_it_plays_and_freed_by_the_terminator() -> frames = _beacon() _play(sc, ingress, frames[:3]) slot = sc.protocols["SYSTEM"].STATUS[2] - assert slot["TX_TYPE"] == HBPF_SLT_VHEAD and slot["TX_RFS"] == bytes_3(BEACON_ID) + assert slot["RX_RFS"] == bytes_3(BEACON_ID) and slot["RX_TYPE"] != HBPF_SLT_VTERM + # routed voice for another TG finds the slot busy while the beacon plays + assert hbp_slot_blocks_group_voice(slot, bytes_3(214), b"\x07" * 4, sc.clock.time(), 0) _play(sc, ingress, frames[3:]) - assert slot["TX_TYPE"] == HBPF_SLT_VTERM + assert slot["RX_TYPE"] == HBPF_SLT_VTERM def test_a_radio_talking_on_the_slot_stops_the_beacon() -> None: @@ -234,7 +237,8 @@ def test_a_beacon_cut_by_a_radio_still_ends_on_the_monitor() -> None: sc.protocols["SYSTEM"].STATUS[2].update(RX_TYPE=HBPF_SLT_VHEAD, RX_TIME=sc.clock.time() + 0.06, RX_STREAM_ID=b"\x09" * 4) assert _play(sc, ingress, frames[3:4]) == [False] assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1 - assert sc.protocols["SYSTEM"].STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM + slot = sc.protocols["SYSTEM"].STATUS[2] + assert slot["RX_STREAM_ID"] == b"\x09" * 4 and slot["RX_TYPE"] == HBPF_SLT_VHEAD # the radio's state is left alone def test_a_stream_the_plugin_abandons_ends_after_a_silent_second() -> None: @@ -246,4 +250,34 @@ def test_a_stream_the_plugin_abandons_ends_after_a_silent_second() -> None: sc.clock.advance(1.5) assert ingress.voice_slot_for_tg(TG) == 2 # the next query sweeps it assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1 - assert sc.protocols["SYSTEM"].STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM + assert sc.protocols["SYSTEM"].STATUS[2]["RX_TYPE"] == HBPF_SLT_VTERM # released by the sweep + + +def test_a_beacon_on_a_slot_with_stale_rx_state_is_forwarded_whole_with_its_own_lc() -> None: + """Regression (2131, 25-sep-2026): the ingress slot still held the last radio's RX + state. From the second frame on, routing took the beacon for a colliding QSO, + suppressed the uplink, and forwarded only 2 frames; the LC came from RX_LC.""" + from adn_server.domain.dmr import decode + + sc, ingress, _ = _voice_scenario() + sc.protocols["SYSTEM"].STATUS[2].update( + RX_TYPE=HBPF_SLT_VTERM, RX_STREAM_ID=b"\x0e" * 4, RX_RFS=bytes_3(3120001), RX_TGID=bytes_3(214), + RX_PEER=bytes_4(1001), RX_TIME=sc.clock.time() - 30, RX_LC=b"\x00\x00\x20" + bytes_3(214) + bytes_3(3120001), + ) + frames = _beacon() + assert _play(sc, ingress, frames) == [True] * len(frames) + to_b = sc.capture.for_system("SYSTEM-B") + assert len(to_b) == len(frames) + lc = decode.voice_head_term(to_b[0].packet[20:53])["LC"] + assert lc[3:6] == bytes_3(TG) and lc[6:9] == bytes_3(BEACON_ID) + # the embedded LC of the voice bursts B-E must be the beacon's too, not the stale RX_LC + from bitarray import bitarray + + from adn_server.domain.dmr import bptc + + by_vseq = {p.packet[15] & 0x0F: p.packet for p in to_b[1:-1]} + frags = bitarray(endian="big") + for vseq in (1, 2, 3, 4): + frags += decode.voice(by_vseq[vseq][20:53])["EMBED"] + emb = bptc.decode_emblc(frags) + assert emb[3:6] == bytes_3(TG) and emb[6:9] == bytes_3(BEACON_ID)