fix: per-hotspot slot contention and ingress drop log dedup

Enforce one QSO per RF slot for normal hotspots, clear peer_voice_slots
on disconnect, allow matching VTERM through slot gates, exempt lab
witnesses with many static TGs, and log ingress TG-busy drops once per stream.
pull/29/head
Rodrigo Pérez 3 months ago
parent b914dda2f8
commit de385660c1

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

@ -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 <FROM> SUB: %s PEER: %s <TO> 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"]

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

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

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

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

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

@ -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": {}}

@ -0,0 +1,79 @@
# ADN DMR Peer Server - self rf_src downlink filter
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
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)
Loading…
Cancel
Save

Powered by TurnKey Linux.