From e87332a1d66b7eac36ec16f179a577f958ccedbf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 26 Jun 2026 01:43:14 -0400 Subject: [PATCH] fix: hard one-QSO slot rule for all peers and stale session cleanup Remove the lab-witness exemption so every hotspot obeys one group QSO per RF slot. Expire same-TG zombie peer_voice_slots after STREAM_TO, allow orphan VTERM through the slot gate, and route ingress slot tracking via track_peer_group_dmrd for consistent session lifecycle. --- .../application/routing/downlink.py | 49 ++++++--- src/adn_server/application/routing/helpers.py | 100 ++++++++++++------ .../twisted_adapters/udp_hbp.py | 49 ++++----- tests/application/test_slot_contention.py | 12 +-- 4 files changed, 129 insertions(+), 81 deletions(-) diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index 8b886a4..e779e03 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -41,6 +41,7 @@ from .helpers import ( parse_dmrd_route_fields, peer_downlink_voice_slot, peer_is_simplex, + peer_options_static_tg_slot, peer_receives_group_tgid, peer_should_receive_group_voice, peer_single_exclusive_tgid, @@ -115,12 +116,6 @@ def peer_listen_slots(peer: dict[str, Any], tgid: int) -> list[int]: return [] -def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None: - from adn_server.application.routing.helpers import peer_options_static_tg_slot as _slot - - return _slot(peer, tgid) - - def peer_accepts_group_downlink( ctx: DownlinkContext, peer_id: bytes, @@ -217,9 +212,11 @@ def peer_slot_blocks_downlink( if not rx_listening: return True active = per_slot.get(int(listen_slot)) - if not isinstance(active, dict) or active.get("stream_id") != stream_id: - return True - return False + if not isinstance(active, dict): + return False + if active.get("stream_id") == stream_id: + return False + return True if not ctx.per_peer_contention(): return False hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0) @@ -357,15 +354,19 @@ def track_peer_group_dmrd( *, pkt_time: float | None = None, from_ingress: bool = False, + voice_slot: int | None = None, ) -> None: """Update per-hotspot slot state from downlink or ingress DMRD.""" burst = parse_dmrd_burst_fields(packet) if burst is None: return wire_slot, frame_type, dtype_vseq, stream_id, dst_id, _call_type = burst - voice_slot = peer_downlink_voice_slot( - peer, wire_slot, int_id(dst_id), ctx.sys_cfg, peer_id=peer_id, - ) + if voice_slot is None: + voice_slot = peer_downlink_voice_slot( + peer, wire_slot, int_id(dst_id), ctx.sys_cfg, peer_id=peer_id, + ) + else: + voice_slot = int(voice_slot) if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: if ( not from_ingress @@ -436,9 +437,23 @@ def track_peer_group_dmrd( ) -def build_dmra_route_packet(slot: int, tgid: int, stream_id: bytes | None = None) -> bytes: - """Synthetic DMRD for DMRA / monitor fan-out lookup.""" - return synthetic_group_dmrd_route_packet(slot, tgid, stream_id) +def peer_accepts_group_dmrd_packet( + ctx: DownlinkContext, + peer_id: bytes, + peer: dict[str, Any], + packet: bytes, + *, + routed: bool = False, +) -> bool: + """True when group/vcsbk DMRD passes per-hotspot slot gate (OPTIONS checked separately).""" + if not ctx.per_peer_contention(): + return True + route_pkt = ( + packet + if routed + else remap_dmrd_for_peer(packet, peer, ctx.sys_cfg, peer_id=peer_id) + ) + return not peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt) def peer_accepts_dmra( @@ -451,7 +466,7 @@ def peer_accepts_dmra( peer = ctx.peers.get(peer_id) if not isinstance(peer, dict): return False - route_pkt = build_dmra_route_packet(slot, tgid) + route_pkt = synthetic_group_dmrd_route_packet(slot, tgid) if not peer_accepts_group_downlink(ctx, peer_id, peer, slot, tgid): return False remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id) @@ -477,7 +492,7 @@ def peer_would_show_group_voice_on_monitor( return False if ctx is None or not ctx.per_peer_contention(): return True - route_pkt = build_dmra_route_packet(wire_slot, tgid, stream_id) + route_pkt = synthetic_group_dmrd_route_packet(wire_slot, tgid, stream_id) remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id) return not peer_slot_blocks_downlink(ctx, peer_id, peer, remapped) diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index e32bca7..9a1467b 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -271,6 +271,31 @@ def _peer_status_rx_hangtime_blocks( return (pkt_time - rx_t) < hang +def _peer_single_locked_tgid( + peer: dict[str, Any], + sys_cfg: dict[str, Any], + *, + peer_id: bytes | None = None, + now: float | None = None, + prefer_slot: int | None = None, +) -> int | None: + """First active SINGLE=1 session TG (``prefer_slot`` checked before the other TS).""" + slots: list[int] = [] + if prefer_slot is not None: + slots.append(int(prefer_slot)) + for voice_slot in (1, 2): + if prefer_slot is not None and voice_slot == int(prefer_slot): + continue + slots.append(voice_slot) + for voice_slot in slots: + locked = peer_single_exclusive_tgid( + peer, voice_slot, sys_cfg, peer_id=peer_id, now=now, + ) + if locked is not None: + return locked + return None + + def peer_single_same_tg_foreign_tx_blocks( peer: dict[str, Any], peer_id: bytes, @@ -285,13 +310,9 @@ def peer_single_same_tg_foreign_tx_blocks( if not sys_cfg or not peer_single_mode(peer, sys_cfg): return False incoming = int_id(incoming_tgid_b) - locked = None - for voice_slot in (1, 2): - locked = peer_single_exclusive_tgid( - peer, voice_slot, sys_cfg, peer_id=peer_id, now=pkt_time, - ) - if locked is not None: - break + locked = _peer_single_locked_tgid( + peer, sys_cfg, peer_id=peer_id, now=pkt_time, + ) if locked is None or int(locked) != incoming: return False if not slot_has_active_voice(slot_st, pkt_time): @@ -339,10 +360,22 @@ def peer_hotspot_voice_slot_busy( if isinstance(active, dict): if active.get("ingress"): return True - if peer_hotspot_hard_slot_one_qso(peer): - active_stream = active.get("stream_id") - if not (stream_id and active_stream and active_stream == stream_id): - return True + active_stream = active.get("stream_id") + incoming_tgid = int_id(incoming_tgid_b) + active_tgid = int(active.get("tgid", 0) or 0) + if stream_id and active_stream and active_stream == stream_id: + pass + elif ( + stream_id + and active_stream + and active_stream != stream_id + and active_tgid + and active_tgid == incoming_tgid + and (pkt_time - float(active.get("time", 0) or 0)) >= STREAM_TO + ): + peer_slots.pop(int(voice_slot), None) + else: + return True if isinstance(peer, dict) and peer_single_blocks_foreign_same_tg_downlink( peer, pk, voice_slot, incoming_tgid_b, peer_slots, sys_cfg, now=pkt_time, ): @@ -1302,18 +1335,9 @@ def peer_single_blocks_foreign_same_tg_downlink( if not sys_cfg or not peer_single_mode(peer, sys_cfg): return False incoming = int_id(incoming_tgid_b) - locked = peer_single_exclusive_tgid( - peer, voice_slot, sys_cfg, peer_id=peer_id, now=now, + locked = _peer_single_locked_tgid( + peer, sys_cfg, peer_id=peer_id, now=now, prefer_slot=voice_slot, ) - if locked is None: - for alt_slot in (1, 2): - if alt_slot == voice_slot: - continue - locked = peer_single_exclusive_tgid( - peer, alt_slot, sys_cfg, peer_id=peer_id, now=now, - ) - if locked is not None: - break if locked is None or int(locked) != incoming: return False active = (peer_slots or {}).get(int(voice_slot)) @@ -1340,13 +1364,6 @@ def peer_wants_downlink_single_listen_lock(peer: dict[str, Any], sys_cfg: dict[s return peer_static_options_tg_count(peer) <= 6 -def peer_hotspot_hard_slot_one_qso(peer: dict[str, Any] | None) -> bool: - """True when the one-QSO-per-RF-slot rule applies (normal hotspots, not lab witnesses).""" - if not isinstance(peer, dict): - return True - return peer_static_options_tg_count(peer) <= 6 - - def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool: """True when peer RPTO OPTIONS list the group TG on TS1 or TS2 (legacy REPEAT parity). @@ -1378,17 +1395,36 @@ def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None: return None -def synthetic_group_dmrd_route_packet( +def synthetic_group_dmrd_burst_packet( slot: int, tgid: int, - stream_id: bytes | None = None, + stream_id: bytes, + *, + frame_type: int = 0, + dtype_vseq: int = 0, + call_type: str = "group", ) -> bytes: - """Minimal DMRD for downlink/monitor gate lookup (slot, TG, optional stream).""" + """Minimal DMRD with burst header fields for slot-tracking helpers.""" bits = 0x80 if int(slot) == 2 else 0 + if call_type == "vcsbk": + bits |= 0x23 + else: + bits |= (int(frame_type) & 0x3) << 4 + bits |= int(dtype_vseq) & 0xF sid = bytes_4(int_id(stream_id)) if stream_id else b"\x00" * 4 return b"DMRD" + b"\x00" * 4 + bytes_3(tgid) + b"\x00" * 4 + bytes([bits]) + sid + b"\x00" * 34 +def synthetic_group_dmrd_route_packet( + slot: int, + tgid: int, + stream_id: bytes | None = None, +) -> bytes: + """Minimal DMRD for downlink/monitor gate lookup (slot, TG, optional stream).""" + sid = stream_id if stream_id else b"\x00" * 4 + return synthetic_group_dmrd_burst_packet(slot, tgid, sid) + + def peer_downlink_voice_slot( peer: dict[str, Any], wire_slot: int, diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 13ed6be..35c6d41 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -41,11 +41,10 @@ from twisted.internet.protocol import DatagramProtocol from ...application.proxy.deployment import is_proxy_inject_only from ...application.routing.downlink import ( DownlinkContext, - build_dmra_route_packet, iter_downlink_voice_slots, normalize_ua_voice_slot, peer_accepts_dmra, - peer_slot_blocks_downlink, + peer_accepts_group_dmrd_packet, remap_dmrd_for_peer, track_peer_group_dmrd, ) @@ -65,6 +64,8 @@ from ...application.routing.helpers import ( register_peer_ua_session, resolve_voice_peer_id, seed_peer_ua_session_from_status, + synthetic_group_dmrd_route_packet, + synthetic_group_dmrd_burst_packet, tg4000_reset_on_vhead, ) from ...application.routing.peer_downlink_index import ( @@ -619,24 +620,26 @@ class HBPProtocol(DatagramProtocol): peer = self._peers.get(peer_id) if peer is None: return - from ...application.routing.downlink import normalize_ua_voice_slot - - voice_slot = normalize_ua_voice_slot(peer, wire_slot) ctx = self._downlink_ctx() if not ctx.per_peer_contention(): return - if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: - from ...application.routing.downlink import end_peer_voice_slot - - end_peer_voice_slot( - ctx, peer_id, voice_slot, stream_id, dst_id, pkt_time=pkt_time, - ) - else: - from ...application.routing.downlink import touch_peer_voice_slot - - touch_peer_voice_slot( - ctx, peer_id, voice_slot, stream_id, dst_id, pkt_time=pkt_time, ingress=True, - ) + packet = synthetic_group_dmrd_burst_packet( + wire_slot, + int_id(dst_id), + stream_id, + frame_type=frame_type, + dtype_vseq=dtype_vseq, + call_type=call_type, + ) + track_peer_group_dmrd( + ctx, + peer_id, + packet, + peer, + pkt_time=pkt_time, + from_ingress=True, + voice_slot=normalize_ua_voice_slot(peer, wire_slot), + ) def _peer_would_accept_group_dmrd( self, @@ -651,15 +654,9 @@ class HBPProtocol(DatagramProtocol): peer = self._peers.get(peer_id) if peer is None: return True - ctx = self._downlink_ctx() - if not ctx.per_peer_contention(): - return True - route_pkt = ( - packet - if routed - else remap_dmrd_for_peer(packet, peer, self._config, peer_id=peer_id) + return peer_accepts_group_dmrd_packet( + self._downlink_ctx(), peer_id, peer, packet, routed=routed, ) - return not peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt) def _downlink_drop_key(self, peer_id: bytes, route_pkt: bytes) -> tuple[bytes, bytes]: pk = bytes_4(int_id(peer_id)) @@ -844,7 +841,7 @@ class HBPProtocol(DatagramProtocol): return 0 ctx = self._downlink_ctx() - route_pkt = build_dmra_route_packet(slot, tgid) + route_pkt = synthetic_group_dmrd_route_packet(slot, tgid) sent = 0 for peer_id in self._iter_downlink_peers(route_pkt): if exclude_peer and peer_id == exclude_peer: diff --git a/tests/application/test_slot_contention.py b/tests/application/test_slot_contention.py index 27fa8bf..0f03516 100644 --- a/tests/application/test_slot_contention.py +++ b/tests/application/test_slot_contention.py @@ -249,8 +249,8 @@ def test_per_peer_obp_tx_stamp_blocked_when_other_stream_active() -> None: ) is True -def test_lab_witness_many_static_tgs_allows_second_tg_on_slot() -> None: - """Full-table lab witness (>6 static TGs) is not subject to one-QSO hard slot lock.""" +def test_lab_witness_many_static_tgs_blocks_second_tg_on_slot() -> None: + """Full-table lab witness (>6 static TGs) still obeys one-QSO-per-RF-slot.""" now = 1_000_000.0 witness_id = bytes_4(730039257) tg_list = ",".join(str(730500 + i) for i in range(13)) @@ -259,7 +259,7 @@ def test_lab_witness_many_static_tgs_allows_second_tg_on_slot() -> None: 2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 1.0}, } slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} - assert not peer_hotspot_voice_slot_busy( + assert peer_hotspot_voice_slot_busy( witness_id, 2, _STREAM_B, @@ -518,8 +518,8 @@ def test_single0_ingress_tx_blocks_other_static_tg_on_slot() -> None: ) -def test_lab_witness_nine_static_tgs_allows_second_tg_on_slot() -> None: - """Lab witness with >6 static TGs is not hard-locked to one QSO per RF slot.""" +def test_lab_witness_nine_static_tgs_blocks_second_tg_on_slot() -> None: + """Lab witness with >6 static TGs is hard-locked to one QSO per RF slot.""" now = 1_000_000.0 hs = bytes_4(730039257) peer = { @@ -536,7 +536,7 @@ def test_lab_witness_nine_static_tgs_allows_second_tg_on_slot() -> None: peer_slots = { 2: {"stream_id": bytes_4(0x11111111), "tgid": 730507, "time": now, "ingress": False}, } - assert not peer_hotspot_voice_slot_busy( + assert peer_hotspot_voice_slot_busy( hs, 2, bytes_4(0x22222222),