From ecca56febdb6698aee939099d0c37f1205ee941c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 22 Jun 2026 18:09:22 -0400 Subject: [PATCH] fix: block downlink DMRD to busy hotspot RF slot per peer Per-hotspot voice sessions now persist until VTERM so inter-burst gaps no longer admit foreign TGs from OBP or REPEAT on the same timeslot. --- .../application/report/monitor_topology.py | 23 ++- .../application/routing/hbp_forward.py | 11 +- src/adn_server/application/routing/helpers.py | 166 ++++++++++++++++-- .../application/routing_use_cases.py | 8 + .../twisted_adapters/udp_hbp.py | 151 +++++++++++++++- tests/application/test_monitor_topology.py | 41 +++++ .../test_obp_hbp_cross_slot_downlink.py | 16 ++ tests/application/test_slot_contention.py | 64 +++++++ 8 files changed, 457 insertions(+), 23 deletions(-) diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py index 4777f68..3cf517b 100644 --- a/src/adn_server/application/report/monitor_topology.py +++ b/src/adn_server/application/report/monitor_topology.py @@ -35,6 +35,7 @@ from typing import Any from adn_server.application.proxy.deployment import is_proxy_inject_only, proxy_target_system from adn_server.application.routing.helpers import ( hbp_slot_blocks_group_voice_for_peer, + master_per_peer_slot_contention, is_special_tg, peer_downlink_voice_slot, peer_should_receive_group_voice, @@ -396,7 +397,9 @@ def remap_inject_proxy_voice_events( connected = _connected_peers(peers) slot_map = _resolve_slot_map(connected, peer_slots, max_slots=max_slots) trx = parts[2].strip() if len(parts) > 2 else "" - per_peer = is_proxy_inject_only(config, target) + per_peer = master_per_peer_slot_contention( + config, target, sys_cfg, connected_count=len(connected), + ) stream_id = _voice_event_stream_id(parts) if trx == "TX": @@ -405,6 +408,22 @@ def remap_inject_proxy_voice_events( if echo_peer is not None: slot = slot_map.get(echo_peer) if slot is not None: + peer = peers.get(echo_peer) + if ( + peer is not None + and tgid_slot is not None + and monitor_downlink_blocked_by_slot_contention( + _peer_key_from_int(echo_peer), + peer, + tgid_slot[1], + tgid_slot[0], + stream_id, + master_status, + sys_cfg, + per_peer=per_peer, + ) + ): + return [] tx_parts = list(parts) if tgid_slot is not None: tgid, _ = tgid_slot @@ -464,7 +483,7 @@ def remap_inject_proxy_voice_events( peer_parts, target=target, slot=mapped_slot, peer_key=peer_key ) ) - return remapped if remapped else [event] + return remapped peer_key = _peer_key_from_voice_csv(parts, peers) if peer_key is None: diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 3c353f3..dc519f6 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -47,8 +47,8 @@ import logging from hashlib import blake2b from ...domain import int_id -from ..proxy.deployment import is_proxy_inject_only -from .helpers import hbp_ingress_new_stream_collision +from .helpers import hbp_ingress_new_stream_collision, master_per_peer_slot_contention +from .peer_downlink_index import count_connected_peers logger = logging.getLogger(__name__) @@ -88,7 +88,12 @@ class HbpForwardMixin: _slot_st.pop("_bcsq", None) _slot_st["lastSeq"] = False _slot_st["lastData"] = False - per_peer = is_proxy_inject_only(self._config, system_name) + 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 + per_peer = master_per_peer_slot_contention( + self._config, system_name, sys_cfg, connected_count=connected, + ) if hbp_ingress_new_stream_collision( _slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=per_peer, ): diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 1cc40f3..811b4b8 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -49,6 +49,9 @@ from typing import Any from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, bytes_3, bytes_4, int_id from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO +PeerVoiceSlotRow = dict[str, Any] +PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow] + RF_MODE_SIMPLEX = "simplex" RF_MODE_DUPLEX = "duplex" # MMDVMHost DMO: downlink DMRD with TS1 bit set is dropped; only TS2 passes (DMRNetwork.cpp). @@ -186,27 +189,130 @@ def slot_status_peer_owner(slot_st: dict[str, Any]) -> bytes | None: return None +def peer_key_in_peers(peer_id: bytes, peers: dict[Any, Any] | None) -> bool: + """True when ``peer_id`` is a connected hotspot key in ``PEERS``.""" + if not peers: + return False + pk = bytes_4(int_id(peer_id)) + if pk in peers: + return True + for key in peers: + try: + if bytes_4(int_id(key)) == pk: + return True + except (TypeError, ValueError): + continue + return False + + +def slot_status_hotspot_owner( + slot_st: dict[str, Any], + peers: dict[Any, Any] | None = None, +) -> bytes | None: + """Connected hotspot owning this slot row; ignores bridge ``TX_PEER`` (e.g. OBP 73010).""" + for field in ("RX_PEER", "TX_PEER"): + raw = slot_st.get(field) + if raw is None or int_id(raw) == 0: + continue + pk = bytes_4(int_id(raw)) + if peers is not None and not peer_key_in_peers(pk, peers): + continue + return pk + return None + + +def peer_hotspot_voice_slot_busy( + peer_id: bytes, + voice_slot: int, + stream_id: bytes, + incoming_tgid_b: bytes, + slot_st: dict[str, Any], + peer_slots: PeerVoiceSlotMap | None, + hang_row: tuple[int, float] | None, + pkt_time: float, + group_hangtime: float, +) -> bool: + """True when this hotspot must not receive another group stream on ``voice_slot``.""" + pk = bytes_4(int_id(peer_id)) + active = (peer_slots or {}).get(int(voice_slot)) + if isinstance(active, dict): + active_stream = active.get("stream_id") + if stream_id and active_stream == stream_id: + return False + # Session stays open until VTERM clears ``peer_slots``; DMR voice has + # inter-burst gaps longer than STREAM_TO so time-since-last-packet must + # not release the slot to another TG mid-QSO. + return True + 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 + if slot_in_group_hangtime(slot_st, incoming_tgid_b, pkt_time, group_hangtime): + return True + if hang_row is not None and float(group_hangtime or 0) > 0: + last_tg, last_t = hang_row + if int_id(incoming_tgid_b) != int(last_tg) and (pkt_time - float(last_t)) < float(group_hangtime): + return True + return False + + +def master_per_peer_slot_contention( + config: dict[str, Any], + system_name: str, + system_cfg: dict[str, Any], + *, + connected_count: int = 0, +) -> bool: + """True when slot busy/hangtime applies per hotspot, not globally on the MASTER row.""" + if system_cfg.get("MODE") != "MASTER": + return False + from ..proxy.deployment import is_proxy_inject_only + + if is_proxy_inject_only(config, system_name): + return True + return connected_count > 1 + + def inject_only_defer_obp_hbp_slot_contention( config: dict[str, Any], target_system: str, target_system_cfg: dict[str, Any], *, source_is_obp: bool, + connected_count: int = 0, ) -> bool: """Whether OBP→MASTER ``to_target`` should skip global slot STATUS contention. - Inject-only proxies defer slot checks to ``send_peer`` (same as REPEAT): - per-peer ``hbp_slot_blocks_group_voice_for_peer`` + OPTIONS/UA slot remap. - Global STATUS on the bridge wire TS would block cross-slot downlink while - another peer is active on that TS even though recipients listen on the other TS. + Defer to ``send_peer`` (same as REPEAT): per-peer ``hbp_slot_blocks_group_voice_for_peer`` + + OPTIONS/UA slot remap. Global STATUS on the bridge wire TS would block cross-slot + downlink while another peer is active on that TS even though recipients listen elsewhere. """ if not source_is_obp: return False if target_system_cfg.get("MODE") != "MASTER": return False - from ..proxy.deployment import is_proxy_inject_only + return master_per_peer_slot_contention( + config, target_system, target_system_cfg, connected_count=connected_count, + ) + - return is_proxy_inject_only(config, target_system) +def _downlink_same_stream_for_peer( + slot_st: dict[str, Any], + peer_id: bytes, + stream_id: bytes, +) -> bool: + """True when ``stream_id`` continues an RF leg owned by ``peer_id`` (downlink fan-out).""" + if not stream_id: + return False + pid = bytes_4(int_id(peer_id)) + if stream_id == slot_st.get("RX_STREAM_ID"): + rx_peer = slot_st.get("RX_PEER") + if rx_peer is not None and bytes_4(int_id(rx_peer)) == pid: + return True + if stream_id == slot_st.get("TX_STREAM_ID"): + tx_peer = slot_st.get("TX_PEER") + if tx_peer is not None and bytes_4(int_id(tx_peer)) == pid: + return True + return False def hbp_slot_blocks_group_voice_for_peer( @@ -218,16 +324,41 @@ def hbp_slot_blocks_group_voice_for_peer( group_hangtime: float, *, per_peer: bool, + peers: dict[Any, Any] | None = None, + peer_slots: PeerVoiceSlotMap | None = None, + peer_hang_row: tuple[int, float] | None = None, + voice_slot: int | None = None, ) -> bool: """Slot contention scoped to one hotspot when ``per_peer`` (inject-only multi-HS). ``STATUS[slot]`` is shared at the MASTER, but each connected hotspot has an independent RF timeslot. Another peer's active QSO must not block this peer. + + Inject-only OBP→HBP defers global contention and stamps bridge ``TX_*`` on the + shared slot row before ``send_peer``. Same-stream exemption must not treat that + bridge TX stamp as the hotspot's own leg while the peer is still on the air (RX). """ if per_peer: - owner = slot_status_peer_owner(slot_st) + if voice_slot is not None and peer_hotspot_voice_slot_busy( + peer_id, + int(voice_slot), + stream_id, + incoming_tgid_b, + slot_st, + peer_slots, + peer_hang_row, + pkt_time, + group_hangtime, + ): + return True + owner = slot_status_hotspot_owner(slot_st, peers) if owner is not None and bytes_4(int_id(owner)) != bytes_4(int_id(peer_id)): return False + if _downlink_same_stream_for_peer(slot_st, peer_id, stream_id): + return False + if slot_has_active_voice(slot_st, pkt_time): + return True + return slot_in_group_hangtime(slot_st, incoming_tgid_b, pkt_time, group_hangtime) return hbp_slot_blocks_group_voice( slot_st, incoming_tgid_b, stream_id, pkt_time, group_hangtime, ) @@ -429,17 +560,30 @@ def is_special_tg(relay_table_key: str) -> bool: def parse_dmrd_route_fields(packet: bytes) -> tuple[int, int, str] | None: """Parse HBP DMRD slot, destination TG, and call type for downlink OPTIONS filter.""" - if len(packet) < 17 or packet[:4] != b"DMRD": + burst = parse_dmrd_burst_fields(packet) + if burst is None: + return None + slot, _, _, _, dst_id, call_type = burst + return slot, int_id(dst_id), call_type + + +def parse_dmrd_burst_fields( + packet: bytes, +) -> tuple[int, int, int, bytes, bytes, str] | None: + """Parse wire slot, frame type, dtype, stream id, dst, call type from group DMRD.""" + if len(packet) < 20 or packet[:4] != b"DMRD": return None bits = packet[15] slot = 2 if (bits & 0x80) else 1 if bits & 0x40: - call_type = "unit" - elif (bits & 0x23) == 0x23: + return None + if (bits & 0x23) == 0x23: call_type = "vcsbk" else: call_type = "group" - return slot, int_id(packet[8:11]), call_type + frame_type = (bits & 0x30) >> 4 + dtype_vseq = bits & 0xF + return slot, frame_type, dtype_vseq, packet[16:20], packet[8:11], call_type def _system_has_active_bridge_leg( diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index ace70dc..8ac717a 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -42,6 +42,7 @@ from ..domain.dmr import bptc from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id from .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder from .talker_alias_use_cases import TalkerAliasUseCases +from .routing.peer_downlink_index import count_connected_peers from .routing.helpers import ( hbp_slot_blocks_group_voice, inject_only_defer_obp_hbp_slot_contention, @@ -617,11 +618,18 @@ class RoutingUseCases( _ts_st["TX_TYPE"] = HBPF_SLT_VTERM # Slot contention: active QSO blocks any other stream; post-VTERM uses GROUP_HANGTIME. _group_hangtime = float(_target_system.get("GROUP_HANGTIME", 0) or 0) + _tgt_peers = getattr(tgt_proto, "_peers", None) if tgt_proto else None + if _tgt_peers is None: + _tgt_peers = _target_system.get("PEERS", {}) + _target_connected = ( + count_connected_peers(_tgt_peers) if isinstance(_tgt_peers, dict) else 0 + ) _defer_slot_contention = inject_only_defer_obp_hbp_slot_contention( self._config, entry["SYSTEM"], _target_system, source_is_obp=source_is_obp, + connected_count=_target_connected, ) if ( not _closing_bridge_leg diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index c9349e0..5720686 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -45,10 +45,12 @@ from ...application.routing.helpers import ( clear_peer_rx_status_slots, clear_peer_ua_sessions, hbp_slot_blocks_group_voice_for_peer, + master_per_peer_slot_contention, + parse_dmrd_burst_fields, + peer_downlink_voice_slot, is_special_tg, is_unit_data_ingress, parse_dmrd_route_fields, - peer_downlink_voice_slot, peer_matches_rf_source, peer_should_receive_group_voice, peer_single_exclusive_tgid, @@ -241,6 +243,8 @@ class HBPProtocol(DatagramProtocol): self._config_push_delayed = None self._config_push_throttle = ConfigPushThrottle() self._repeat_downlink_report_start: dict[bytes, float] = {} + self._peer_voice_slots: dict[bytes, dict[int, dict[str, Any]]] = {} + self._peer_voice_hangtime: dict[bytes, dict[int, tuple[int, float]]] = {} self._refresh_connected_peer_count() else: self._peers = {} @@ -593,11 +597,17 @@ class HBPProtocol(DatagramProtocol): ) if voice_slot not in self.STATUS: return True - per_peer = self._inject_multi_peer_options_filter() + connected = self._cached_connected_peer_count() + per_peer = master_per_peer_slot_contention( + self._CONFIG, self._system, self._config, connected_count=connected, + ) if not per_peer: return True - stream_id = packet[4:8] if len(packet) >= 8 else b"" + burst = parse_dmrd_burst_fields(packet) + stream_id = burst[3] if burst is not None else b"" hang = float(self._config.get("GROUP_HANGTIME", 0) or 0) + pk = bytes_4(int_id(peer_id)) + hang_row = self._peer_voice_hangtime.get(pk, {}).get(voice_slot) return not hbp_slot_blocks_group_voice_for_peer( self.STATUS[voice_slot], peer_id, @@ -606,20 +616,122 @@ class HBPProtocol(DatagramProtocol): time.time(), hang, per_peer=per_peer, + peers=self._peers, + peer_slots=self._peer_voice_slots.get(pk), + peer_hang_row=hang_row, + voice_slot=voice_slot, + ) + + def _peer_voice_slot_pk(self, peer_id: bytes) -> bytes: + return bytes_4(int_id(peer_id)) + + def _touch_peer_voice_slot( + self, + peer_id: bytes, + voice_slot: int, + stream_id: bytes, + tgid: bytes, + pkt_time: float, + ) -> None: + pk = self._peer_voice_slot_pk(peer_id) + self._peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = { + "stream_id": stream_id, + "tgid": int_id(tgid), + "time": float(pkt_time), + } + self._peer_voice_hangtime.get(pk, {}).pop(int(voice_slot), None) + + def _end_peer_voice_slot( + self, + peer_id: bytes, + voice_slot: int, + stream_id: bytes, + pkt_time: float, + ) -> None: + pk = self._peer_voice_slot_pk(peer_id) + per_slot = self._peer_voice_slots.get(pk, {}) + active = per_slot.get(int(voice_slot)) + if not isinstance(active, dict): + return + if stream_id and active.get("stream_id") not in (stream_id, None): + return + ended = per_slot.pop(int(voice_slot), None) + if isinstance(ended, dict): + self._peer_voice_hangtime.setdefault(pk, {})[int(voice_slot)] = ( + int(ended.get("tgid", 0) or 0), + float(pkt_time), + ) + + def _track_peer_group_dmrd(self, peer_id: bytes, packet: bytes, *, pkt_time: float | None = None) -> None: + """Record per-hotspot downlink/ingress voice so a second stream is dropped.""" + burst = parse_dmrd_burst_fields(packet) + if burst is None: + return + wire_slot, frame_type, dtype_vseq, stream_id, dst_id, _call_type = burst + peer = self._peers.get(peer_id) + if peer is None: + return + voice_slot = peer_downlink_voice_slot( + peer, wire_slot, int_id(dst_id), self._config, peer_id=peer_id, + ) + now = time.time() if pkt_time is None else float(pkt_time) + if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: + self._end_peer_voice_slot(peer_id, voice_slot, stream_id, now) + return + self._touch_peer_voice_slot(peer_id, voice_slot, stream_id, dst_id, now) + + def _sync_peer_voice_from_ingress( + self, + peer_id: bytes, + wire_slot: int, + dst_id: bytes, + stream_id: bytes, + *, + call_type: str, + frame_type: int, + dtype_vseq: int, + pkt_time: float, + ) -> None: + if call_type not in ("group", "vcsbk"): + return + if not master_per_peer_slot_contention( + self._CONFIG, + self._system, + self._config, + connected_count=self._cached_connected_peer_count(), + ): + return + peer = self._peers.get(peer_id) + if peer is None: + return + voice_slot = peer_downlink_voice_slot( + peer, wire_slot, int_id(dst_id), self._config, peer_id=peer_id, ) + if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: + self._end_peer_voice_slot(peer_id, voice_slot, stream_id, pkt_time) + else: + self._touch_peer_voice_slot(peer_id, voice_slot, stream_id, dst_id, pkt_time) def send_peer(self, _peer: bytes, _packet: bytes) -> None: if _packet[:4] == DMRD: if not self._peer_should_receive_dmrd(_peer, _packet): return - if not self._peer_would_accept_group_dmrd(_peer, _packet): - return peer = self._peers.get(_peer) + route_pkt = _packet if peer is not None: - _packet = remap_dmrd_to_peer_static_slot( + route_pkt = remap_dmrd_to_peer_static_slot( _packet, peer, self._config, peer_id=_peer, ) - _packet = b"".join([_packet[:11], _peer, _packet[15:]]) + if not self._peer_would_accept_group_dmrd(_peer, route_pkt): + return + _packet = b"".join([route_pkt[:11], _peer, route_pkt[15:]]) + if master_per_peer_slot_contention( + self._CONFIG, + self._system, + self._config, + connected_count=self._cached_connected_peer_count(), + ): + self._track_peer_group_dmrd(_peer, _packet) self.transport.write(_packet, self._peers[_peer]["SOCKADDR"]) def _ta_buffer_enabled(self) -> bool: @@ -775,6 +887,8 @@ class HBPProtocol(DatagramProtocol): continue if not self._peer_should_receive_dmrd(peer, route_pkt): continue + if not self._peer_would_accept_group_dmrd(peer, route_pkt): + continue for pkt in packets: self.send_peer(peer, pkt) sent += 1 @@ -983,6 +1097,9 @@ class HBPProtocol(DatagramProtocol): if isinstance(sessions, dict): sessions.clear() clear_peer_rx_status_slots(self.STATUS, peer_id) + pk = self._peer_voice_slot_pk(peer_id) + self._peer_voice_slots.pop(pk, None) + self._peer_voice_hangtime.pop(pk, None) def _remove_peer(self, peer_id: bytes) -> None: self._on_peer_disconnected(peer_id) @@ -1230,6 +1347,16 @@ class HBPProtocol(DatagramProtocol): self.STATUS[_slot]["RX_TGID"] = _dst_id self.STATUS[_slot]["RX_TIME"] = pkt_time self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id + self._sync_peer_voice_from_ingress( + _peer_id, + _slot, + _dst_id, + _stream_id, + call_type=_call_type, + frame_type=_frame_type, + dtype_vseq=_dtype_vseq, + pkt_time=pkt_time, + ) _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): dmrpkt = _data[20:53] if len(_data) >= 53 else _data[20:] @@ -1654,6 +1781,16 @@ class HBPProtocol(DatagramProtocol): self.STATUS[_slot]["RX_TGID"] = _dst_id self.STATUS[_slot]["RX_TIME"] = pkt_time self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id + self._sync_peer_voice_from_ingress( + _peer_id, + _slot, + _dst_id, + _stream_id, + call_type=_call_type, + frame_type=_frame_type, + dtype_vseq=_dtype_vseq, + pkt_time=pkt_time, + ) _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): dmrpkt = _data[20:53] if len(_data) >= 53 else _data[20:] diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py index 644b238..4f5e64a 100644 --- a/tests/application/test_monitor_topology.py +++ b/tests/application/test_monitor_topology.py @@ -391,3 +391,44 @@ def test_obp_tx_fanout_suppressed_when_peer_slot_busy_on_other_tg() -> None: ) systems = {ev.split(",")[3] for ev in events} assert systems == {"SYSTEM-1"} + + +def test_obp_tx_fully_suppressed_when_all_receivers_busy() -> None: + """OBP downlink must not reach monitor when the sole eligible hotspot is slot-busy.""" + hs = bytes_4(7300444) + peers = {hs: _peer(options=b"TS2=730444,7144;")} + config = _proxy_config(peers) + peer_slots = {hs: 1} + now = time.time() + master_status = { + 2: { + "RX_TYPE": HBPF_SLT_VHEAD, + "TX_TYPE": HBPF_SLT_VTERM, + "RX_PEER": hs, + "RX_TGID": bytes_3(7144), + "RX_STREAM_ID": bytes_4(0x11111111), + "RX_TIME": now, + "TX_TIME": 0.0, + } + } + raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,73010,7000002,2,730444" + bridges = { + "730444": [ + { + "SYSTEM": "SYSTEM", + "TS": 2, + "TGID": 730444, + "ACTIVE": True, + "TO_TYPE": "ON", + } + ], + } + events = remap_inject_proxy_voice_events( + raw, + config, + config["SYSTEMS"], + peer_slots, + bridges, + master_status=master_status, + ) + assert events == [] diff --git a/tests/application/test_obp_hbp_cross_slot_downlink.py b/tests/application/test_obp_hbp_cross_slot_downlink.py index 8054eec..8ee2942 100644 --- a/tests/application/test_obp_hbp_cross_slot_downlink.py +++ b/tests/application/test_obp_hbp_cross_slot_downlink.py @@ -33,6 +33,22 @@ def test_defer_helper_requires_inject_only_obp_to_master() -> None: ) +def test_defer_helper_multi_peer_master_without_inject_proxy() -> None: + sys_cfg = { + "MODE": "MASTER", + "PEERS": { + bytes_4(1): {"CONNECTION": "YES"}, + bytes_4(2): {"CONNECTION": "YES"}, + }, + } + assert inject_only_defer_obp_hbp_slot_contention( + {}, "MASTER-A", sys_cfg, source_is_obp=True, connected_count=2, + ) + assert not inject_only_defer_obp_hbp_slot_contention( + {}, "MASTER-A", sys_cfg, source_is_obp=True, connected_count=1, + ) + + def _seed_busy_slot_2(scenario: DeterministicScenario, peer_id: int) -> None: proto = scenario.protocols["MASTER-A"] t = scenario.clock.time() diff --git a/tests/application/test_slot_contention.py b/tests/application/test_slot_contention.py index 32848b8..5c81364 100644 --- a/tests/application/test_slot_contention.py +++ b/tests/application/test_slot_contention.py @@ -9,8 +9,10 @@ from adn_server.application.routing.helpers import ( hbp_ingress_new_stream_collision, hbp_slot_blocks_group_voice, hbp_slot_blocks_group_voice_for_peer, + peer_hotspot_voice_slot_busy, slot_has_active_voice, slot_in_group_hangtime, + slot_status_hotspot_owner, ) from adn_server.domain import bytes_3, bytes_4 from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO @@ -71,6 +73,68 @@ def test_per_peer_scope_ignores_other_hotspot_busy_slot() -> None: ) is True +def test_per_peer_blocks_own_rx_when_bridge_tx_stream_stamped() -> None: + """OBP defer path stamps TX_STREAM on STATUS before send_peer; must still block.""" + now = 1_000_000.0 + peer = bytes_4(714002301) + slot = _active_rx_slot() + slot["RX_PEER"] = peer + slot["TX_STREAM_ID"] = _STREAM_B + slot["TX_PEER"] = bytes_4(73010) + slot["TX_TYPE"] = HBPF_SLT_VHEAD + slot["TX_TIME"] = now + 0.05 + peers = {peer: {"CONNECTION": "YES"}} + assert hbp_slot_blocks_group_voice_for_peer( + slot, peer, _TG_B, _STREAM_B, now + 0.1, 0.0, per_peer=True, peers=peers, + peer_slots={2: {"stream_id": _STREAM_A, "tgid": 7144, "time": now}}, + voice_slot=2, + ) is True + + +def test_peer_slot_session_blocks_other_tg_after_burst_gap() -> None: + """Per-peer session must survive DMR inter-burst gaps (> STREAM_TO).""" + now = 1_000_000.0 + hs = bytes_4(714002301) + slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} + peer_slots = { + 2: {"stream_id": _STREAM_A, "tgid": 7141, "time": now - STREAM_TO - 2.0}, + } + assert peer_hotspot_voice_slot_busy( + hs, 2, _STREAM_B, _TG_B, slot, peer_slots, None, now, 5.0, + ) + assert hbp_slot_blocks_group_voice_for_peer( + slot, hs, _TG_B, _STREAM_B, now, 5.0, per_peer=True, + peer_slots=peer_slots, voice_slot=2, + ) + + +def test_bridge_tx_peer_does_not_clear_hotspot_contention() -> None: + """Bridge TX_PEER=73010 must not disable per-hotspot slot busy checks.""" + now = 1_000_000.0 + hs = bytes_4(714002301) + slot = _active_rx_slot(stream_id=_STREAM_A) + slot["TX_PEER"] = bytes_4(73010) + slot["TX_STREAM_ID"] = _STREAM_A + slot["TX_TYPE"] = HBPF_SLT_VHEAD + slot["TX_TIME"] = now + peers = {hs: {"CONNECTION": "YES"}} + assert slot_status_hotspot_owner(slot, peers) is None + assert peer_hotspot_voice_slot_busy( + hs, 2, _STREAM_B, _TG_B, slot, + {2: {"stream_id": _STREAM_A, "tgid": 7144, "time": now}}, + None, now + 0.05, 5.0, + ) + + +def test_peer_hotspot_hangtime_blocks_other_tg() -> None: + now = 1_000_000.0 + hs = bytes_4(714002301) + slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} + assert peer_hotspot_voice_slot_busy( + hs, 2, _STREAM_B, _TG_B, slot, None, (7144, now), now + 2.0, 5.0, + ) + + def test_global_scope_still_blocks_any_peer_on_busy_slot() -> None: now = 1_000_000.0 peer_b = bytes_4(714002301)