diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index e447e2c..8b886a4 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -31,6 +31,7 @@ from adn_server.domain.hbp_protocol import STREAM_TO from .helpers import ( SIMPLEX_VOICE_SLOT, + _peer_ua_session_entry, clear_peer_ua_sessions, hbp_slot_blocks_group_voice_for_peer, is_special_tg, @@ -218,6 +219,7 @@ def peer_slot_blocks_downlink( 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 ctx.per_peer_contention(): return False hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0) @@ -307,7 +309,8 @@ def touch_peer_voice_slot( "tgid": int_id(tgid), "time": float(now), } - if ingress: + prev = ctx.peer_voice_slots.get(pk, {}).get(int(voice_slot)) + if ingress or (isinstance(prev, dict) and prev.get("ingress")): row["ingress"] = True ctx.peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = row if clear_hangtime: @@ -332,7 +335,9 @@ def end_peer_voice_slot( if isinstance(active, dict): active_stream = active.get("stream_id") if stream_id and active_stream and active_stream != stream_id: - return + active_time = float(active.get("time", 0) or 0) + if (now - active_time) < STREAM_TO: + return active = per_slot.pop(int(voice_slot), None) if not apply_hangtime: return @@ -375,7 +380,9 @@ def track_peer_group_dmrd( peer, voice_slot, ctx.sys_cfg, peer_id=peer_id, now=pkt_time, ) if locked is not None and locked == ended_tg: - clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot) + entry = _peer_ua_session_entry(ctx.sys_cfg, peer_id, voice_slot) + if not isinstance(entry, dict) or entry.get("source") != "local": + clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot) end_peer_voice_slot( ctx, peer_id, @@ -391,16 +398,32 @@ def track_peer_group_dmrd( and frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD and peer_wants_downlink_single_listen_lock(peer, ctx.sys_cfg) - and is_ua_session_tgid(int_id(dst_id)) ): - pk = bytes_4(int_id(peer_id)) - per_slot = ctx.peer_voice_slots.get(pk, {}) - active = per_slot.get(int(voice_slot)) - if not isinstance(active, dict) or active.get("stream_id") != stream_id: - now = time.time() if pkt_time is None else float(pkt_time) - register_peer_ua_session( - peer, peer_id, voice_slot, int_id(dst_id), ctx.sys_cfg, now=now, + dst_tgid = int_id(dst_id) + listen_lock_tg = is_ua_session_tgid(dst_tgid) or peer_receives_group_tgid( + peer, wire_slot, dst_tgid, + ) + if listen_lock_tg: + pk = bytes_4(int_id(peer_id)) + per_slot = ctx.peer_voice_slots.get(pk, {}) + active = per_slot.get(int(voice_slot)) + locked = peer_single_exclusive_tgid( + peer, voice_slot, ctx.sys_cfg, peer_id=peer_id, now=pkt_time, + ) + entry = _peer_ua_session_entry(ctx.sys_cfg, peer_id, voice_slot) + local_lock = ( + locked is not None + and int(locked) == int(dst_tgid) + and isinstance(entry, dict) + and entry.get("source") == "local" ) + if local_lock: + pass + elif not isinstance(active, dict) or active.get("stream_id") != stream_id: + now = time.time() if pkt_time is None else float(pkt_time) + register_peer_ua_session( + peer, peer_id, voice_slot, dst_tgid, ctx.sys_cfg, now=now, source="listen", + ) touch_peer_voice_slot( ctx, peer_id, diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 7069305..72710de 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -46,8 +46,13 @@ from __future__ import annotations import logging from hashlib import blake2b -from ...domain import int_id -from .helpers import hbp_ingress_new_stream_collision, master_per_peer_slot_contention +from ...domain import bytes_4, int_id +from .helpers import ( + group_voice_tg_ingress_collision, + hbp_ingress_downlink_session_blocks_tx, + hbp_ingress_new_stream_collision, + master_per_peer_slot_contention, +) from .peer_downlink_index import count_connected_peers logger = logging.getLogger(__name__) @@ -56,6 +61,63 @@ logger = logging.getLogger(__name__) class HbpForwardMixin: """routerHBP group voice ingress controls and sendDataToHBP.""" + def _ingress_drop_log_cache(self) -> set[tuple]: + cache = getattr(self, "_ingress_drop_logged", None) + if cache is None: + cache = set() + self._ingress_drop_logged = cache + return cache + + def _ingress_drop_key( + self, + kind: str, + system_name: str, + peer_id: bytes, + dst_id: bytes, + stream_id: bytes, + *, + slot: int | None = None, + ) -> tuple: + key: tuple = ( + kind, + system_name, + bytes_4(int_id(peer_id)), + dst_id, + stream_id, + ) + if slot is not None: + return (*key, int(slot)) + return key + + def _log_ingress_warning_once( + self, + key: tuple, + msg: str, + *args: object, + ) -> None: + cache = self._ingress_drop_log_cache() + if key in cache: + return + cache.add(key) + logger.warning(msg, *args) + + def _clear_ingress_drop_log( + self, + system_name: str, + peer_id: bytes, + dst_id: bytes, + stream_id: bytes, + slot: int, + ) -> None: + cache = self._ingress_drop_log_cache() + pk = bytes_4(int_id(peer_id)) + tg = dst_id + sid = stream_id + sl = int(slot) + for kind in ("slot_collision", "tg_busy", "downlink_tx"): + cache.discard((kind, system_name, pk, tg, sid)) + cache.discard((kind, system_name, pk, tg, sid, sl)) + def _hbp_group_voice_ingress_controls( self, system_name: str, @@ -81,13 +143,6 @@ 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["packets"] = 0 - _slot_st["loss"] = 0 - _slot_st["crcs"] = set() - _slot_st["LOOPLOG"] = False - _slot_st.pop("_bcsq", None) - _slot_st["lastSeq"] = False - _slot_st["lastData"] = False 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 @@ -97,11 +152,51 @@ class HbpForwardMixin: if hbp_ingress_new_stream_collision( _slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=per_peer, ): - logger.warning( + self._log_ingress_warning_once( + self._ingress_drop_key( + "slot_collision", system_name, peer_id, dst_id, stream_id, slot=slot, + ), "(%s) Packet received with STREAM ID: %s SUB: %s PEER: %s TGID %s, SLOT %s collided with existing call", system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, ) return False + 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 + from .downlink import normalize_ua_voice_slot + + peer = peers.get(peer_id) if isinstance(peers, dict) else None + if peer is None and isinstance(peers, dict): + peer = peers.get(bytes_4(int_id(peer_id))) + voice_slot = normalize_ua_voice_slot(peer, slot) if isinstance(peer, dict) else slot + peer_slots = getattr(src_proto, "_peer_voice_slots", {}).get(bytes_4(int_id(peer_id))) + if hbp_ingress_downlink_session_blocks_tx( + voice_slot, dst_id, peer_slots, + ): + self._log_ingress_warning_once( + self._ingress_drop_key( + "downlink_tx", system_name, peer_id, dst_id, stream_id, slot=voice_slot, + ), + "(%s) Packet dropped: peer %s already receiving on TG %s slot %s", + system_name, int_id(peer_id), int_id(dst_id), voice_slot, + ) + return False + self._clear_ingress_drop_log(system_name, peer_id, dst_id, stream_id, slot) + _slot_st["packets"] = 0 + _slot_st["loss"] = 0 + _slot_st["crcs"] = set() + _slot_st["LOOPLOG"] = False + _slot_st.pop("_bcsq", None) + _slot_st["lastSeq"] = False + _slot_st["lastData"] = False _slot_st["RX_START"] = pkt_time _slot_st["packets"] = _slot_st.get("packets", 0) + 1 _pkts = _slot_st["packets"] diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 7ef30fc..e32bca7 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -271,24 +271,42 @@ def _peer_status_rx_hangtime_blocks( return (pkt_time - rx_t) < hang -def peer_cross_static_tg_allowed_on_slot( +def peer_single_same_tg_foreign_tx_blocks( peer: dict[str, Any], - voice_slot: int, + peer_id: bytes, incoming_tgid_b: bytes, + stream_id: bytes, + slot_st: dict[str, Any], sys_cfg: dict[str, Any] | None, *, - peer_id: bytes | None = None, - now: float | None = None, + pkt_time: float, ) -> bool: - """True when another static OPTIONS TG may share this RF slot (SINGLE=0 or no UA lock).""" - if not sys_cfg: + """SINGLE=1 UA on TG T: block another peer's stream on the same TG (slot busy).""" + if not sys_cfg or not peer_single_mode(peer, sys_cfg): return False incoming = int_id(incoming_tgid_b) - if not peer_receives_group_tgid(peer, voice_slot, incoming): + 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 + if locked is None or int(locked) != incoming: return False - return not peer_single_blocks_group_voice( - peer, voice_slot, incoming, sys_cfg, peer_id=peer_id, now=now, - ) + if not slot_has_active_voice(slot_st, pkt_time): + return False + slot_tg = int_id(slot_st.get("RX_TGID", b"\x00\x00\x00")) + if slot_tg != incoming: + return False + owner = slot_st.get("RX_PEER") or slot_st.get("TX_PEER") + pk = bytes_4(int_id(peer_id)) + if owner is None or int_id(owner) == 0 or bytes_4(int_id(owner)) == pk: + return False + leg_stream = slot_st.get("RX_STREAM_ID") or slot_st.get("TX_STREAM_ID") + if stream_id and leg_stream == stream_id: + return False + return True def peer_hotspot_voice_slot_busy( @@ -306,78 +324,38 @@ def peer_hotspot_voice_slot_busy( peer: dict[str, Any] | None = None, sys_cfg: dict[str, Any] | None = None, ) -> bool: - """True when this hotspot must not receive another group stream on ``voice_slot``.""" + """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. + """ pk = bytes_4(int_id(peer_id)) if peer is None and peers is not None: peer = peers.get(peer_id) or peers.get(pk) - cross_static = ( - isinstance(peer, dict) - and peer_cross_static_tg_allowed_on_slot( - peer, voice_slot, incoming_tgid_b, sys_cfg, peer_id=peer_id, now=pkt_time, - ) - ) - # Ingress-owned post-VTERM window: must win over OBP bridge TX stamp on STATUS[slot]. if _peer_transmit_hangtime_blocks(hang_row, incoming_tgid_b, pkt_time, group_hangtime): return True active = (peer_slots or {}).get(int(voice_slot)) - if isinstance(active, dict) and active.get("ingress"): - active_stream = active.get("stream_id") - if stream_id and active_stream == stream_id: - return False - if cross_static: - return False - # Local RF TX: block foreign streams even when bridge TX stamp matches. - return True if isinstance(active, dict): - active_stream = active.get("stream_id") - if stream_id and active_stream == stream_id: - return False - if stream_id and active_stream and active_stream != stream_id: - active_time = float(active.get("time", 0) or 0) - if (pkt_time - active_time) < STREAM_TO: - active_tg = int(active.get("tgid", 0) or 0) - incoming_tg = int_id(incoming_tgid_b) - if ( - active_tg - and incoming_tg - and active_tg != incoming_tg - and not active.get("ingress") - ): - if cross_static: - return False - owner = slot_status_hotspot_owner(slot_st, peers) - rx_listening = ( - owner is not None - and bytes_4(int_id(owner)) == pk - and slot_has_active_voice(slot_st, pkt_time) - and int_id(slot_st.get("RX_TGID", b"")) == active_tg - ) - if not rx_listening: - return True - tx_matches = ( - stream_id == slot_st.get("TX_STREAM_ID") - and slot_st.get("TX_PEER") is not None - and int_id(slot_st.get("TX_PEER")) != 0 - and (peers is None or not peer_key_in_peers(slot_st.get("TX_PEER"), peers)) - ) - if not tx_matches: - return True - # OBP bridge TX stamp on shared STATUS[slot] overrides stale downlink-only session rows. + 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 + 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, + ): + return True + if isinstance(peer, dict) and peer_single_same_tg_foreign_tx_blocks( + 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 isinstance(active, dict): - active_stream = active.get("stream_id") - if stream_id and active_stream == stream_id: - return False - if cross_static: - 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 @@ -507,26 +485,153 @@ def hbp_ingress_new_stream_collision( ) -> bool: """True when a new group-voice stream must drop on ingress (legacy routerHBP). - Legacy blocks only when the prior stream is still open (STREAM_TO), the RF - source differs (another subscriber), and — in inject-only — the slot row - belongs to this hotspot. Same-subscriber rekey with a new stream id is allowed. + 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). """ + 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 slot_st.get("RX_TYPE") == HBPF_SLT_VTERM: + if stream_id and stream_id == slot_st.get("TX_STREAM_ID"): return False - rx_time = float(slot_st.get("RX_TIME", 0) or 0) - if pkt_time >= rx_time + STREAM_TO: + for leg in ("RX", "TX"): + type_key = f"{leg}_TYPE" + time_key = f"{leg}_TIME" + rfs_key = f"{leg}_RFS" + stream_key = f"{leg}_STREAM_ID" + dtype = slot_st.get(type_key) + if dtype is None or dtype == HBPF_SLT_VTERM: + continue + leg_time = float(slot_st.get(time_key, 0) or 0) + if pkt_time >= leg_time + STREAM_TO: + continue + if stream_id and stream_id == slot_st.get(stream_key): + continue + prev_rfs = slot_st.get(rfs_key, b"\x00\x00\x00") + if int_id(rf_src) != 0 and bytes_4(int_id(rf_src)) == bytes_4(int_id(prev_rfs)): + continue + return True + return False + + +def _same_rf_source(a: bytes, b: bytes) -> bool: + return int_id(a) != 0 and bytes_4(int_id(a)) == bytes_4(int_id(b)) + + +def _hbp_slot_active_tgid(slot_st: dict[str, Any], pkt_time: float) -> bytes | None: + if not slot_has_active_voice(slot_st, pkt_time): + return None + rx_type = slot_st.get("RX_TYPE") + if rx_type is not None and rx_type != HBPF_SLT_VTERM: + return slot_st.get("RX_TGID") + tx_type = slot_st.get("TX_TYPE") + if tx_type is not None and tx_type != HBPF_SLT_VTERM: + return slot_st.get("TX_TGID") or slot_st.get("RX_TGID") + return slot_st.get("RX_TGID") + + +def _obp_stream_active(st: dict[str, Any], pkt_time: float) -> bool: + if st.get("_fin"): return False - prev_rfs = slot_st.get("RX_RFS", b"\x00\x00\x00") - if int_id(rf_src) != 0 and bytes_4(int_id(rf_src)) == bytes_4(int_id(prev_rfs)): + start = float(st.get("START", 0) or 0) + if start <= 0 or start + 180 < pkt_time: + return False + return True + + +def group_voice_tg_ingress_collision( + protocols: dict[str, Any], + systems_cfg: dict[str, Any], + tgid_b: bytes, + stream_id: bytes, + rf_src: bytes, + pkt_time: float, +) -> bool: + """True when another active group-voice leg already owns this TG (HBP or OBP).""" + 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 + 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 hbp_ingress_downlink_session_blocks_tx( + voice_slot: int, + incoming_tgid_b: bytes, + peer_slots: PeerVoiceSlotMap | None, +) -> bool: + """True when a hotspot mid downlink QSO on this TG/slot must not ingress TX.""" + active = (peer_slots or {}).get(int(voice_slot)) + if not isinstance(active, dict) or active.get("ingress"): + return False + active_tg = int(active.get("tgid", 0) or 0) + incoming_tg = int_id(incoming_tgid_b) + return bool(active_tg and incoming_tg and active_tg == incoming_tg) + + +def hbp_master_ingress_repeat_allowed( + slot_st: dict[str, Any], + peer_id: bytes, + rf_src: bytes, + dst_id: bytes, + stream_id: bytes, + pkt_time: float, + *, + protocols: dict[str, Any] | None = None, + systems_cfg: dict[str, Any] | None = None, +) -> bool: + """True when MASTER REPEAT may fan this ingress packet to other peers.""" + if stream_id and stream_id == slot_st.get("RX_STREAM_ID"): + owner = slot_st.get("RX_PEER") + return ( + owner is not None + and int_id(owner) != 0 + and bytes_4(int_id(owner)) == bytes_4(int_id(peer_id)) + ) + if hbp_ingress_new_stream_collision( + slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=False, + ): + return False + if protocols and systems_cfg and group_voice_tg_ingress_collision( + protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time, + ): return False - if per_peer: - owner = slot_status_peer_owner(slot_st) - if owner is not None and bytes_4(int_id(owner)) != bytes_4(int_id(peer_id)): - return False return True @@ -799,8 +904,10 @@ def _write_peer_ua_session( tgid: int, expires: float, sys_cfg: dict[str, Any], + *, + source: str = "local", ) -> None: - entry = {"tgid": int(tgid), "expires": float(expires)} + entry = {"tgid": int(tgid), "expires": float(expires), "source": str(source)} pk = bytes_4(int_id(peer_id)) sys_cfg.setdefault("_PEER_UA_SESSIONS", {}).setdefault(pk, {})[slot] = entry peer.setdefault("_UA_SESSION", {})[slot] = entry @@ -880,6 +987,7 @@ def register_peer_ua_session( sys_cfg: dict[str, Any], *, now: float | None = None, + source: str = "local", ) -> None: """Track UA TG for this hotspot (SINGLE=1 exclusive; SINGLE=0 multi-dynamic set).""" if not is_ua_session_tgid(tgid): @@ -906,6 +1014,7 @@ def register_peer_ua_session( int(tgid), expires_at, sys_cfg, + source=source, ) @@ -1179,6 +1288,41 @@ def peer_single_blocks_group_voice( return False +def peer_single_blocks_foreign_same_tg_downlink( + peer: dict[str, Any], + peer_id: bytes, + voice_slot: int, + incoming_tgid_b: bytes, + peer_slots: PeerVoiceSlotMap | None, + sys_cfg: dict[str, Any] | None, + *, + now: float | None = None, +) -> bool: + """SINGLE=1 local UA on TG T: block network downlink on T unless hotspot is TX on T.""" + 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, + ) + 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)) + if isinstance(active, dict) and active.get("ingress"): + if int(active.get("tgid", 0) or 0) == incoming: + return False + return True + + def peer_static_options_tg_count(peer: dict[str, Any]) -> int: """Count distinct static group TGs listed in peer OPTIONS (TS1 ∪ TS2).""" from adn_server.application.report.payloads import parse_peer_options_static @@ -1189,11 +1333,20 @@ def peer_static_options_tg_count(peer: dict[str, Any]) -> int: def peer_wants_downlink_single_listen_lock(peer: dict[str, Any], sys_cfg: dict[str, Any]) -> bool: """SINGLE=1 downlink listen lock for overlap — not full-table lab witnesses.""" + if not sys_cfg: + return False if not peer_single_mode(peer, sys_cfg): return False 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). diff --git a/src/adn_server/application/routing/obp_forward.py b/src/adn_server/application/routing/obp_forward.py index 4c10c4a..66cd09f 100644 --- a/src/adn_server/application/routing/obp_forward.py +++ b/src/adn_server/application/routing/obp_forward.py @@ -52,6 +52,7 @@ from typing import Any from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, int_id from ...domain.dmr import decode from ...domain.dmr.const import LC_OPT +from .helpers import group_voice_tg_ingress_collision logger = logging.getLogger(__name__) @@ -162,6 +163,17 @@ class ObpForwardMixin: return True if stream_id not in status: + 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 OBP stream %s from peer %s", + system_name, int_id(dst_id), int_id(stream_id), int_id(peer_id), + ) + return False st: dict[str, Any] = { "START": pkt_time, "CONTENTION": False, diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 590264d..13ed6be 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -52,6 +52,7 @@ from ...application.routing.downlink import ( from ...application.routing.helpers import ( clear_peer_rx_status_slots, clear_peer_ua_sessions, + hbp_master_ingress_repeat_allowed, is_on_demand_service_dst, is_server_originated_voice, is_special_tg, @@ -576,6 +577,17 @@ class HBPProtocol(DatagramProtocol): if len(packet) >= 8 and peer_matches_rf_source(peer_id, packet[5:8], self._peers): return True return self._cached_connected_peer_count() == 1 + if len(packet) >= 8 and peer_matches_rf_source(peer_id, packet[5:8], self._peers): + ctx = self._downlink_ctx() + if ctx.per_peer_contention(): + peer = self._peers[peer_id] + voice_slot = peer_downlink_voice_slot( + peer, slot, tgid, self._config, peer_id=peer_id, + ) + pk = bytes_4(int_id(peer_id)) + active = ctx.peer_voice_slots.get(pk, {}).get(int(voice_slot)) + if isinstance(active, dict) and active.get("ingress"): + return False connected = self._cached_connected_peer_count() store = self._get_subscription_store() if self._get_subscription_store else None return peer_should_receive_group_voice( @@ -626,7 +638,13 @@ class HBPProtocol(DatagramProtocol): ctx, peer_id, voice_slot, stream_id, dst_id, pkt_time=pkt_time, ingress=True, ) - def _peer_would_accept_group_dmrd(self, peer_id: bytes, packet: bytes) -> bool: + def _peer_would_accept_group_dmrd( + self, + peer_id: bytes, + packet: bytes, + *, + routed: bool = False, + ) -> bool: """True when group/vcsbk DMRD passes OPTIONS filter and per-hotspot slot gate.""" if packet[:4] != DMRD or not self._peer_should_receive_dmrd(peer_id, packet): return False @@ -636,7 +654,11 @@ class HBPProtocol(DatagramProtocol): ctx = self._downlink_ctx() if not ctx.per_peer_contention(): return True - route_pkt = remap_dmrd_for_peer(packet, peer, self._config, peer_id=peer_id) + route_pkt = ( + packet + if routed + else remap_dmrd_for_peer(packet, peer, self._config, peer_id=peer_id) + ) 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]: @@ -666,7 +688,9 @@ class HBPProtocol(DatagramProtocol): route_pkt = remap_dmrd_for_peer( _packet, peer, self._config, peer_id=_peer, ) - if not self._peer_would_accept_group_dmrd(_peer, route_pkt): + if not self._peer_would_accept_group_dmrd( + _peer, _packet if not _skip_dual_expand else route_pkt, routed=_skip_dual_expand, + ): self._log_downlink_drop_once(_peer, route_pkt) return self._downlink_drop_logged.discard(self._downlink_drop_key(_peer, route_pkt)) @@ -1034,6 +1058,9 @@ class HBPProtocol(DatagramProtocol): sessions = peer.get("_UA_SESSION") if isinstance(sessions, dict): sessions.clear() + pk = bytes_4(int_id(peer_id)) + self._peer_voice_slots.pop(pk, None) + self._peer_voice_hangtime.pop(pk, None) clear_peer_rx_status_slots(self.STATUS, peer_id) def _remove_peer(self, peer_id: bytes) -> None: @@ -1123,16 +1150,6 @@ class HBPProtocol(DatagramProtocol): if sub_map is not None: sub_map[_rf_src] = (self._system, _slot, pkt_time) self.note_dmrd_stream(_peer_id, _rf_src, _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, - ) if ( _call_type in ("group", "vcsbk") and _frame_type == HBPF_DATA_SYNC @@ -1176,8 +1193,37 @@ class HBPProtocol(DatagramProtocol): self.store_ta_from_voice_burst( _peer_id, _rf_src, _stream_id, _dtype_vseq, _data[20:53], ) + if _call_type in ("group", "vcsbk"): + 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, + ) + _slot_st = self.STATUS.get(_slot, {}) + _protocols: dict[str, Any] = {} + if self._router and getattr(self._router, "_get_protocols", None): + _protocols = self._router._get_protocols() or {} + _repeat_ok = ( + _call_type not in ("group", "vcsbk") + or hbp_master_ingress_repeat_allowed( + _slot_st, + _peer_id, + _rf_src, + _dst_id, + _stream_id, + pkt_time, + protocols=_protocols, + systems_cfg=self._CONFIG.get("SYSTEMS", {}), + ) + ) if ( - self._config.get("REPEAT", True) + _repeat_ok + and self._config.get("REPEAT", True) and _call_type in ("group", "vcsbk") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD @@ -1190,7 +1236,11 @@ class HBPProtocol(DatagramProtocol): self._on_talker_alias_local_repeat( self._system, _peer_id, _rf_src, _stream_id, ) - if self._config.get("REPEAT", True) and _call_type in ("group", "vcsbk"): + if ( + _repeat_ok + and self._config.get("REPEAT", True) + and _call_type in ("group", "vcsbk") + ): _repeat_tail = _data[15:] if ( _dtype_vseq in (1, 2, 3, 4) @@ -1226,17 +1276,6 @@ class HBPProtocol(DatagramProtocol): _call_type, _dtype_vseq, _stream_id, self.STATUS.get(_slot, {}).get("RX_STREAM_ID") if _slot in self.STATUS else None, ) - if _slot in self.STATUS and not _unit_data: - if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): - 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"] - except Exception: - self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src - else: - self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src _accepted = False if self._dmrd_received: _accepted = self._dmrd_received( @@ -1244,6 +1283,25 @@ class HBPProtocol(DatagramProtocol): _call_type, _frame_type, _dtype_vseq, _stream_id, _data, ingress_pkt_time=pkt_time, ) + if _accepted: + if _slot in self.STATUS and not _unit_data: + if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"): + 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"] + except Exception: + 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 + 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): dmrpkt = _data[20:53] if len(_data) >= 53 else _data[20:] @@ -1260,7 +1318,8 @@ class HBPProtocol(DatagramProtocol): ): reactor.callInThread(self._on_play_file_request, str(_int_dst_id), self._system) if ( - _call_type in ("group", "vcsbk") + _accepted + and _call_type in ("group", "vcsbk") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM and _slot in self.STATUS @@ -1269,21 +1328,13 @@ class HBPProtocol(DatagramProtocol): ): self._on_in_band_signalling(self._system, _slot, _dst_id, pkt_time) if ( - _call_type in ("group", "vcsbk") + _accepted + and _call_type in ("group", "vcsbk") and _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VTERM and self._on_talker_alias_stream_end ): self._on_talker_alias_stream_end(self._system, _stream_id) - # Legacy routerHBP: unit data (_data_call) does not mark slot RX busy. - if _slot in self.STATUS and not _unit_data: - 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 - self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id elif _command == RPTL: _peer_id = _data[4:8] diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index b70d0fd..72db18d 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -25,6 +25,7 @@ from __future__ import annotations from adn_server.application.routing.helpers import ( clear_peer_rx_status_slots, clear_peer_ua_sessions, + peer_hotspot_voice_slot_busy, peer_options_static_tg_slot, peer_receives_group_tgid, peer_should_receive_group_voice, @@ -509,3 +510,59 @@ def test_single_lock_persists_after_vterm_until_timer_expires() -> None: assert peer_should_receive_group_voice( peer, 2, 730, peer_id=peer_id, connected_count=3, sys_cfg=sys_cfg, now=now + 61, ) + + +def test_single_blocks_foreign_same_tg_while_local_ua() -> None: + """J39JQ UA on 730502: slot gate blocks HP3ICC downlink on the same TG.""" + from adn_server.application.routing.helpers import peer_single_blocks_foreign_same_tg_downlink + + peer = {"OPTIONS": b"TS2=730500,730508;SINGLE=1;TIMER=60;"} + sys_cfg = _sys_cfg() + peer_id = _peer_id() + now = 1_000_000.0 + register_peer_ua_session(peer, peer_id, 2, 730502, sys_cfg, now=now) + assert peer_single_blocks_foreign_same_tg_downlink( + peer, peer_id, 2, bytes_3(730502), None, sys_cfg, now=now + 10, + ) + assert peer_hotspot_voice_slot_busy( + peer_id, + 2, + bytes_4(0x22222222), + bytes_3(730502), + {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}, + None, + None, + now + 0.1, + 5.0, + peer=peer, + sys_cfg=sys_cfg, + ) + + +def test_downlink_vterm_does_not_clear_local_ua_session() -> None: + """Network VTERM on session TG must not wipe a local PTT lock (TIMER session).""" + from adn_server.application.routing.downlink import DownlinkContext, track_peer_group_dmrd + + sys_cfg = {"SINGLE_MODE": False, "DEFAULT_UA_TIMER": 10, "MODE": "MASTER", "MAX_PEERS": 8} + config = {"PROXY": {"TARGET_SYSTEM": "MASTER-A"}, "SYSTEMS": {"MASTER-A": sys_cfg}} + peer_id = _peer_id() + peer = {"OPTIONS": b"TS2=730500,730508;SINGLE=1;TIMER=60;"} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={peer_id: peer}, + status={1: {}, 2: {}}, + connected_count=3, + ) + now = 1_000_000.0 + register_peer_ua_session(peer, peer_id, 2, 730502, sys_cfg, now=now, source="local") + stream = bytes_4(0x11111111) + vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(100), 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, peer_id, vterm, peer, pkt_time=now + 8) + assert not peer_should_receive_group_voice( + peer, 2, 730500, peer_id=peer_id, connected_count=3, sys_cfg=sys_cfg, now=now + 9, + ) diff --git a/tests/application/test_slot_contention.py b/tests/application/test_slot_contention.py index a7d39fb..27fa8bf 100644 --- a/tests/application/test_slot_contention.py +++ b/tests/application/test_slot_contention.py @@ -5,10 +5,14 @@ from __future__ import annotations from adn_server.application.routing.helpers import ( + group_voice_tg_ingress_collision, + hbp_ingress_downlink_session_blocks_tx, hbp_ingress_new_stream_collision, + hbp_master_ingress_repeat_allowed, hbp_slot_blocks_group_voice, hbp_slot_blocks_group_voice_for_peer, peer_hotspot_voice_slot_busy, + register_peer_ua_session, slot_has_active_voice, slot_in_group_hangtime, slot_status_hotspot_owner, @@ -57,6 +61,142 @@ 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.""" + 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( + slot, + bytes_4(730039265), + bytes_3(730039265), + _STREAM_B, + now + 0.1, + per_peer=True, + ) + + +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.""" + now = 1_000_000.0 + slot = { + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": now, + "TX_RFS": bytes_3(730039256), + "TX_STREAM_ID": _STREAM_A, + } + assert hbp_ingress_new_stream_collision( + slot, + bytes_4(730039267), + bytes_3(730039267), + _STREAM_B, + now + 0.1, + per_peer=True, + ) + + +def test_ingress_downlink_session_blocks_tx_on_same_tg() -> None: + peer_slots = { + 1: {"stream_id": _STREAM_A, "tgid": 730500, "time": 1_000_000.0}, + } + assert hbp_ingress_downlink_session_blocks_tx(1, bytes_3(730500), peer_slots) + assert not hbp_ingress_downlink_session_blocks_tx(1, bytes_3(730502), peer_slots) + + +def test_ingress_downlink_session_allows_tx_after_own_ingress() -> None: + peer_slots = { + 2: {"stream_id": _STREAM_A, "tgid": 730502, "time": 1_000_000.0, "ingress": True}, + } + assert not hbp_ingress_downlink_session_blocks_tx(2, bytes_3(730502), peer_slots) + + +def test_master_repeat_denied_for_colliding_stream() -> None: + now = 1_000_000.0 + slot = _active_rx_slot() + slot["RX_RFS"] = bytes_3(730039264) + slot["RX_PEER"] = bytes_4(730039264) + assert not hbp_master_ingress_repeat_allowed( + slot, + bytes_4(730039265), + bytes_3(730039265), + _TG_A, + _STREAM_B, + now + 0.1, + ) + + +def test_master_repeat_allowed_for_slot_owner_continuation() -> None: + now = 1_000_000.0 + slot = _active_rx_slot() + slot["RX_PEER"] = bytes_4(730039264) + assert hbp_master_ingress_repeat_allowed( + slot, + bytes_4(730039264), + bytes_3(730039264), + _TG_A, + _STREAM_A, + now + 0.1, + ) + + +def test_group_voice_tg_collision_rejects_second_hbp_stream() -> None: + now = 1_000_000.0 + master_status = { + 2: _active_rx_slot(tgid=_TG_A, 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 group_voice_tg_ingress_collision( + protocols, systems, _TG_A, _STREAM_B, bytes_3(730039265), now + 0.1, + ) + + +def test_group_voice_tg_collision_rejects_obp_while_hbp_active() -> None: + now = 1_000_000.0 + master_status = {1: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now)} + master_status[1]["RX_RFS"] = bytes_3(730039266) + obp_status: dict[bytes, dict] = {} + protocols = { + "M1": type("P", (), {"STATUS": master_status})(), + "OBP": type("P", (), {"STATUS": obp_status})(), + } + systems = {"M1": {"MODE": "MASTER"}, "OBP": {"MODE": "OPENBRIDGE"}} + assert group_voice_tg_ingress_collision( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039256), now + 0.1, + ) + + +def test_group_voice_tg_collision_across_hbp_slots() -> None: + """Active TG on slot 1 must block a new stream on slot 2 (dual-slot OBP peers).""" + now = 1_000_000.0 + master_status = { + 1: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now), + } + master_status[1]["RX_RFS"] = bytes_3(730039266) + protocols = {"M1": type("P", (), {"STATUS": master_status})()} + systems = {"M1": {"MODE": "MASTER"}} + assert group_voice_tg_ingress_collision( + protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039267), 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 + hs = bytes_4(730039265) + slot = _active_rx_slot() + slot["RX_PEER"] = bytes_4(730039264) + peer_slots = { + 2: {"stream_id": _STREAM_B, "tgid": 730502, "time": now, "ingress": True}, + } + assert peer_hotspot_voice_slot_busy( + hs, 2, _STREAM_A, _TG_A, slot, peer_slots, None, now + 0.05, 5.0, + ) + + def test_per_peer_scope_ignores_other_hotspot_busy_slot() -> None: """Inject-only: peer B must not inherit slot contention from peer A.""" now = 1_000_000.0 @@ -91,8 +231,8 @@ def test_bridge_tx_stamp_same_stream_allows_obp_downlink() -> None: ) -def test_per_peer_obp_tx_stamp_clears_stale_peer_slot_block() -> None: - """OBP TX stamp + same stream must deliver despite stale per-peer session row.""" +def test_per_peer_obp_tx_stamp_blocked_when_other_stream_active() -> None: + """OBP TX stamp must not override an active per-peer session on another stream.""" now = 1_000_000.0 peer = bytes_4(714002301) slot = _active_rx_slot() @@ -106,7 +246,63 @@ def test_per_peer_obp_tx_stamp_clears_stale_peer_slot_block() -> None: 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 False + ) 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.""" + now = 1_000_000.0 + witness_id = bytes_4(730039257) + 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}, + } + slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} + assert not peer_hotspot_voice_slot_busy( + witness_id, + 2, + _STREAM_B, + _TG_B, + slot, + peer_slots, + None, + now, + 5.0, + peer=witness, + sys_cfg={"GROUP_HANGTIME": 5.0}, + ) + + +def test_same_stream_vterm_not_blocked_by_peer_slot_busy() -> None: + """VTERM for the active stream must reach the hotspot to clear peer_voice_slots.""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + touch_peer_voice_slot, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8} + config = {"PROXY": {"TARGET_SYSTEM": "MASTER-A"}, "SYSTEMS": {"MASTER-A": sys_cfg}} + hs = bytes_4(714002301) + peer = {"OPTIONS": b"TS2=7141,71442;SINGLE=1;"} + 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=2, + ) + now = 1_000_000.0 + stream = bytes_4(0x11111111) + touch_peer_voice_slot(ctx, hs, 2, stream, bytes_3(7141), pkt_time=now) + vterm = b"".join([ + b"DMRD", b"\x00", bytes_3(100), bytes_3(7141), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), + stream, + ] + [b"\x00"] * 33) + assert not peer_slot_blocks_downlink(ctx, hs, peer, vterm, pkt_time=now + 0.5) def test_peer_slot_session_blocks_other_tg_after_burst_gap() -> None: @@ -151,6 +347,42 @@ def test_obp_tx_stamp_does_not_override_peer_hangtime() -> None: ) +def test_downlink_track_preserves_ingress_tx_flag() -> None: + """Delivered downlink DMRD must not clear ingress while hotspot is still TX.""" + from adn_server.application.routing.downlink import ( + DownlinkContext, + peer_slot_blocks_downlink, + touch_peer_voice_slot, + track_peer_group_dmrd, + ) + + sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8} + config = {"SYSTEMS": {"MASTER-A": sys_cfg}} + peer_id = bytes_4(730039264) + peer = {"OPTIONS": b"TS2=730502;SINGLE=1;"} + ctx = DownlinkContext( + config=config, + system_name="MASTER-A", + sys_cfg=sys_cfg, + peers={peer_id: peer}, + status={1: {}, 2: {}}, + connected_count=2, + ) + now = 1_000_000.0 + stream_a = bytes_4(0x11111111) + stream_b = bytes_4(0x22222222) + touch_peer_voice_slot( + ctx, peer_id, 2, stream_a, bytes_3(730502), pkt_time=now, ingress=True, + ) + foreign_vhead = b"".join([ + b"DMRD", b"\x00", bytes_3(730039265), bytes_3(730502), b"\x00\x00\x00\x00", + bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), stream_b, + ] + [b"\x00"] * 33) + assert peer_slot_blocks_downlink(ctx, peer_id, peer, foreign_vhead, pkt_time=now + 2.1) + track_peer_group_dmrd(ctx, peer_id, foreign_vhead, peer, pkt_time=now + 2) + assert ctx.peer_voice_slots[peer_id][2].get("ingress") is True + + def test_downlink_vterm_does_not_reset_transmit_hangtime() -> None: """Delivered OBP VTERM must not replace ingress GROUP_HANGTIME ownership.""" from adn_server.application.routing.downlink import ( @@ -210,8 +442,8 @@ def test_obp_tx_stamp_does_not_bypass_fresh_chile_downlink_session() -> None: ) -def test_obp_bridge_tx_overrides_stale_peer_slot_session() -> None: - """Stale HS session must not block OBP bridged downlink on the same TG.""" +def test_obp_bridge_tx_blocked_when_other_stream_session_active() -> None: + """Active per-peer session on another stream blocks OBP bridged downlink (one QSO per slot).""" now = 1_000_000.0 hs = bytes_4(0x2B83833D) slot = { @@ -222,7 +454,7 @@ def test_obp_bridge_tx_overrides_stale_peer_slot_session() -> None: "RX_TYPE": HBPF_SLT_VTERM, } peers = {hs: {"CONNECTION": "YES"}} - assert not peer_hotspot_voice_slot_busy( + assert peer_hotspot_voice_slot_busy( hs, 2, _STREAM_B, _TG_B, slot, {2: {"stream_id": _STREAM_A, "tgid": 7305, "time": now - 5.0}}, None, now, 5.0, peers=peers, @@ -256,8 +488,8 @@ def test_peer_hotspot_hangtime_blocks_other_tg() -> None: ) -def test_single0_ingress_tx_allows_other_static_tg() -> None: - """SINGLE=0: local TX on one static TG must not block RX on another in OPTIONS.""" +def test_single0_ingress_tx_blocks_other_static_tg_on_slot() -> None: + """SINGLE=0: local TX on one TG blocks any other downlink on the same RF slot.""" now = 1_000_000.0 hs = bytes_4(730039253) peer = {"OPTIONS": b"TS2=730507,730508;SINGLE=0;"} @@ -271,7 +503,7 @@ def test_single0_ingress_tx_allows_other_static_tg() -> None: "ingress": True, }, } - assert not peer_hotspot_voice_slot_busy( + assert peer_hotspot_voice_slot_busy( hs, 2, bytes_4(0x22222222), @@ -286,8 +518,8 @@ def test_single0_ingress_tx_allows_other_static_tg() -> None: ) -def test_monitor_peer_allows_second_static_tg_during_listen() -> None: - """Lab witness (many static TGs) must hear concurrent calls on different TGs.""" +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.""" now = 1_000_000.0 hs = bytes_4(730039257) peer = { @@ -319,6 +551,32 @@ def test_monitor_peer_allows_second_static_tg_during_listen() -> None: ) +def test_ingress_tx_blocks_same_stream_obp_echo() -> None: + """OBP loopback with same stream_id must not downlink while hotspot is still TX.""" + now = 1_000_000.0 + hs = bytes_4(730039264) + stream = bytes_4(0x22222222) + slot = { + "TX_PEER": bytes_4(73010), + "TX_STREAM_ID": stream, + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": now, + "RX_TYPE": HBPF_SLT_VTERM, + } + peers = {hs: {"CONNECTION": "YES", "OPTIONS": b"TS2=730502;SINGLE=1;"}} + peer_slots = { + 2: { + "stream_id": stream, + "tgid": 730502, + "time": now, + "ingress": True, + }, + } + assert peer_hotspot_voice_slot_busy( + hs, 2, stream, bytes_3(730502), slot, peer_slots, None, now, 5.0, peers=peers, + ) + + def test_ingress_tx_blocks_same_tg_foreign_stream_despite_bridge_stamp() -> None: """Local RF TX must block downlink even when OBP bridge TX stamp matches incoming.""" now = 1_000_000.0 @@ -344,6 +602,67 @@ def test_ingress_tx_blocks_same_tg_foreign_stream_despite_bridge_stamp() -> None ) +def test_single1_listen_blocks_other_static_tg_during_downlink() -> None: + """Lab J39JQ: SINGLE=1 listening on 730502 must not RX 730500 on same slot.""" + now = 1_000_000.0 + hs = bytes_4(730039270) + peer = {"OPTIONS": b"TS2=730500,730508;SINGLE=1;TIMER=5;"} + sys_cfg = {"SINGLE_MODE": False, "DEFAULT_UA_TIMER": 10} + slot = { + "RX_PEER": bytes_4(730039269), + "RX_TGID": bytes_3(730500), + "RX_STREAM_ID": bytes_4(0x11111111), + "RX_TIME": now, + "RX_TYPE": HBPF_SLT_VHEAD, + } + peer_slots = { + 2: {"stream_id": bytes_4(0x22222222), "tgid": 730502, "time": now, "ingress": False}, + } + assert peer_hotspot_voice_slot_busy( + hs, + 2, + bytes_4(0x11111111), + bytes_3(730500), + slot, + peer_slots, + None, + now + 0.1, + 5.0, + peer=peer, + sys_cfg=sys_cfg, + ) + + +def test_single1_ua_blocks_same_tg_foreign_tx() -> None: + """SINGLE=1 UA on 730502: block HP3ICC downlink on same TG while slot is busy.""" + now = 1_000_000.0 + listener = bytes_4(730039270) + tx_peer = bytes_4(730039269) + peer = {"OPTIONS": b"TS2=730500,730508;SINGLE=1;TIMER=5;"} + sys_cfg = {"SINGLE_MODE": False, "DEFAULT_UA_TIMER": 10} + register_peer_ua_session(peer, listener, 2, 730502, sys_cfg, now=now) + slot = { + "RX_PEER": tx_peer, + "RX_TGID": bytes_3(730502), + "RX_STREAM_ID": bytes_4(0x11111111), + "RX_TIME": now, + "RX_TYPE": HBPF_SLT_VHEAD, + } + assert peer_hotspot_voice_slot_busy( + listener, + 2, + bytes_4(0x22222222), + bytes_3(730502), + slot, + None, + None, + now + 0.1, + 5.0, + peer=peer, + sys_cfg=sys_cfg, + ) + + def test_global_scope_still_blocks_any_peer_on_busy_slot() -> None: now = 1_000_000.0 peer_b = bytes_4(714002301) diff --git a/tests/infrastructure/test_peer_disconnect_report.py b/tests/infrastructure/test_peer_disconnect_report.py index 2e1b128..73139b0 100644 --- a/tests/infrastructure/test_peer_disconnect_report.py +++ b/tests/infrastructure/test_peer_disconnect_report.py @@ -71,6 +71,21 @@ def test_disconnect_keeps_sys_cfg_sessions() -> None: assert export_peer_ua_sessions(sys_cfg, peer_id, now=1_000_100.0)["2"]["tgid"] == 7305 +def test_disconnect_clears_peer_voice_slots() -> None: + proto = _master_protocol() + peer_id = bytes_4(730039101) + pk = bytes_4(730039101) + proto._peer_voice_slots[pk] = { + 2: {"stream_id": bytes_4(0x12345678), "tgid": 7305, "time": 1_000_000.0}, + } + proto._peer_voice_hangtime[pk] = {2: (7305, 1_000_000.0)} + + proto._on_peer_disconnected(peer_id) + + assert pk not in proto._peer_voice_slots + assert pk not in proto._peer_voice_hangtime + + def test_export_peer_ua_sessions_omits_expired() -> None: peer_id = bytes_4(730039101) sys_cfg: dict = {"_PEER_UA_SESSIONS": {}} diff --git a/tests/infrastructure/test_self_rf_downlink_filter.py b/tests/infrastructure/test_self_rf_downlink_filter.py new file mode 100644 index 0000000..e874552 --- /dev/null +++ b/tests/infrastructure/test_self_rf_downlink_filter.py @@ -0,0 +1,79 @@ +# ADN DMR Peer Server - self rf_src downlink filter +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY + +from __future__ import annotations + +from adn_server.application.routing.downlink import touch_peer_voice_slot +from adn_server.application.routing.helpers import peer_matches_rf_source, synthetic_group_dmrd_route_packet +from adn_server.domain import bytes_3, bytes_4 +from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config +from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol + + +def _master_config() -> dict: + config = { + "GLOBAL": {"USE_ACL": False}, + "SYSTEMS": { + "TEST": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 8, + "GROUP_HANGTIME": 0, + } + }, + } + ensure_system_runtime_config(config) + return config + + +def test_group_downlink_blocked_when_self_rf_during_ingress_tx() -> None: + """Base rf_src echo must not downlink to a hotspot that is still transmitting.""" + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_id = bytes_4(730039264) + proto._peers = { + peer_id: { + "CONNECTION": "YES", + "OPTIONS": b"TS2=730502;", + "SOCKADDR": ("127.0.0.1", 62031), + }, + } + ctx = proto._downlink_ctx() + touch_peer_voice_slot( + ctx, + peer_id, + 2, + bytes_4(0x11111111), + bytes_3(730502), + ingress=True, + ) + base_rf = bytes_3(730039264 // 100) + pkt = synthetic_group_dmrd_route_packet(2, 730502) + pkt = pkt[:5] + base_rf + pkt[8:] + assert peer_matches_rf_source(peer_id, base_rf, proto._peers) + assert not proto._peer_should_receive_dmrd(peer_id, pkt) + + +def test_group_downlink_allowed_for_shared_base_rf_when_not_tx() -> None: + """Lab peers sharing rf_src base id may RX each other when not transmitting.""" + config = _master_config() + proto = HBPProtocol("TEST", config) + peer_id = bytes_4(730039264) + other = bytes_4(730039265) + proto._peers = { + peer_id: { + "CONNECTION": "YES", + "OPTIONS": b"TS2=730502;", + "SOCKADDR": ("127.0.0.1", 62031), + }, + other: { + "CONNECTION": "YES", + "OPTIONS": b"TS2=730502;", + "SOCKADDR": ("127.0.0.1", 62032), + }, + } + foreign_rf = bytes_3(730039265 // 100) + pkt = synthetic_group_dmrd_route_packet(2, 730502) + pkt = pkt[:5] + foreign_rf + pkt[8:] + assert proto._peer_should_receive_dmrd(peer_id, pkt)