From 8820d2ac999ba77a480f2c6767210dbff7b87333 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 29 Jun 2026 19:22:30 -0400 Subject: [PATCH] fix: per-peer slot contention, silent activation, and multi-hotspot REPEAT gate Enforce one active group QSO per peer/slot (SINGLE=0 and SINGLE=1) at the downlink layer instead of treating the shared MASTER STATUS[slot] as a single RF slot. Concurrent streams from different hotspots to different TGs are now legitimate on ingress; contention is decided per-peer. - Ingress: hbp_ingress_new_stream_collision honours per_peer=True so a second hotspot is not silenced by the first; hbp_master_ingress_repeat_allowed applies the same multi-hotspot rule and always passes VTERMs so listener sessions close and mid-call join works. - Silent activation: TX onto a TG with an active QSO activates the TG dynamically, suppresses the uplink, and delivers the in-progress QSO downlink. Only a TG in GROUP_HANGTIME (no live stream) is rejected. - Downlink: peer_voice_slots tracks one listen TG per peer/slot with GROUP_HANGTIME and bridge-hold semantics; stale sessions expire after STREAM_TO; duplex slots stay independent. - Non-regression tests for slot contention, silent activation, stale session expiry, and duplex independence. --- .../application/routing/downlink.py | 60 +- .../application/routing/hbp_forward.py | 34 +- src/adn_server/application/routing/helpers.py | 169 ++++- .../application/routing_use_cases.py | 56 +- .../subscription/subscription_queries.py | 49 ++ .../twisted_adapters/udp_hbp.py | 72 +- tests/application/test_slot_contention.py | 676 +++++++++++++++++- tests/hbp/test_timeout_collision.py | 13 +- 8 files changed, 1032 insertions(+), 97 deletions(-) diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index ef95407..8eff898 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -46,6 +46,7 @@ from .helpers import ( peer_receives_group_tgid, peer_should_receive_group_voice, peer_single_exclusive_tgid, + peer_single_mode, peer_wants_downlink_single_listen_lock, register_peer_ua_session, remap_dmrd_to_peer_static_slot, @@ -320,6 +321,8 @@ def touch_peer_voice_slot( "time": float(now), } prev = ctx.peer_voice_slots.get(pk, {}).get(int(voice_slot)) + if not ingress and isinstance(prev, dict) and prev.get("ingress"): + return if ingress or (isinstance(prev, dict) and prev.get("ingress")): row["ingress"] = True ctx.peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = row @@ -364,6 +367,7 @@ def end_peer_voice_slot( *, pkt_time: float | None = None, apply_hangtime: bool = True, + from_ingress: bool = False, ) -> None: """Close session on VTERM; ingress VTERM may start GROUP_HANGTIME window.""" now = time.time() if pkt_time is None else float(pkt_time) @@ -377,14 +381,49 @@ def end_peer_voice_slot( voice_slot = int(vs) active = per_slot.pop(vs, None) break + peer = ctx.peers.get(peer_id) or ctx.peers.get(pk) + if ( + isinstance(peer, dict) + and ended_tg + and ctx.sys_cfg is not None + and not peer_single_mode(peer, ctx.sys_cfg) + and ctx.subscription_store is not None + and ctx.system_name + ): + from adn_server.application.subscription.subscription_queries import ( + system_has_active_leg_in_store, + ) + + voice_ended = isinstance(active, dict) and ( + active.get("stream_id") or active.get("ingress") + ) + if voice_ended and system_has_active_leg_in_store( + ctx.subscription_store, + ctx.system_name, + int(voice_slot), + int(ended_tg), + ): + per_slot[int(voice_slot)] = { + "stream_id": b"", + "tgid": int(ended_tg), + "time": float(now), + "bridge_hold": True, + "bridge_hold_ingress": bool(from_ingress), + } + ctx.peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = per_slot[int(voice_slot)] + return if not apply_hangtime: return if isinstance(active, dict): - ended_tg = int(active.get("tgid", 0) or 0) - elif not ended_tg: + active_tg = int(active.get("tgid", 0) or 0) + # Only seed hangtime when the VTERM TG matches the session TG. A VTERM for a + # different TG (e.g. foreign stream that this peer never heard) must not seed + # hangtime from an unrelated active/bridge_hold session. + if active_tg and active_tg == ended_tg: + apply_hangtime_after_vterm(ctx, peer_id, voice_slot, active_tg, pkt_time=now) return - if ended_tg: - apply_hangtime_after_vterm(ctx, peer_id, voice_slot, ended_tg, pkt_time=now) + # No active session for this peer: VTERM for a stream it never received — do not + # seed GROUP_HANGTIME (would block a fresh PTT on a different TG). def track_peer_group_dmrd( @@ -429,6 +468,10 @@ def track_peer_group_dmrd( clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot) elif listen_lock and active_tg and active_tg == ended_tg: clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot) + pk = bytes_4(int_id(peer_id)) + apply_hangtime = from_ingress + if not from_ingress and not peer_single_mode(peer, ctx.sys_cfg): + apply_hangtime = True end_peer_voice_slot( ctx, peer_id, @@ -436,7 +479,8 @@ def track_peer_group_dmrd( stream_id, dst_id, pkt_time=pkt_time, - apply_hangtime=from_ingress, + apply_hangtime=apply_hangtime, + from_ingress=from_ingress, ) return if ( @@ -506,6 +550,8 @@ def peer_accepts_dmra( peer_id: bytes, slot: int, tgid: int, + *, + pkt_time: float | None = None, ) -> bool: """P1: DMRA uses same accept + slot-busy rules as DMRD.""" peer = ctx.peers.get(peer_id) @@ -515,7 +561,9 @@ def peer_accepts_dmra( 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) - return not peer_slot_blocks_downlink(ctx, peer_id, peer, remapped) + return not peer_slot_blocks_downlink( + ctx, peer_id, peer, remapped, pkt_time=pkt_time, + ) def peer_would_show_group_voice_on_monitor( diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 72710de..b6257cb 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -52,6 +52,7 @@ from .helpers import ( hbp_ingress_downlink_session_blocks_tx, hbp_ingress_new_stream_collision, master_per_peer_slot_contention, + tg_has_active_conversation, ) from .peer_downlink_index import count_connected_peers @@ -143,6 +144,8 @@ class HbpForwardMixin: _slot_st = getattr(src_proto, "STATUS", {}).get(slot, {}) _is_new_stream = stream_id != _slot_st.get("RX_STREAM_ID") if _is_new_stream: + _slot_st.pop("_suppress_uplink", None) + _slot_st.pop("_silent_activation_tg", None) sys_cfg = systems_cfg.get(system_name, {}) peers = getattr(src_proto, "_peers", None) or sys_cfg.get("PEERS", {}) connected = count_connected_peers(peers) if isinstance(peers, dict) else 0 @@ -163,14 +166,29 @@ class HbpForwardMixin: if group_voice_tg_ingress_collision( protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time, ): - self._log_ingress_warning_once( - self._ingress_drop_key( - "tg_busy", system_name, peer_id, dst_id, stream_id, - ), - "(%s) TG %s busy — dropping stream %s from peer %s", - system_name, int_id(dst_id), int_id(stream_id), int_id(peer_id), - ) - return False + # TX onto a TG with an active (in-progress) QSO is not rejected. + # The TG is activated dynamically, the user's uplink audio is suppressed + # (not forwarded to the network), and the downlink of the active QSO is + # delivered to the user. Only a TG that is merely in GROUP_HANGTIME + # (no live stream) is rejected as busy (legacy parity). + if tg_has_active_conversation( + protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time, + ): + logger.info( + "(%s) TG %s has active QSO — activating dynamic TG silently for peer %s (uplink suppressed)", + system_name, int_id(dst_id), int_id(peer_id), + ) + _slot_st["_suppress_uplink"] = True + _slot_st["_silent_activation_tg"] = int_id(dst_id) + else: + self._log_ingress_warning_once( + self._ingress_drop_key( + "tg_busy", system_name, peer_id, dst_id, stream_id, + ), + "(%s) TG %s busy — dropping stream %s from peer %s", + system_name, int_id(dst_id), int_id(stream_id), int_id(peer_id), + ) + return False from .downlink import normalize_ua_voice_slot peer = peers.get(peer_id) if isinstance(peers, dict) else None diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 3204db3..797c785 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -52,6 +52,11 @@ from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO PeerVoiceSlotRow = dict[str, Any] PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow] +# A per-peer downlink voice session with no frames for this long is considered +# dead (VTERM lost or stream abandoned). Matches the legacy bridge_master idle +# loop threshold (bridge_master.py ~607: RX_TIME < now - 5). +_STALE_PEER_SESSION_TIMEOUT = 5.0 + RF_MODE_SIMPLEX = "simplex" RF_MODE_DUPLEX = "duplex" # MMDVMHost DMO: downlink DMRD with TS1 bit set is dropped; only TS2 passes (DMRNetwork.cpp). @@ -347,9 +352,14 @@ def peer_hotspot_voice_slot_busy( ) -> bool: """True when this hotspot must not receive another group stream on ``voice_slot``. - Hard rule (SINGLE=0 and SINGLE=1): at most one group QSO per RF slot per hotspot. - A second TG or stream is dropped until VTERM clears ``peer_voice_slots`` (and - GROUP_HANGTIME where applicable). Same ``stream_id`` continues one call leg. + Hard rules (SINGLE=0 and SINGLE=1): + + - **Transmitting** (``ingress``): drop every downlink byte until VTERM clears the slot. + - **Listening** on TG *T*: drop every byte for TG *U* ≠ *T* on this RF slot (no hangtime + exception for another OPTIONS/UA TG). + - **Bridge hold**: while an ACTIVE bridge leg keeps *T* on the slot, foreign legs on + *U* ≠ *T* stay dropped (OBP or HBP). + - Same ``stream_id`` on the same TG continues one call leg. """ pk = bytes_4(int_id(peer_id)) if peer is None and peers is not None: @@ -364,17 +374,18 @@ def peer_hotspot_voice_slot_busy( active_time = float(active.get("time", 0) or 0) age = pkt_time - active_time if active.get("ingress"): - if ( - age < STREAM_TO - and isinstance(peer, dict) - and sys_cfg is not None - and not peer_single_mode(peer, sys_cfg) - and active_tgid - and incoming_tgid != active_tgid - ): - return False return True - if stream_id and active_stream: + if active.get("bridge_hold") and active_tgid and incoming_tgid != active_tgid: + # Bridge hold (ingress from own TX or listener) blocks a foreign TG only for + # GROUP_HANGTIME, matching legacy bridge_master contention on TX_TIME/RX_TIME. + if age <= group_hangtime: + return True + peer_slots.pop(int(voice_slot), None) + active = None + if isinstance(active, dict) and active_time > 0 and age >= _STALE_PEER_SESSION_TIMEOUT: + peer_slots.pop(int(voice_slot), None) + active = None + if isinstance(active, dict) and stream_id and active_stream: if active_stream == stream_id: pass elif active_tgid and active_tgid == incoming_tgid: @@ -386,8 +397,11 @@ def peer_hotspot_voice_slot_busy( return True else: return True - else: - return True + elif isinstance(active, dict): + if active_tgid and active_tgid == incoming_tgid: + pass + 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, ): @@ -396,11 +410,6 @@ def peer_hotspot_voice_slot_busy( peer, pk, incoming_tgid_b, stream_id, slot_st, sys_cfg, pkt_time=pkt_time, ): return True - if stream_id and stream_id == slot_st.get("TX_STREAM_ID"): - tx_peer = slot_st.get("TX_PEER") - if tx_peer is not None and int_id(tx_peer) != 0: - if peers is None or not peer_key_in_peers(tx_peer, peers): - return False if bytes_4(int_id(slot_st.get("RX_PEER", b""))) == pk: if slot_has_active_voice(slot_st, pkt_time) and stream_id != slot_st.get("RX_STREAM_ID"): return True @@ -530,17 +539,25 @@ def hbp_ingress_new_stream_collision( ) -> bool: """True when a new group-voice stream must drop on ingress (legacy routerHBP). - MASTER ``STATUS[slot]`` is shared: only one live stream per wire timeslot. - Same-subscriber rekey with a new stream id is allowed. ``per_peer`` applies - to downlink slot gates only, not ingress collision (legacy bridge_master). + When ``per_peer`` is False the MASTER ``STATUS[slot]`` is treated as a shared + RF slot: only one live stream per wire timeslot (legacy bridge_master). + + When ``per_peer`` is True the MASTER fronts multiple hotspots, each on its + own frequency, so concurrent streams to different TGs on the same slot are + legitimate. Per-hotspot contention (one listen TG per peer/slot) is enforced + at downlink time by ``peer_hotspot_voice_slot_busy``; it must not be applied + here on ingress, otherwise a second hotspot is silenced by the first. + Same-subscriber rekey with a new stream id is always allowed. """ - del per_peer, peer_id from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO if stream_id and stream_id == slot_st.get("RX_STREAM_ID"): return False if stream_id and stream_id == slot_st.get("TX_STREAM_ID"): return False + if per_peer: + return False + del peer_id for leg in ("RX", "TX"): type_key = f"{leg}_TYPE" time_key = f"{leg}_TIME" @@ -589,6 +606,65 @@ def _obp_stream_active(st: dict[str, Any], pkt_time: float) -> bool: return True +def tg_has_active_conversation( + protocols: dict[str, Any], + systems_cfg: dict[str, Any], + tgid_b: bytes, + stream_id: bytes, + rf_src: bytes, + pkt_time: float, +) -> bool: + """True when TG *tgid_b* has an active (in-progress) voice conversation. + + Distinct from ``group_voice_tg_ingress_collision``: that helper also matches a TG + that is merely in ``GROUP_HANGTIME`` (idle but recent). This one only matches a TG + with a live stream (within ``STREAM_TO`` of the last frame), so the ingress gate + can activate the TG silently (no uplink, deliver downlink of the active QSO) + instead of rejecting the stream. + """ + tgid = int_id(tgid_b) + if tgid < 5 or tgid in (9, 4000, 5000): + return False + for sys_name, proto in protocols.items(): + mode = systems_cfg.get(sys_name, {}).get("MODE") + status = getattr(proto, "STATUS", None) + if not isinstance(status, dict): + continue + if mode in ("MASTER", "PEER"): + for slot_key in (1, 2): + slot_st = status.get(slot_key) + if not isinstance(slot_st, dict): + continue + if not slot_has_active_voice(slot_st, pkt_time): + continue + active_tg = _hbp_slot_active_tgid(slot_st, pkt_time) + if active_tg is None or int_id(active_tg) != tgid: + continue + leg_stream = slot_st.get("RX_STREAM_ID") or slot_st.get("TX_STREAM_ID") + leg_rfs = slot_st.get("RX_RFS") or slot_st.get("TX_RFS") or b"\x00\x00\x00" + if stream_id and leg_stream == stream_id: + continue + if _same_rf_source(rf_src, leg_rfs): + continue + return True + elif mode == "OPENBRIDGE": + for key, st in status.items(): + if isinstance(key, int) or not isinstance(st, dict): + continue + if "TGID" not in st or int_id(st.get("TGID", b"")) != tgid: + continue + if not _obp_stream_active(st, pkt_time): + continue + leg_stream = key if isinstance(key, (bytes, bytearray)) else b"" + leg_rfs = st.get("RFS", b"\x00\x00\x00") + if stream_id and leg_stream == stream_id: + continue + if _same_rf_source(rf_src, leg_rfs): + continue + return True + return False + + def group_voice_tg_ingress_collision( protocols: dict[str, Any], systems_cfg: dict[str, Any], @@ -663,8 +739,18 @@ def hbp_master_ingress_repeat_allowed( *, protocols: dict[str, Any] | None = None, systems_cfg: dict[str, Any] | None = None, + is_vterm: bool = False, + system_cfg: dict[str, Any] | None = None, ) -> bool: - """True when MASTER REPEAT may fan this ingress packet to other peers.""" + """True when MASTER REPEAT may fan this ingress packet to other peers. + + In a multi-hotspot MASTER, concurrent streams from different peers to + different TGs on the same wire timeslot are legitimate: each hotspot is + on its own RF frequency. Contention is enforced per-peer at downlink + time (``peer_slot_blocks_downlink``). Applying a global shared-slot gate + here would drop voice packets of an active stream whenever a second + stream updates ``RX_STREAM_ID``, causing audible gaps on the first call. + """ if stream_id and stream_id == slot_st.get("RX_STREAM_ID"): owner = slot_st.get("RX_PEER") return ( @@ -672,8 +758,16 @@ def hbp_master_ingress_repeat_allowed( and int_id(owner) != 0 and bytes_4(int_id(owner)) == bytes_4(int_id(peer_id)) ) + # A VTERM closes an existing stream; it must reach peers that were hearing + # that stream even when another stream is now active on the shared slot + # (multi-hotspot MASTER). Blocking it leaves per-peer sessions open forever. + if is_vterm: + return True + per_peer = bool(system_cfg and master_per_peer_slot_contention( + systems_cfg or {}, "", system_cfg, connected_count=0, + )) if hbp_ingress_new_stream_collision( - slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=False, + slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=per_peer, ): return False if protocols and systems_cfg and group_voice_tg_ingress_collision( @@ -1321,19 +1415,22 @@ def peer_single_blocks_group_voice( ) -> bool: """True when SINGLE=1 peer must not receive downlink for ``tgid``. - With an active session on TG X, every other TG (static or dynamic) is blocked - until the TIMER expires or a new local TX replaces the session. + With an active SINGLE session on TG *X* in the peer's RF listen slot for + *tgid*, every other TG on **that same RF slot** is blocked until the TIMER + expires or a new local TX replaces the session. + + Duplex hotspots have independent RF timeslots: a listen lock on TS1 must not + block a static TG on TS2 (and vice versa). Simplex hotspots always collapse + to ``SIMPLEX_VOICE_SLOT``, so both TGs share the same RF slot and blocking + still applies. """ - del slot if not sys_cfg: return False - for voice_slot in (1, 2): - locked = peer_single_exclusive_tgid( - peer, voice_slot, sys_cfg, peer_id=peer_id, now=now, - ) - if locked is not None and int(tgid) != locked: - return True - return False + voice_slot = peer_downlink_voice_slot(peer, int(slot), int(tgid), sys_cfg, peer_id=peer_id) + locked = peer_single_exclusive_tgid( + peer, voice_slot, sys_cfg, peer_id=peer_id, now=now, + ) + return locked is not None and int(tgid) != locked def peer_single_blocks_foreign_same_tg_downlink( diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index b0ed0bf..202119b 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -276,6 +276,7 @@ class RoutingUseCases( if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: if not _obp_grp: _is_new_rx_stream = True + _suppress_uplink_rx = False if not source_is_obp: protocols = self._get_protocols() if self._get_protocols else {} src_proto = protocols.get(system_name) if protocols else None @@ -283,7 +284,8 @@ class RoutingUseCases( slot_st = src_proto.STATUS.get(slot, {}) if isinstance(slot_st, dict): _is_new_rx_stream = stream_id != slot_st.get("RX_STREAM_ID") - if _is_new_rx_stream: + _suppress_uplink_rx = bool(slot_st.get("_suppress_uplink")) + if _is_new_rx_stream and not _suppress_uplink_rx: _rx_report_peer = peer_id if not source_is_obp: _rx_report_peer = resolve_voice_peer_id( @@ -305,6 +307,7 @@ class RoutingUseCases( elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: if not _obp_grp: duration = 0.0 + _suppress_uplink_vterm = False protocols = self._get_protocols() if self._get_protocols else {} src_proto = protocols.get(system_name) if protocols else None if src_proto and getattr(src_proto, "STATUS", None): @@ -315,25 +318,29 @@ class RoutingUseCases( start = st.get(slot, {}).get("RX_START") if start is not None: duration = pkt_time - start - _rx_report_peer = peer_id - if not source_is_obp: - _rx_report_peer = resolve_voice_peer_id( - peer_id, - rf_src, - system_name, - systems_cfg, - ) - self._send_routing_event( - "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( - system_name, - int_id(stream_id), - int_id(_rx_report_peer), - int_id(rf_src), - slot, - int_id(dst_id), - duration, + slot_st_vterm = st.get(slot, {}) + if isinstance(slot_st_vterm, dict): + _suppress_uplink_vterm = bool(slot_st_vterm.get("_suppress_uplink")) + if not _suppress_uplink_vterm: + _rx_report_peer = peer_id + if not source_is_obp: + _rx_report_peer = resolve_voice_peer_id( + peer_id, + rf_src, + system_name, + systems_cfg, + ) + self._send_routing_event( + "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( + system_name, + int_id(stream_id), + int_id(_rx_report_peer), + int_id(rf_src), + slot, + int_id(dst_id), + duration, + ) ) - ) has_source = bool( self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int) ) @@ -435,6 +442,15 @@ class RoutingUseCases( resolve_voice_peer_id(peer_id, rf_src, system_name, systems_cfg) ) forwarded = [] + # When a TX arrived on a TG with an active QSO, the ingress gate set + # ``_suppress_uplink`` on the source slot status. The TG was activated dynamically + # (so downlink of the active QSO reaches this peer), but this stream's audio must + # not be forwarded upstream to avoid disrupting the in-progress conversation. + _suppress_uplink = False + if not source_is_obp and src_proto and getattr(src_proto, "STATUS", None): + _src_slot_st = src_proto.STATUS.get(slot, {}) + if isinstance(_src_slot_st, dict) and _src_slot_st.get("_suppress_uplink"): + _suppress_uplink = True _leg_iter: list[tuple[str, dict[str, Any]]] = [ ( forward_tables[0] if forward_tables else str(dst_int), @@ -449,6 +465,8 @@ class RoutingUseCases( ] for _relay_table_key, entry in _leg_iter: + if _suppress_uplink: + continue if not systems_cfg.get(entry["SYSTEM"], {}).get("ENABLED", True): continue _target_system = systems_cfg.get(entry["SYSTEM"], {}) diff --git a/src/adn_server/application/subscription/subscription_queries.py b/src/adn_server/application/subscription/subscription_queries.py index 43b2298..7daea34 100644 --- a/src/adn_server/application/subscription/subscription_queries.py +++ b/src/adn_server/application/subscription/subscription_queries.py @@ -74,6 +74,55 @@ def active_system_slots_for_tg_in_store( return tuple(sorted(slots)) +def system_has_other_active_bridge_on_slot( + store: SubscriptionStore, + system: str, + slot: int, + incoming_tgid: int, +) -> bool: + """True when another ACTIVE bridge leg occupies ``(system, slot)`` on a different TG. + + Diagnostic / table export only — static OPTIONS bridges stay ACTIVE idle and must + not gate per-peer downlink (see ``peer_voice_slots`` / hangtime in ``downlink.py``). + """ + incoming = int(incoming_tgid) + for sub in store.snapshot(): + if sub.system.value != system: + continue + if int(sub.channel.slot) != int(slot): + continue + if not sub.is_active(): + continue + if int(sub.target_tgid) != incoming: + return True + return False + + +def bridge_timer_active_on_slot( + store: SubscriptionStore, + system: str, + slot: int, + tgid: int, + *, + now: float, +) -> bool: + """True when an ACTIVE bridge leg for ``tgid`` on ``slot`` still has a live timer.""" + tg = int(tgid) + for sub in store.snapshot(): + if sub.system.value != system: + continue + if int(sub.channel.slot) != int(slot): + continue + if int(sub.target_tgid) != tg: + continue + if not sub.is_active(): + continue + exp = sub.state.timer_expires_at + if exp is not None and float(exp) > float(now): + return True + return False + + def store_legs_for_table( store: SubscriptionStore, table_key: str, diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index e57a644..5c377aa 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -1216,6 +1216,11 @@ class HBPProtocol(DatagramProtocol): pkt_time, protocols=_protocols, systems_cfg=self._CONFIG.get("SYSTEMS", {}), + is_vterm=( + _frame_type == HBPF_DATA_SYNC + and _dtype_vseq == HBPF_SLT_VTERM + ), + system_cfg=self._config, ) ) if ( @@ -1282,22 +1287,31 @@ class HBPProtocol(DatagramProtocol): ) if _accepted: if _slot in self.STATUS and not _unit_data: + _suppress = self.STATUS[_slot].get("_suppress_uplink") if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): - self.STATUS[_slot]["RX_START"] = pkt_time + if not _suppress: + self.STATUS[_slot]["RX_START"] = pkt_time if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt) >= 33: try: decoded_slot = decode.voice_head_term(dmrpkt) - self.STATUS[_slot]["RX_LC"] = decoded_slot["LC"] + if not _suppress: + self.STATUS[_slot]["RX_LC"] = decoded_slot["LC"] except Exception: - self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src + if not _suppress: + self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src else: - self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src - self.STATUS[_slot]["RX_PEER"] = _peer_id - self.STATUS[_slot]["RX_SEQ"] = _seq - self.STATUS[_slot]["RX_RFS"] = _rf_src - self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq - self.STATUS[_slot]["RX_TGID"] = _dst_id - self.STATUS[_slot]["RX_TIME"] = pkt_time + if not _suppress: + self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src + if not _suppress: + self.STATUS[_slot]["RX_PEER"] = _peer_id + self.STATUS[_slot]["RX_SEQ"] = _seq + self.STATUS[_slot]["RX_RFS"] = _rf_src + self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq + self.STATUS[_slot]["RX_TGID"] = _dst_id + self.STATUS[_slot]["RX_TIME"] = pkt_time + # Always track RX_STREAM_ID so subsequent frames of a suppressed + # stream are not re-evaluated as a new stream (which would clear + # _suppress_uplink and leak the uplink to the network). self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id _voice = self._CONFIG.get("VOICE", {}) if _accepted and self._on_handle_recording and _voice.get("RECORDING_ENABLED") and int_id(_dst_id) == _voice.get("RECORDING_TG") and _slot == _voice.get("RECORDING_TIMESLOT", 2): @@ -1329,9 +1343,20 @@ class HBPProtocol(DatagramProtocol): and _call_type in ("group", "vcsbk") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM + and _slot in self.STATUS and self._on_talker_alias_stream_end ): self._on_talker_alias_stream_end(self._system, _stream_id) + # Clean up silent-activation markers when the suppressed stream ends. + if ( + _accepted + and _call_type in ("group", "vcsbk") + and _frame_type == HBPF_DATA_SYNC + and _dtype_vseq == HBPF_SLT_VTERM + and _slot in self.STATUS + ): + self.STATUS[_slot].pop("_suppress_uplink", None) + self.STATUS[_slot].pop("_silent_activation_tg", None) elif _command == RPTL: _peer_id = _data[4:8] @@ -1723,15 +1748,28 @@ class HBPProtocol(DatagramProtocol): ): self._on_in_band_signalling(self._system, _slot, _dst_id, pkt_time) if _slot in self.STATUS and not _unit_data: - if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): + _suppress = self.STATUS[_slot].get("_suppress_uplink") + if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID") and not _suppress: self.STATUS[_slot]["RX_START"] = pkt_time - self.STATUS[_slot]["RX_PEER"] = _peer_id - self.STATUS[_slot]["RX_SEQ"] = _seq - self.STATUS[_slot]["RX_RFS"] = _rf_src - self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq - self.STATUS[_slot]["RX_TGID"] = _dst_id - self.STATUS[_slot]["RX_TIME"] = pkt_time + if not _suppress: + self.STATUS[_slot]["RX_PEER"] = _peer_id + self.STATUS[_slot]["RX_SEQ"] = _seq + self.STATUS[_slot]["RX_RFS"] = _rf_src + self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq + self.STATUS[_slot]["RX_TGID"] = _dst_id + self.STATUS[_slot]["RX_TIME"] = pkt_time + # Always track RX_STREAM_ID so subsequent frames of a suppressed + # stream are not re-evaluated as a new stream (leak of uplink). self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id + # Clean up silent-activation markers when the suppressed stream ends. + if ( + _call_type in ("group", "vcsbk") + and _frame_type == HBPF_DATA_SYNC + and _dtype_vseq == HBPF_SLT_VTERM + and _slot in self.STATUS + ): + self.STATUS[_slot].pop("_suppress_uplink", None) + self.STATUS[_slot].pop("_silent_activation_tg", None) elif _command == DMRA: if len(_data) >= DMRA_PACKET_LEN: diff --git a/tests/application/test_slot_contention.py b/tests/application/test_slot_contention.py index f37a5c0..4d0d895 100644 --- a/tests/application/test_slot_contention.py +++ b/tests/application/test_slot_contention.py @@ -16,6 +16,7 @@ from adn_server.application.routing.helpers import ( slot_has_active_voice, slot_in_group_hangtime, slot_status_hotspot_owner, + tg_has_active_conversation, ) from adn_server.domain import HBPF_DATA_SYNC, bytes_3, bytes_4 from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO @@ -61,13 +62,18 @@ def test_ingress_same_subscriber_rekey_allowed() -> None: ) -def test_ingress_collides_when_other_hotspot_owns_busy_slot() -> None: - """MASTER slot is global: second hotspot must not open a new stream on a busy TS.""" +def test_ingress_allows_other_hotspot_when_per_peer() -> None: + """Multi-hotspot MASTER: second hotspot may open a new stream on a busy TS. + + Per-peer contention (one listen TG per hotspot/slot) is enforced at + downlink, not at ingress; a global slot collision here would silence a + second hotspot that operates on its own frequency. + """ now = 1_000_000.0 slot = _active_rx_slot() slot["RX_RFS"] = bytes_3(730039264) slot["RX_PEER"] = bytes_4(730039264) - assert hbp_ingress_new_stream_collision( + assert not hbp_ingress_new_stream_collision( slot, bytes_4(730039265), bytes_3(730039265), @@ -78,7 +84,11 @@ def test_ingress_collides_when_other_hotspot_owns_busy_slot() -> None: def test_ingress_collides_when_obp_bridge_tx_leg_active() -> None: - """OBP bridge TX stamp on STATUS[slot] must block a new HBP ingress stream.""" + """OBP bridge TX stamp on STATUS[slot] must block a new HBP ingress stream. + + Only applies to the legacy shared-slot model (``per_peer=False``); with + ``per_peer=True`` the decision is deferred to per-peer downlink gates. + """ now = 1_000_000.0 slot = { "RX_TYPE": HBPF_SLT_VTERM, @@ -93,7 +103,7 @@ def test_ingress_collides_when_obp_bridge_tx_leg_active() -> None: bytes_3(730039267), _STREAM_B, now + 0.1, - per_peer=True, + per_peer=False, ) @@ -221,6 +231,94 @@ def test_group_voice_tg_collision_across_hbp_slots() -> None: ) +def test_tg_has_active_conversation_detects_live_hbp_qso() -> None: + """Spec §3: a TG with an active (in-progress) HBP QSO is detected as active conversation.""" + now = 1_000_000.0 + master_status = { + 2: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now), + } + master_status[2]["RX_RFS"] = bytes_3(730039264) + protocols = {"M1": type("P", (), {"STATUS": master_status})()} + systems = {"M1": {"MODE": "MASTER"}} + assert tg_has_active_conversation( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039265), now + 0.1, + ) + + +def test_tg_has_active_conversation_detects_live_obp_qso() -> None: + """Spec §3: a TG with an active (in-progress) OBP stream is detected as active conversation.""" + now = 1_000_000.0 + obp_stream = bytes_4(0x3E7A0F77) + obp_status = { + obp_stream: { + "TGID": bytes_3(730500), + "START": now - 5.0, + "LAST": now - 0.1, + "RFS": bytes_3(730039266), + }, + } + protocols = {"OBP": type("P", (), {"STATUS": obp_status})()} + systems = {"OBP": {"MODE": "OPENBRIDGE"}} + assert tg_has_active_conversation( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039256), now, + ) + + +def test_tg_has_active_conversation_false_for_hangtime_only() -> None: + """Spec §3 vs §4: a TG that is merely in GROUP_HANGTIME (no live stream) is NOT active conversation.""" + now = 1_000_000.0 + # Slot idle (VTERM) but recent — within hangtime + master_status = { + 2: { + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VTERM, + "RX_TGID": bytes_3(730500), + "RX_TIME": now - 1.0, + "RX_STREAM_ID": _STREAM_A, + "TX_TIME": 0.0, + }, + } + protocols = {"M1": type("P", (), {"STATUS": master_status})()} + systems = {"M1": {"MODE": "MASTER"}} + assert not tg_has_active_conversation( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039265), now, + ) + + +def test_tg_has_active_conversation_false_for_stale_obp() -> None: + """Spec §3: a truncated OBP stream (no recent packets) is NOT an active conversation.""" + now = 1_000_000.0 + obp_stream = bytes_4(0x3E7A0F77) + obp_status = { + obp_stream: { + "TGID": bytes_3(730500), + "START": now - 60.0, + "LAST": now - 45.0, + "RFS": bytes_3(730039266), + }, + } + protocols = {"OBP": type("P", (), {"STATUS": obp_status})()} + systems = {"OBP": {"MODE": "OPENBRIDGE"}} + assert not tg_has_active_conversation( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039256), now, + ) + + +def test_tg_has_active_conversation_false_for_same_rf_source() -> None: + """Spec §3: same RF source rekeying is not an 'active conversation' collision.""" + now = 1_000_000.0 + rf = bytes_3(730039264) + master_status = { + 2: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now), + } + master_status[2]["RX_RFS"] = rf + protocols = {"M1": type("P", (), {"STATUS": master_status})()} + systems = {"M1": {"MODE": "MASTER"}} + assert not tg_has_active_conversation( + protocols, systems, bytes_3(730500), _STREAM_B, rf, now + 0.1, + ) + + def test_peer_hotspot_voice_slot_busy_blocks_downlink_during_local_ingress() -> None: """Rejected ingress still marks local TX — downlink must not arrive on that slot.""" now = 1_000_000.0 @@ -294,7 +392,7 @@ def test_lab_witness_many_static_tgs_blocks_second_tg_on_slot() -> None: tg_list = ",".join(str(730500 + i) for i in range(13)) witness = {"OPTIONS": f"TS2={tg_list};SINGLE=1;".encode()} peer_slots = { - 2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 1.0}, + 2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 0.1}, } slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} assert peer_hotspot_voice_slot_busy( @@ -644,8 +742,8 @@ def test_peer_hotspot_hangtime_blocks_other_tg() -> None: ) -def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None: - """SINGLE=0: RX on another OPTIONS TG while local TX on the same RF slot.""" +def test_single0_ingress_tx_blocks_other_tg_on_slot() -> None: + """SINGLE=0: no downlink bytes on another TG while local TX on the same RF slot.""" now = 1_000_000.0 hs = bytes_4(730039253) peer = {"OPTIONS": b"TS2=730507,730508;SINGLE=0;"} @@ -659,7 +757,7 @@ def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None: "ingress": True, }, } - assert not peer_hotspot_voice_slot_busy( + assert peer_hotspot_voice_slot_busy( hs, 2, bytes_4(0x22222222), @@ -671,6 +769,7 @@ def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None: 0.0, peer=peer, sys_cfg=sys_cfg, + peers={hs: peer, bytes_4(730039254): {"OPTIONS": b"TS2=730508;"}}, ) @@ -1042,3 +1141,562 @@ def test_global_slot_blocks_foreign_vterm_during_active_rx() -> None: slot, bytes_3(71442), panama_stream, now + 0.1, 0.0, is_vterm=True, ) assert not hbp_slot_blocks_group_voice(slot, _TG_A, chile_stream, now + 0.1, 0.0) + + +def test_single0_listen_vterm_hangtime_blocks_obp_and_dmra() -> None: + """SINGLE=0 listen-only: post-VTERM GROUP_HANGTIME blocks OBP voice and TA on another TG.""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_accepts_dmra, + peer_slot_blocks_downlink, + track_peer_group_dmrd, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039269) + peer = {"OPTIONS": b"TS2=730501,730504;"} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={ + 1: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}, + 2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}, + }, + connected_count=5, + ) + now = 1_000_000.0 + j39jq_stream = bytes_4(0x11111111) + obp_stream = bytes_4(0x22222222) + vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), j39jq_stream, + ] + [b"\x00"] * 33) + vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), j39jq_stream, + ] + [b"\x00"] * 33) + track_peer_group_dmrd(ctx, hs, vhead, peer, pkt_time=now) + track_peer_group_dmrd(ctx, hs, vterm, peer, pkt_time=now + 4.5) + assert ctx.peer_voice_hangtime[hs][2] == (730502, now + 4.5) + obp_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), obp_stream, + ] + [b"\x00"] * 33) + overlap = now + 9.0 + assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=overlap) + assert not peer_accepts_dmra(ctx, hs, 2, 730504, pkt_time=overlap) + assert not peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 11.0) + + +def test_live_listen_session_blocks_foreign_tg_downlink() -> None: + """Mid-QSO downlink listen on TG A blocks foreign TG B on same slot (730502 vs 730504).""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_accepts_dmra, + peer_slot_blocks_downlink, + touch_peer_voice_slot, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039269) + peer = {"OPTIONS": b"TS2=730501,730504;"} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}, + connected_count=5, + ) + now = 1_000_000.0 + j39jq_stream = bytes_4(0x11111111) + touch_peer_voice_slot(ctx, hs, 2, j39jq_stream, bytes_3(730502), pkt_time=now) + obp_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222), + ] + [b"\x00"] * 33) + assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 0.5) + assert not peer_accepts_dmra(ctx, hs, 2, 730504, pkt_time=now + 0.5) + + +def test_idle_static_bridges_do_not_block_options_fanout() -> None: + """Static ACTIVE bridge rows must not drop OPTIONS fan-out for another TG on the slot.""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + ) + from adn_server.application.subscription.routing_table_import import ( + subscriptions_from_routing_table, + ) + from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + def _row(*, system: str, ts: int, tgid: int) -> dict: + return { + "SYSTEM": system, + "TS": ts, + "TGID": bytes_3(tgid), + "ACTIVE": True, + "TIMEOUT": 3600.0, + "TO_TYPE": "OFF", + "ON": [bytes_3(tgid)], + "OFF": [], + "RESET": [], + "TIMER": 0.0, + } + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": True} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039252) + peer = {"OPTIONS": b"TS2=730500,730504;"} + store = InMemorySubscriptionStore() + store.replace_all( + subscriptions_from_routing_table( + { + "730501": [_row(system="MASTER-A", ts=2, tgid=730501)], + "730502": [_row(system="MASTER-A", ts=2, tgid=730502)], + } + ) + ) + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}, + connected_count=6, + subscription_store=store, + ) + ref_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039251), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (1 << 4) | HBPF_SLT_VHEAD]), bytes_4(0x33333333), + ] + [b"\x00"] * 33) + assert not peer_slot_blocks_downlink(ctx, hs, peer, ref_vhead, pkt_time=1_000_002.0) + + +def test_foreign_obp_blocked_when_master_slot_carries_other_tg() -> None: + """OBP (non-hotspot rf_src) blocked when this hotspot still holds the prior TG on slot.""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039269) + peer = {"OPTIONS": b"TS2=730501,730504;"} + now = 1_000_000.0 + j39jq_stream = bytes_4(0x11111111) + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={ + 2: { + "RX_PEER": bytes_4(730039270), + "RX_TGID": bytes_3(730502), + "RX_STREAM_ID": j39jq_stream, + "RX_TIME": now, + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VTERM, + }, + }, + peer_voice_slots={ + hs: { + 2: { + "stream_id": b"", + "tgid": 730502, + "time": now, + "bridge_hold": True, + }, + }, + }, + connected_count=5, + ) + obp_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222), + ] + [b"\x00"] * 33) + hbp_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039270), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x33333333), + ] + [b"\x00"] * 33) + assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 0.5) + assert peer_slot_blocks_downlink(ctx, hs, peer, hbp_vhead, pkt_time=now + 0.5) + + +def test_single0_ingress_vterm_bridge_hold_blocks_within_hangtime() -> None: + """SINGLE=0 ingress TX: bridge_hold blocks a foreign TG only within GROUP_HANGTIME. + + Matches legacy bridge_master contention on TX_TIME: after hangtime expires a fresh + PTT on a different TG is delivered. + """ + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + track_peer_group_dmrd, + ) + from adn_server.application.subscription.routing_table_import import ( + subscriptions_from_routing_table, + ) + from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + def _row(*, system: str, ts: int, tgid: int) -> dict: + return { + "SYSTEM": system, + "TS": ts, + "TGID": bytes_3(tgid), + "ACTIVE": True, + "TIMEOUT": 3600.0, + "TO_TYPE": "OFF", + "ON": [bytes_3(tgid)], + "OFF": [], + "RESET": [], + "TIMER": 0.0, + } + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039270) + peer = {"OPTIONS": b"TS2=730500,730508;"} + store = InMemorySubscriptionStore() + store.replace_all( + subscriptions_from_routing_table( + {"730502": [_row(system="MASTER-A", ts=2, tgid=730502)]}, + ) + ) + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}, + connected_count=5, + subscription_store=store, + ) + now = 1_000_000.0 + stream = bytes_4(0x11111111) + vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), stream, + ] + [b"\x00"] * 33) + vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), stream, + ] + [b"\x00"] * 33) + track_peer_group_dmrd(ctx, hs, vhead, peer, pkt_time=now, from_ingress=True, voice_slot=2) + track_peer_group_dmrd(ctx, hs, vterm, peer, pkt_time=now + 3, from_ingress=True, voice_slot=2) + assert ctx.peer_voice_slots[hs][2].get("bridge_hold") is True + obp_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222), + ] + [b"\x00"] * 33) + # Within GROUP_HANGTIME of the VTERM (age < 5s): blocked + assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 4.0) + # After GROUP_HANGTIME (age > 5s): fresh PTT on a different TG is delivered + assert not peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 9.0) + + +def test_single0_listener_bridge_hold_expires_after_hangtime_allows_fresh_ptt() -> None: + """SINGLE=0 listener (downlink RX, not ingress TX): bridge_hold must expire after + GROUP_HANGTIME so a fresh PTT on a different TG is delivered. + + Reproduces the hs1/hs2/hs3 rule: hs1 receives 730500; after it ends and hangtime + clears, a NEW call on 730501 must reach hs1 from the start (no stale bridge_hold). + """ + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + track_peer_group_dmrd, + ) + from adn_server.application.subscription.routing_table_import import ( + subscriptions_from_routing_table, + ) + from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + def _row(*, system: str, ts: int, tgid: int) -> dict: + return { + "SYSTEM": system, + "TS": ts, + "TGID": bytes_3(tgid), + "ACTIVE": True, + "TIMEOUT": 3600.0, + "TO_TYPE": "OFF", + "ON": [bytes_3(tgid)], + "OFF": [], + "RESET": [], + "TIMER": 0.0, + } + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039251) + peer = {"OPTIONS": b"TS2=730500,730501;"} + store = InMemorySubscriptionStore() + store.replace_all( + subscriptions_from_routing_table( + { + "730500": [_row(system="MASTER-A", ts=2, tgid=730500)], + "730501": [_row(system="MASTER-A", ts=2, tgid=730501)], + } + ) + ) + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}, + connected_count=5, + subscription_store=store, + ) + now = 1_000_000.0 + ref_stream = bytes_4(0xAAAAAAAA) + ref_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream, + ] + [b"\x00"] * 33) + ref_vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), ref_stream, + ] + [b"\x00"] * 33) + track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, from_ingress=False, voice_slot=2) + track_peer_group_dmrd(ctx, hs, ref_vterm, peer, pkt_time=now + 8.0, from_ingress=False, voice_slot=2) + fresh_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0xBBBBBBBB), + ] + [b"\x00"] * 33) + # During hangtime: fresh PTT on different TG must be blocked (listening on 730500). + assert peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 9.0) + # After hangtime clears: fresh PTT on different TG must be delivered. + assert not peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 20.0) + + +def test_single0_listener_fresh_ptt_after_blocked_second_stream() -> None: + """Full hs1/hs2/hs3 flow: hs_a receives 730500; a second stream 730501 transits + the server (delivered to witness, blocked for hs_a); after 730500 ends and + hangtime clears, a fresh PTT on 730501 must reach hs_a from the start. + + This reproduces the live scenario where a concurrent second stream leaves + residual slot state that blocked the fresh PTT in production. + """ + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + track_peer_group_dmrd, + ) + from adn_server.application.subscription.routing_table_import import ( + subscriptions_from_routing_table, + ) + from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + def _row(*, system: str, ts: int, tgid: int) -> dict: + return { + "SYSTEM": system, + "TS": ts, + "TGID": bytes_3(tgid), + "ACTIVE": True, + "TIMEOUT": 3600.0, + "TO_TYPE": "OFF", + "ON": [bytes_3(tgid)], + "OFF": [], + "RESET": [], + "TIMER": 0.0, + } + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039251) + peer = {"OPTIONS": b"TS2=730500,730501;"} + store = InMemorySubscriptionStore() + store.replace_all( + subscriptions_from_routing_table( + { + "730500": [_row(system="MASTER-A", ts=2, tgid=730500)], + "730501": [_row(system="MASTER-A", ts=2, tgid=730501)], + } + ) + ) + now = 1_000_000.0 + ref_stream = bytes_4(0xAAAAAAAA) + second_stream = bytes_4(0xCCCCCCCC) + fresh_stream = bytes_4(0xDDDDDDDD) + status = {2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status=status, + connected_count=5, + subscription_store=store, + ) + ref_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream, + ] + [b"\x00"] * 33) + ref_vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), ref_stream, + ] + [b"\x00"] * 33) + second_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), second_stream, + ] + [b"\x00"] * 33) + second_vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), second_stream, + ] + [b"\x00"] * 33) + fresh_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), fresh_stream, + ] + [b"\x00"] * 33) + # t=0: hs_a starts receiving 730500. + track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, from_ingress=False, voice_slot=2) + # t=3: second stream 730501 arrives — blocked for hs_a (busy on 730500). + assert peer_slot_blocks_downlink(ctx, hs, peer, second_vhead, pkt_time=now + 3.0) + # While the second stream is active on the wire, STATUS reflects the second + # peer as RX owner of slot 2 (the ingress updates RX_PEER/TGID/TIME/TYPE). + status[2] = { + "RX_PEER": bytes_4(730039252), + "RX_TGID": bytes_3(730501), + "RX_TIME": now + 5.0, + "RX_TYPE": HBPF_SLT_VHEAD, + "RX_STREAM_ID": second_stream, + "TX_TYPE": HBPF_SLT_VTERM, + "TX_TIME": 0.0, + } + # t=8: 730500 ends for hs_a. + track_peer_group_dmrd(ctx, hs, ref_vterm, peer, pkt_time=now + 8.0, from_ingress=False, voice_slot=2) + # t=12: second stream 730501 still in progress — no mid-join (hangtime clear path). + assert peer_slot_blocks_downlink(ctx, hs, peer, second_vhead, pkt_time=now + 12.0) + # t=17: second stream ends. STATUS global still carries the second peer as RX_OWNER. + track_peer_group_dmrd(ctx, hs, second_vterm, peer, pkt_time=now + 17.0, from_ingress=False, voice_slot=2) + status[2]["RX_TYPE"] = HBPF_SLT_VTERM + status[2]["RX_TIME"] = now + 17.0 + # t=20: fresh PTT on 730501 — must reach hs_a (all prior hangtime cleared). + blocked = peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 20.0) + assert not blocked, f"fresh 730501 PTT blocked after all streams ended: slots={ctx.peer_voice_slots.get(hs)}, status={status.get(2)}" + + +def test_peer_stale_session_different_tg_expires_after_stream_to() -> None: + """A stale per-peer voice session (no VTERM seen) on TG A must not block a + fresh stream on TG B once ``STREAM_TO`` has elapsed. + + Reproduces the production case where a listener's downlink session on 730500 + was left in ``peer_voice_slots`` with a non-empty stream_id after the stream + ended without VTERM, and blocked a subsequent 730501 PTT indefinitely. + """ + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(730039251) + peer = {"OPTIONS": b"TS2=730500,730501;"} + now = 1_000_000.0 + status = {2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status=status, + connected_count=5, + subscription_store=None, + ) + ctx.peer_voice_slots[hs] = { + 2: {"stream_id": bytes_4(0x952AFD04), "tgid": 730500, "time": now}, + } + fresh_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0xDDDDDDDD), + ] + [b"\x00"] * 33) + # While within STREAM_TO window: different TG is blocked (active QSO). + assert peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 0.1) + # After the stale-session timeout elapses: dead session expires, fresh PTT on 730501 passes. + assert not peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 6.0) + + +def test_single1_duplex_listen_lock_does_not_block_other_rf_slot() -> None: + """SINGLE=1 duplex hotspot: a listen lock on TS1 must not block TS2. + + Reproduces the intermittent ``single1-overlap-second-longer`` failure where + ``hs_a`` (TS1=730500, TS2=730501, SINGLE=1) receives 730500 on TS1, then a + concurrent 730501 stream on TS2 was incorrectly blocked because + ``peer_single_blocks_group_voice`` iterated both RF slots instead of + scoping the lock to the peer's listen slot for the incoming TG. + + Duplex hotspots have independent RF timeslots; a SINGLE listen lock on one + timeslot must not deny voice on the other. + """ + from adn_server.application.routing.downlink import ( + DownlinkContext, + track_peer_group_dmrd, + ) + from adn_server.application.routing.helpers import ( + peer_should_receive_group_voice, + ) + from adn_server.application.subscription.routing_table_import import ( + subscriptions_from_routing_table, + ) + from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + + def _row(*, system: str, ts: int, tgid: int) -> dict: + return { + "SYSTEM": system, + "TS": ts, + "TGID": bytes_3(tgid), + "ACTIVE": True, + "TIMEOUT": 3600.0, + "TO_TYPE": "OFF", + "ON": [bytes_3(tgid)], + "OFF": [], + "RESET": [], + "TIMER": 0.0, + } + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False} + hs = bytes_4(730039251) + peer = {"OPTIONS": b"TS1=730500;TS2=730501;SINGLE=1;TIMER=60;"} + store = InMemorySubscriptionStore() + store.replace_all( + subscriptions_from_routing_table( + { + "730500": [_row(system="MASTER-A", ts=1, tgid=730500)], + "730501": [_row(system="MASTER-A", ts=2, tgid=730501)], + } + ) + ) + ctx = DownlinkContext( + config={"SYSTEMS": {"MASTER-A": sys_cfg}}, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={hs: peer}, + status={ + 1: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}, + 2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}, + }, + connected_count=5, + subscription_store=store, + ) + now = 1_000_000.0 + ref_stream = bytes_4(0xAAAAAAAA) + ref_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream, + ] + [b"\x00"] * 33) + # hs_a receives 730500 on TS1 — registers a SINGLE listen lock on slot 1. + track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, voice_slot=1) + assert peer["_UA_SESSION"][1]["tgid"] == 730500 + # A 730501 stream arrives on TS2 — different RF slot, must not be blocked. + assert peer_should_receive_group_voice( + peer, 2, 730501, peer_id=hs, system="MASTER-A", + subscription_store=store, connected_count=5, sys_cfg=sys_cfg, now=now + 9.0, + ) diff --git a/tests/hbp/test_timeout_collision.py b/tests/hbp/test_timeout_collision.py index 255d4d1..9a7d718 100644 --- a/tests/hbp/test_timeout_collision.py +++ b/tests/hbp/test_timeout_collision.py @@ -51,7 +51,13 @@ def test_hbp_source_timeout_drops_after_180_seconds() -> None: assert scenario.protocols["MASTER-A"].STATUS[2].get("LOOPLOG") is True -def test_hbp_stream_collision_drops_conflicting_new_stream() -> None: +def test_hbp_stream_collision_silent_activation_on_busy_tg() -> None: + """Spec §3 divergence: TX onto a TG with an active QSO is not rejected. + + The stream passes the ingress gate (silent activation), but uplink audio is + suppressed — no forwarding to other systems. The downlink of the active QSO + continues to reach the peer independently. + """ bridges = active_routing_table(91, (("MASTER-A", 2), ("MASTER-B", 2))) scenario = DeterministicScenario(routing_table=bridges) t0 = scenario.clock.time() @@ -74,8 +80,11 @@ def test_hbp_stream_collision_drops_conflicting_new_stream() -> None: ingress_pkt_time=t0 + 0.1, ) - assert not ok + assert_inject_ok(ok) + # Uplink is suppressed: no packets forwarded to MASTER-B assert len(scenario.capture.for_system("MASTER-B")) == 0 + # Silent activation marker is set on the source slot + assert scenario.protocols["MASTER-A"].STATUS[2].get("_suppress_uplink") is True @pytest.mark.behavior