fix(plugins): plugin voice records RX state like a hotspot frame

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 <noreply@anthropic.com>
pull/106/head
yo 4 days ago
parent 1ab9036345
commit 0e61d5d473

@ -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

@ -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)

Loading…
Cancel
Save

Powered by TurnKey Linux.