fix: per-peer slot contention, silent activation, and multi-hotspot REPEAT gate

Enforce one active group QSO per peer/slot (SINGLE=0 and SINGLE=1) at the
downlink layer instead of treating the shared MASTER STATUS[slot] as a single
RF slot. Concurrent streams from different hotspots to different TGs are now
legitimate on ingress; contention is decided per-peer.

- Ingress: hbp_ingress_new_stream_collision honours per_peer=True so a second
  hotspot is not silenced by the first; hbp_master_ingress_repeat_allowed
  applies the same multi-hotspot rule and always passes VTERMs so listener
  sessions close and mid-call join works.
- Silent activation: TX onto a TG with an active QSO activates the TG
  dynamically, suppresses the uplink, and delivers the in-progress QSO
  downlink. Only a TG in GROUP_HANGTIME (no live stream) is rejected.
- Downlink: peer_voice_slots tracks one listen TG per peer/slot with
  GROUP_HANGTIME and bridge-hold semantics; stale sessions expire after
  STREAM_TO; duplex slots stay independent.
- Non-regression tests for slot contention, silent activation, stale session
  expiry, and duplex independence.
pull/29/head
Rodrigo Pérez 3 months ago
parent e09b58a05c
commit 8820d2ac99

@ -46,6 +46,7 @@ from .helpers import (
peer_receives_group_tgid,
peer_should_receive_group_voice,
peer_single_exclusive_tgid,
peer_single_mode,
peer_wants_downlink_single_listen_lock,
register_peer_ua_session,
remap_dmrd_to_peer_static_slot,
@ -320,6 +321,8 @@ def touch_peer_voice_slot(
"time": float(now),
}
prev = ctx.peer_voice_slots.get(pk, {}).get(int(voice_slot))
if not ingress and isinstance(prev, dict) and prev.get("ingress"):
return
if ingress or (isinstance(prev, dict) and prev.get("ingress")):
row["ingress"] = True
ctx.peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = row
@ -364,6 +367,7 @@ def end_peer_voice_slot(
*,
pkt_time: float | None = None,
apply_hangtime: bool = True,
from_ingress: bool = False,
) -> None:
"""Close session on VTERM; ingress VTERM may start GROUP_HANGTIME window."""
now = time.time() if pkt_time is None else float(pkt_time)
@ -377,14 +381,49 @@ def end_peer_voice_slot(
voice_slot = int(vs)
active = per_slot.pop(vs, None)
break
peer = ctx.peers.get(peer_id) or ctx.peers.get(pk)
if (
isinstance(peer, dict)
and ended_tg
and ctx.sys_cfg is not None
and not peer_single_mode(peer, ctx.sys_cfg)
and ctx.subscription_store is not None
and ctx.system_name
):
from adn_server.application.subscription.subscription_queries import (
system_has_active_leg_in_store,
)
voice_ended = isinstance(active, dict) and (
active.get("stream_id") or active.get("ingress")
)
if voice_ended and system_has_active_leg_in_store(
ctx.subscription_store,
ctx.system_name,
int(voice_slot),
int(ended_tg),
):
per_slot[int(voice_slot)] = {
"stream_id": b"",
"tgid": int(ended_tg),
"time": float(now),
"bridge_hold": True,
"bridge_hold_ingress": bool(from_ingress),
}
ctx.peer_voice_slots.setdefault(pk, {})[int(voice_slot)] = per_slot[int(voice_slot)]
return
if not apply_hangtime:
return
if isinstance(active, dict):
ended_tg = int(active.get("tgid", 0) or 0)
elif not ended_tg:
active_tg = int(active.get("tgid", 0) or 0)
# Only seed hangtime when the VTERM TG matches the session TG. A VTERM for a
# different TG (e.g. foreign stream that this peer never heard) must not seed
# hangtime from an unrelated active/bridge_hold session.
if active_tg and active_tg == ended_tg:
apply_hangtime_after_vterm(ctx, peer_id, voice_slot, active_tg, pkt_time=now)
return
if ended_tg:
apply_hangtime_after_vterm(ctx, peer_id, voice_slot, ended_tg, pkt_time=now)
# No active session for this peer: VTERM for a stream it never received — do not
# seed GROUP_HANGTIME (would block a fresh PTT on a different TG).
def track_peer_group_dmrd(
@ -429,6 +468,10 @@ def track_peer_group_dmrd(
clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot)
elif listen_lock and active_tg and active_tg == ended_tg:
clear_peer_ua_sessions(peer, ctx.sys_cfg, peer_id, slot=voice_slot)
pk = bytes_4(int_id(peer_id))
apply_hangtime = from_ingress
if not from_ingress and not peer_single_mode(peer, ctx.sys_cfg):
apply_hangtime = True
end_peer_voice_slot(
ctx,
peer_id,
@ -436,7 +479,8 @@ def track_peer_group_dmrd(
stream_id,
dst_id,
pkt_time=pkt_time,
apply_hangtime=from_ingress,
apply_hangtime=apply_hangtime,
from_ingress=from_ingress,
)
return
if (
@ -506,6 +550,8 @@ def peer_accepts_dmra(
peer_id: bytes,
slot: int,
tgid: int,
*,
pkt_time: float | None = None,
) -> bool:
"""P1: DMRA uses same accept + slot-busy rules as DMRD."""
peer = ctx.peers.get(peer_id)
@ -515,7 +561,9 @@ def peer_accepts_dmra(
if not peer_accepts_group_downlink(ctx, peer_id, peer, slot, tgid):
return False
remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id)
return not peer_slot_blocks_downlink(ctx, peer_id, peer, remapped)
return not peer_slot_blocks_downlink(
ctx, peer_id, peer, remapped, pkt_time=pkt_time,
)
def peer_would_show_group_voice_on_monitor(

@ -52,6 +52,7 @@ from .helpers import (
hbp_ingress_downlink_session_blocks_tx,
hbp_ingress_new_stream_collision,
master_per_peer_slot_contention,
tg_has_active_conversation,
)
from .peer_downlink_index import count_connected_peers
@ -143,6 +144,8 @@ class HbpForwardMixin:
_slot_st = getattr(src_proto, "STATUS", {}).get(slot, {})
_is_new_stream = stream_id != _slot_st.get("RX_STREAM_ID")
if _is_new_stream:
_slot_st.pop("_suppress_uplink", None)
_slot_st.pop("_silent_activation_tg", None)
sys_cfg = systems_cfg.get(system_name, {})
peers = getattr(src_proto, "_peers", None) or sys_cfg.get("PEERS", {})
connected = count_connected_peers(peers) if isinstance(peers, dict) else 0
@ -163,14 +166,29 @@ class HbpForwardMixin:
if group_voice_tg_ingress_collision(
protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time,
):
self._log_ingress_warning_once(
self._ingress_drop_key(
"tg_busy", system_name, peer_id, dst_id, stream_id,
),
"(%s) TG %s busy — dropping stream %s from peer %s",
system_name, int_id(dst_id), int_id(stream_id), int_id(peer_id),
)
return False
# TX onto a TG with an active (in-progress) QSO is not rejected.
# The TG is activated dynamically, the user's uplink audio is suppressed
# (not forwarded to the network), and the downlink of the active QSO is
# delivered to the user. Only a TG that is merely in GROUP_HANGTIME
# (no live stream) is rejected as busy (legacy parity).
if tg_has_active_conversation(
protocols, systems_cfg, dst_id, stream_id, rf_src, pkt_time,
):
logger.info(
"(%s) TG %s has active QSO — activating dynamic TG silently for peer %s (uplink suppressed)",
system_name, int_id(dst_id), int_id(peer_id),
)
_slot_st["_suppress_uplink"] = True
_slot_st["_silent_activation_tg"] = int_id(dst_id)
else:
self._log_ingress_warning_once(
self._ingress_drop_key(
"tg_busy", system_name, peer_id, dst_id, stream_id,
),
"(%s) TG %s busy — dropping stream %s from peer %s",
system_name, int_id(dst_id), int_id(stream_id), int_id(peer_id),
)
return False
from .downlink import normalize_ua_voice_slot
peer = peers.get(peer_id) if isinstance(peers, dict) else None

@ -52,6 +52,11 @@ from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO
PeerVoiceSlotRow = dict[str, Any]
PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow]
# A per-peer downlink voice session with no frames for this long is considered
# dead (VTERM lost or stream abandoned). Matches the legacy bridge_master idle
# loop threshold (bridge_master.py ~607: RX_TIME < now - 5).
_STALE_PEER_SESSION_TIMEOUT = 5.0
RF_MODE_SIMPLEX = "simplex"
RF_MODE_DUPLEX = "duplex"
# MMDVMHost DMO: downlink DMRD with TS1 bit set is dropped; only TS2 passes (DMRNetwork.cpp).
@ -347,9 +352,14 @@ def peer_hotspot_voice_slot_busy(
) -> bool:
"""True when this hotspot must not receive another group stream on ``voice_slot``.
Hard rule (SINGLE=0 and SINGLE=1): at most one group QSO per RF slot per hotspot.
A second TG or stream is dropped until VTERM clears ``peer_voice_slots`` (and
GROUP_HANGTIME where applicable). Same ``stream_id`` continues one call leg.
Hard rules (SINGLE=0 and SINGLE=1):
- **Transmitting** (``ingress``): drop every downlink byte until VTERM clears the slot.
- **Listening** on TG *T*: drop every byte for TG *U* ≠ *T* on this RF slot (no hangtime
exception for another OPTIONS/UA TG).
- **Bridge hold**: while an ACTIVE bridge leg keeps *T* on the slot, foreign legs on
*U* ≠ *T* stay dropped (OBP or HBP).
- Same ``stream_id`` on the same TG continues one call leg.
"""
pk = bytes_4(int_id(peer_id))
if peer is None and peers is not None:
@ -364,17 +374,18 @@ def peer_hotspot_voice_slot_busy(
active_time = float(active.get("time", 0) or 0)
age = pkt_time - active_time
if active.get("ingress"):
if (
age < STREAM_TO
and isinstance(peer, dict)
and sys_cfg is not None
and not peer_single_mode(peer, sys_cfg)
and active_tgid
and incoming_tgid != active_tgid
):
return False
return True
if stream_id and active_stream:
if active.get("bridge_hold") and active_tgid and incoming_tgid != active_tgid:
# Bridge hold (ingress from own TX or listener) blocks a foreign TG only for
# GROUP_HANGTIME, matching legacy bridge_master contention on TX_TIME/RX_TIME.
if age <= group_hangtime:
return True
peer_slots.pop(int(voice_slot), None)
active = None
if isinstance(active, dict) and active_time > 0 and age >= _STALE_PEER_SESSION_TIMEOUT:
peer_slots.pop(int(voice_slot), None)
active = None
if isinstance(active, dict) and stream_id and active_stream:
if active_stream == stream_id:
pass
elif active_tgid and active_tgid == incoming_tgid:
@ -386,8 +397,11 @@ def peer_hotspot_voice_slot_busy(
return True
else:
return True
else:
return True
elif isinstance(active, dict):
if active_tgid and active_tgid == incoming_tgid:
pass
else:
return True
if isinstance(peer, dict) and peer_single_blocks_foreign_same_tg_downlink(
peer, pk, voice_slot, incoming_tgid_b, peer_slots, sys_cfg, now=pkt_time,
):
@ -396,11 +410,6 @@ def peer_hotspot_voice_slot_busy(
peer, pk, incoming_tgid_b, stream_id, slot_st, sys_cfg, pkt_time=pkt_time,
):
return True
if stream_id and stream_id == slot_st.get("TX_STREAM_ID"):
tx_peer = slot_st.get("TX_PEER")
if tx_peer is not None and int_id(tx_peer) != 0:
if peers is None or not peer_key_in_peers(tx_peer, peers):
return False
if bytes_4(int_id(slot_st.get("RX_PEER", b""))) == pk:
if slot_has_active_voice(slot_st, pkt_time) and stream_id != slot_st.get("RX_STREAM_ID"):
return True
@ -530,17 +539,25 @@ def hbp_ingress_new_stream_collision(
) -> bool:
"""True when a new group-voice stream must drop on ingress (legacy routerHBP).
MASTER ``STATUS[slot]`` is shared: only one live stream per wire timeslot.
Same-subscriber rekey with a new stream id is allowed. ``per_peer`` applies
to downlink slot gates only, not ingress collision (legacy bridge_master).
When ``per_peer`` is False the MASTER ``STATUS[slot]`` is treated as a shared
RF slot: only one live stream per wire timeslot (legacy bridge_master).
When ``per_peer`` is True the MASTER fronts multiple hotspots, each on its
own frequency, so concurrent streams to different TGs on the same slot are
legitimate. Per-hotspot contention (one listen TG per peer/slot) is enforced
at downlink time by ``peer_hotspot_voice_slot_busy``; it must not be applied
here on ingress, otherwise a second hotspot is silenced by the first.
Same-subscriber rekey with a new stream id is always allowed.
"""
del per_peer, peer_id
from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO
if stream_id and stream_id == slot_st.get("RX_STREAM_ID"):
return False
if stream_id and stream_id == slot_st.get("TX_STREAM_ID"):
return False
if per_peer:
return False
del peer_id
for leg in ("RX", "TX"):
type_key = f"{leg}_TYPE"
time_key = f"{leg}_TIME"
@ -589,6 +606,65 @@ def _obp_stream_active(st: dict[str, Any], pkt_time: float) -> bool:
return True
def tg_has_active_conversation(
protocols: dict[str, Any],
systems_cfg: dict[str, Any],
tgid_b: bytes,
stream_id: bytes,
rf_src: bytes,
pkt_time: float,
) -> bool:
"""True when TG *tgid_b* has an active (in-progress) voice conversation.
Distinct from ``group_voice_tg_ingress_collision``: that helper also matches a TG
that is merely in ``GROUP_HANGTIME`` (idle but recent). This one only matches a TG
with a live stream (within ``STREAM_TO`` of the last frame), so the ingress gate
can activate the TG silently (no uplink, deliver downlink of the active QSO)
instead of rejecting the stream.
"""
tgid = int_id(tgid_b)
if tgid < 5 or tgid in (9, 4000, 5000):
return False
for sys_name, proto in protocols.items():
mode = systems_cfg.get(sys_name, {}).get("MODE")
status = getattr(proto, "STATUS", None)
if not isinstance(status, dict):
continue
if mode in ("MASTER", "PEER"):
for slot_key in (1, 2):
slot_st = status.get(slot_key)
if not isinstance(slot_st, dict):
continue
if not slot_has_active_voice(slot_st, pkt_time):
continue
active_tg = _hbp_slot_active_tgid(slot_st, pkt_time)
if active_tg is None or int_id(active_tg) != tgid:
continue
leg_stream = slot_st.get("RX_STREAM_ID") or slot_st.get("TX_STREAM_ID")
leg_rfs = slot_st.get("RX_RFS") or slot_st.get("TX_RFS") or b"\x00\x00\x00"
if stream_id and leg_stream == stream_id:
continue
if _same_rf_source(rf_src, leg_rfs):
continue
return True
elif mode == "OPENBRIDGE":
for key, st in status.items():
if isinstance(key, int) or not isinstance(st, dict):
continue
if "TGID" not in st or int_id(st.get("TGID", b"")) != tgid:
continue
if not _obp_stream_active(st, pkt_time):
continue
leg_stream = key if isinstance(key, (bytes, bytearray)) else b""
leg_rfs = st.get("RFS", b"\x00\x00\x00")
if stream_id and leg_stream == stream_id:
continue
if _same_rf_source(rf_src, leg_rfs):
continue
return True
return False
def group_voice_tg_ingress_collision(
protocols: dict[str, Any],
systems_cfg: dict[str, Any],
@ -663,8 +739,18 @@ def hbp_master_ingress_repeat_allowed(
*,
protocols: dict[str, Any] | None = None,
systems_cfg: dict[str, Any] | None = None,
is_vterm: bool = False,
system_cfg: dict[str, Any] | None = None,
) -> bool:
"""True when MASTER REPEAT may fan this ingress packet to other peers."""
"""True when MASTER REPEAT may fan this ingress packet to other peers.
In a multi-hotspot MASTER, concurrent streams from different peers to
different TGs on the same wire timeslot are legitimate: each hotspot is
on its own RF frequency. Contention is enforced per-peer at downlink
time (``peer_slot_blocks_downlink``). Applying a global shared-slot gate
here would drop voice packets of an active stream whenever a second
stream updates ``RX_STREAM_ID``, causing audible gaps on the first call.
"""
if stream_id and stream_id == slot_st.get("RX_STREAM_ID"):
owner = slot_st.get("RX_PEER")
return (
@ -672,8 +758,16 @@ def hbp_master_ingress_repeat_allowed(
and int_id(owner) != 0
and bytes_4(int_id(owner)) == bytes_4(int_id(peer_id))
)
# A VTERM closes an existing stream; it must reach peers that were hearing
# that stream even when another stream is now active on the shared slot
# (multi-hotspot MASTER). Blocking it leaves per-peer sessions open forever.
if is_vterm:
return True
per_peer = bool(system_cfg and master_per_peer_slot_contention(
systems_cfg or {}, "", system_cfg, connected_count=0,
))
if hbp_ingress_new_stream_collision(
slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=False,
slot_st, peer_id, rf_src, stream_id, pkt_time, per_peer=per_peer,
):
return False
if protocols and systems_cfg and group_voice_tg_ingress_collision(
@ -1321,19 +1415,22 @@ def peer_single_blocks_group_voice(
) -> bool:
"""True when SINGLE=1 peer must not receive downlink for ``tgid``.
With an active session on TG X, every other TG (static or dynamic) is blocked
until the TIMER expires or a new local TX replaces the session.
With an active SINGLE session on TG *X* in the peer's RF listen slot for
*tgid*, every other TG on **that same RF slot** is blocked until the TIMER
expires or a new local TX replaces the session.
Duplex hotspots have independent RF timeslots: a listen lock on TS1 must not
block a static TG on TS2 (and vice versa). Simplex hotspots always collapse
to ``SIMPLEX_VOICE_SLOT``, so both TGs share the same RF slot and blocking
still applies.
"""
del slot
if not sys_cfg:
return False
for voice_slot in (1, 2):
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now,
)
if locked is not None and int(tgid) != locked:
return True
return False
voice_slot = peer_downlink_voice_slot(peer, int(slot), int(tgid), sys_cfg, peer_id=peer_id)
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now,
)
return locked is not None and int(tgid) != locked
def peer_single_blocks_foreign_same_tg_downlink(

@ -276,6 +276,7 @@ class RoutingUseCases(
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
if not _obp_grp:
_is_new_rx_stream = True
_suppress_uplink_rx = False
if not source_is_obp:
protocols = self._get_protocols() if self._get_protocols else {}
src_proto = protocols.get(system_name) if protocols else None
@ -283,7 +284,8 @@ class RoutingUseCases(
slot_st = src_proto.STATUS.get(slot, {})
if isinstance(slot_st, dict):
_is_new_rx_stream = stream_id != slot_st.get("RX_STREAM_ID")
if _is_new_rx_stream:
_suppress_uplink_rx = bool(slot_st.get("_suppress_uplink"))
if _is_new_rx_stream and not _suppress_uplink_rx:
_rx_report_peer = peer_id
if not source_is_obp:
_rx_report_peer = resolve_voice_peer_id(
@ -305,6 +307,7 @@ class RoutingUseCases(
elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
if not _obp_grp:
duration = 0.0
_suppress_uplink_vterm = False
protocols = self._get_protocols() if self._get_protocols else {}
src_proto = protocols.get(system_name) if protocols else None
if src_proto and getattr(src_proto, "STATUS", None):
@ -315,25 +318,29 @@ class RoutingUseCases(
start = st.get(slot, {}).get("RX_START")
if start is not None:
duration = pkt_time - start
_rx_report_peer = peer_id
if not source_is_obp:
_rx_report_peer = resolve_voice_peer_id(
peer_id,
rf_src,
system_name,
systems_cfg,
)
self._send_routing_event(
"GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format(
system_name,
int_id(stream_id),
int_id(_rx_report_peer),
int_id(rf_src),
slot,
int_id(dst_id),
duration,
slot_st_vterm = st.get(slot, {})
if isinstance(slot_st_vterm, dict):
_suppress_uplink_vterm = bool(slot_st_vterm.get("_suppress_uplink"))
if not _suppress_uplink_vterm:
_rx_report_peer = peer_id
if not source_is_obp:
_rx_report_peer = resolve_voice_peer_id(
peer_id,
rf_src,
system_name,
systems_cfg,
)
self._send_routing_event(
"GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format(
system_name,
int_id(stream_id),
int_id(_rx_report_peer),
int_id(rf_src),
slot,
int_id(dst_id),
duration,
)
)
)
has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
)
@ -435,6 +442,15 @@ class RoutingUseCases(
resolve_voice_peer_id(peer_id, rf_src, system_name, systems_cfg)
)
forwarded = []
# When a TX arrived on a TG with an active QSO, the ingress gate set
# ``_suppress_uplink`` on the source slot status. The TG was activated dynamically
# (so downlink of the active QSO reaches this peer), but this stream's audio must
# not be forwarded upstream to avoid disrupting the in-progress conversation.
_suppress_uplink = False
if not source_is_obp and src_proto and getattr(src_proto, "STATUS", None):
_src_slot_st = src_proto.STATUS.get(slot, {})
if isinstance(_src_slot_st, dict) and _src_slot_st.get("_suppress_uplink"):
_suppress_uplink = True
_leg_iter: list[tuple[str, dict[str, Any]]] = [
(
forward_tables[0] if forward_tables else str(dst_int),
@ -449,6 +465,8 @@ class RoutingUseCases(
]
for _relay_table_key, entry in _leg_iter:
if _suppress_uplink:
continue
if not systems_cfg.get(entry["SYSTEM"], {}).get("ENABLED", True):
continue
_target_system = systems_cfg.get(entry["SYSTEM"], {})

@ -74,6 +74,55 @@ def active_system_slots_for_tg_in_store(
return tuple(sorted(slots))
def system_has_other_active_bridge_on_slot(
store: SubscriptionStore,
system: str,
slot: int,
incoming_tgid: int,
) -> bool:
"""True when another ACTIVE bridge leg occupies ``(system, slot)`` on a different TG.
Diagnostic / table export only — static OPTIONS bridges stay ACTIVE idle and must
not gate per-peer downlink (see ``peer_voice_slots`` / hangtime in ``downlink.py``).
"""
incoming = int(incoming_tgid)
for sub in store.snapshot():
if sub.system.value != system:
continue
if int(sub.channel.slot) != int(slot):
continue
if not sub.is_active():
continue
if int(sub.target_tgid) != incoming:
return True
return False
def bridge_timer_active_on_slot(
store: SubscriptionStore,
system: str,
slot: int,
tgid: int,
*,
now: float,
) -> bool:
"""True when an ACTIVE bridge leg for ``tgid`` on ``slot`` still has a live timer."""
tg = int(tgid)
for sub in store.snapshot():
if sub.system.value != system:
continue
if int(sub.channel.slot) != int(slot):
continue
if int(sub.target_tgid) != tg:
continue
if not sub.is_active():
continue
exp = sub.state.timer_expires_at
if exp is not None and float(exp) > float(now):
return True
return False
def store_legs_for_table(
store: SubscriptionStore,
table_key: str,

@ -1216,6 +1216,11 @@ class HBPProtocol(DatagramProtocol):
pkt_time,
protocols=_protocols,
systems_cfg=self._CONFIG.get("SYSTEMS", {}),
is_vterm=(
_frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VTERM
),
system_cfg=self._config,
)
)
if (
@ -1282,22 +1287,31 @@ class HBPProtocol(DatagramProtocol):
)
if _accepted:
if _slot in self.STATUS and not _unit_data:
_suppress = self.STATUS[_slot].get("_suppress_uplink")
if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"):
self.STATUS[_slot]["RX_START"] = pkt_time
if not _suppress:
self.STATUS[_slot]["RX_START"] = pkt_time
if _frame_type == HBPF_DATA_SYNC and _dtype_vseq == HBPF_SLT_VHEAD and len(dmrpkt) >= 33:
try:
decoded_slot = decode.voice_head_term(dmrpkt)
self.STATUS[_slot]["RX_LC"] = decoded_slot["LC"]
if not _suppress:
self.STATUS[_slot]["RX_LC"] = decoded_slot["LC"]
except Exception:
self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src
if not _suppress:
self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src
else:
self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src
self.STATUS[_slot]["RX_PEER"] = _peer_id
self.STATUS[_slot]["RX_SEQ"] = _seq
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_TIME"] = pkt_time
if not _suppress:
self.STATUS[_slot]["RX_LC"] = LC_OPT + _dst_id + _rf_src
if not _suppress:
self.STATUS[_slot]["RX_PEER"] = _peer_id
self.STATUS[_slot]["RX_SEQ"] = _seq
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_TIME"] = pkt_time
# Always track RX_STREAM_ID so subsequent frames of a suppressed
# stream are not re-evaluated as a new stream (which would clear
# _suppress_uplink and leak the uplink to the network).
self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id
_voice = self._CONFIG.get("VOICE", {})
if _accepted and self._on_handle_recording and _voice.get("RECORDING_ENABLED") and int_id(_dst_id) == _voice.get("RECORDING_TG") and _slot == _voice.get("RECORDING_TIMESLOT", 2):
@ -1329,9 +1343,20 @@ class HBPProtocol(DatagramProtocol):
and _call_type in ("group", "vcsbk")
and _frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VTERM
and _slot in self.STATUS
and self._on_talker_alias_stream_end
):
self._on_talker_alias_stream_end(self._system, _stream_id)
# Clean up silent-activation markers when the suppressed stream ends.
if (
_accepted
and _call_type in ("group", "vcsbk")
and _frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VTERM
and _slot in self.STATUS
):
self.STATUS[_slot].pop("_suppress_uplink", None)
self.STATUS[_slot].pop("_silent_activation_tg", None)
elif _command == RPTL:
_peer_id = _data[4:8]
@ -1723,15 +1748,28 @@ class HBPProtocol(DatagramProtocol):
):
self._on_in_band_signalling(self._system, _slot, _dst_id, pkt_time)
if _slot in self.STATUS and not _unit_data:
if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID"):
_suppress = self.STATUS[_slot].get("_suppress_uplink")
if _stream_id != self.STATUS[_slot].get("RX_STREAM_ID") and not _suppress:
self.STATUS[_slot]["RX_START"] = pkt_time
self.STATUS[_slot]["RX_PEER"] = _peer_id
self.STATUS[_slot]["RX_SEQ"] = _seq
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_TIME"] = pkt_time
if not _suppress:
self.STATUS[_slot]["RX_PEER"] = _peer_id
self.STATUS[_slot]["RX_SEQ"] = _seq
self.STATUS[_slot]["RX_RFS"] = _rf_src
self.STATUS[_slot]["RX_TYPE"] = _dtype_vseq
self.STATUS[_slot]["RX_TGID"] = _dst_id
self.STATUS[_slot]["RX_TIME"] = pkt_time
# Always track RX_STREAM_ID so subsequent frames of a suppressed
# stream are not re-evaluated as a new stream (leak of uplink).
self.STATUS[_slot]["RX_STREAM_ID"] = _stream_id
# Clean up silent-activation markers when the suppressed stream ends.
if (
_call_type in ("group", "vcsbk")
and _frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VTERM
and _slot in self.STATUS
):
self.STATUS[_slot].pop("_suppress_uplink", None)
self.STATUS[_slot].pop("_silent_activation_tg", None)
elif _command == DMRA:
if len(_data) >= DMRA_PACKET_LEN:

@ -16,6 +16,7 @@ from adn_server.application.routing.helpers import (
slot_has_active_voice,
slot_in_group_hangtime,
slot_status_hotspot_owner,
tg_has_active_conversation,
)
from adn_server.domain import HBPF_DATA_SYNC, bytes_3, bytes_4
from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO
@ -61,13 +62,18 @@ def test_ingress_same_subscriber_rekey_allowed() -> None:
)
def test_ingress_collides_when_other_hotspot_owns_busy_slot() -> None:
"""MASTER slot is global: second hotspot must not open a new stream on a busy TS."""
def test_ingress_allows_other_hotspot_when_per_peer() -> None:
"""Multi-hotspot MASTER: second hotspot may open a new stream on a busy TS.
Per-peer contention (one listen TG per hotspot/slot) is enforced at
downlink, not at ingress; a global slot collision here would silence a
second hotspot that operates on its own frequency.
"""
now = 1_000_000.0
slot = _active_rx_slot()
slot["RX_RFS"] = bytes_3(730039264)
slot["RX_PEER"] = bytes_4(730039264)
assert hbp_ingress_new_stream_collision(
assert not hbp_ingress_new_stream_collision(
slot,
bytes_4(730039265),
bytes_3(730039265),
@ -78,7 +84,11 @@ def test_ingress_collides_when_other_hotspot_owns_busy_slot() -> None:
def test_ingress_collides_when_obp_bridge_tx_leg_active() -> None:
"""OBP bridge TX stamp on STATUS[slot] must block a new HBP ingress stream."""
"""OBP bridge TX stamp on STATUS[slot] must block a new HBP ingress stream.
Only applies to the legacy shared-slot model (``per_peer=False``); with
``per_peer=True`` the decision is deferred to per-peer downlink gates.
"""
now = 1_000_000.0
slot = {
"RX_TYPE": HBPF_SLT_VTERM,
@ -93,7 +103,7 @@ def test_ingress_collides_when_obp_bridge_tx_leg_active() -> None:
bytes_3(730039267),
_STREAM_B,
now + 0.1,
per_peer=True,
per_peer=False,
)
@ -221,6 +231,94 @@ def test_group_voice_tg_collision_across_hbp_slots() -> None:
)
def test_tg_has_active_conversation_detects_live_hbp_qso() -> None:
"""Spec §3: a TG with an active (in-progress) HBP QSO is detected as active conversation."""
now = 1_000_000.0
master_status = {
2: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now),
}
master_status[2]["RX_RFS"] = bytes_3(730039264)
protocols = {"M1": type("P", (), {"STATUS": master_status})()}
systems = {"M1": {"MODE": "MASTER"}}
assert tg_has_active_conversation(
protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039265), now + 0.1,
)
def test_tg_has_active_conversation_detects_live_obp_qso() -> None:
"""Spec §3: a TG with an active (in-progress) OBP stream is detected as active conversation."""
now = 1_000_000.0
obp_stream = bytes_4(0x3E7A0F77)
obp_status = {
obp_stream: {
"TGID": bytes_3(730500),
"START": now - 5.0,
"LAST": now - 0.1,
"RFS": bytes_3(730039266),
},
}
protocols = {"OBP": type("P", (), {"STATUS": obp_status})()}
systems = {"OBP": {"MODE": "OPENBRIDGE"}}
assert tg_has_active_conversation(
protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039256), now,
)
def test_tg_has_active_conversation_false_for_hangtime_only() -> None:
"""Spec §3 vs §4: a TG that is merely in GROUP_HANGTIME (no live stream) is NOT active conversation."""
now = 1_000_000.0
# Slot idle (VTERM) but recent — within hangtime
master_status = {
2: {
"RX_TYPE": HBPF_SLT_VTERM,
"TX_TYPE": HBPF_SLT_VTERM,
"RX_TGID": bytes_3(730500),
"RX_TIME": now - 1.0,
"RX_STREAM_ID": _STREAM_A,
"TX_TIME": 0.0,
},
}
protocols = {"M1": type("P", (), {"STATUS": master_status})()}
systems = {"M1": {"MODE": "MASTER"}}
assert not tg_has_active_conversation(
protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039265), now,
)
def test_tg_has_active_conversation_false_for_stale_obp() -> None:
"""Spec §3: a truncated OBP stream (no recent packets) is NOT an active conversation."""
now = 1_000_000.0
obp_stream = bytes_4(0x3E7A0F77)
obp_status = {
obp_stream: {
"TGID": bytes_3(730500),
"START": now - 60.0,
"LAST": now - 45.0,
"RFS": bytes_3(730039266),
},
}
protocols = {"OBP": type("P", (), {"STATUS": obp_status})()}
systems = {"OBP": {"MODE": "OPENBRIDGE"}}
assert not tg_has_active_conversation(
protocols, systems, bytes_3(730500), _STREAM_B, bytes_3(730039256), now,
)
def test_tg_has_active_conversation_false_for_same_rf_source() -> None:
"""Spec §3: same RF source rekeying is not an 'active conversation' collision."""
now = 1_000_000.0
rf = bytes_3(730039264)
master_status = {
2: _active_rx_slot(tgid=bytes_3(730500), stream_id=_STREAM_A, t=now),
}
master_status[2]["RX_RFS"] = rf
protocols = {"M1": type("P", (), {"STATUS": master_status})()}
systems = {"M1": {"MODE": "MASTER"}}
assert not tg_has_active_conversation(
protocols, systems, bytes_3(730500), _STREAM_B, rf, now + 0.1,
)
def test_peer_hotspot_voice_slot_busy_blocks_downlink_during_local_ingress() -> None:
"""Rejected ingress still marks local TX — downlink must not arrive on that slot."""
now = 1_000_000.0
@ -294,7 +392,7 @@ def test_lab_witness_many_static_tgs_blocks_second_tg_on_slot() -> None:
tg_list = ",".join(str(730500 + i) for i in range(13))
witness = {"OPTIONS": f"TS2={tg_list};SINGLE=1;".encode()}
peer_slots = {
2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 1.0},
2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 0.1},
}
slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}
assert peer_hotspot_voice_slot_busy(
@ -644,8 +742,8 @@ def test_peer_hotspot_hangtime_blocks_other_tg() -> None:
)
def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None:
"""SINGLE=0: RX on another OPTIONS TG while local TX on the same RF slot."""
def test_single0_ingress_tx_blocks_other_tg_on_slot() -> None:
"""SINGLE=0: no downlink bytes on another TG while local TX on the same RF slot."""
now = 1_000_000.0
hs = bytes_4(730039253)
peer = {"OPTIONS": b"TS2=730507,730508;SINGLE=0;"}
@ -659,7 +757,7 @@ def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None:
"ingress": True,
},
}
assert not peer_hotspot_voice_slot_busy(
assert peer_hotspot_voice_slot_busy(
hs,
2,
bytes_4(0x22222222),
@ -671,6 +769,7 @@ def test_single0_ingress_tx_allows_other_static_tg_on_slot() -> None:
0.0,
peer=peer,
sys_cfg=sys_cfg,
peers={hs: peer, bytes_4(730039254): {"OPTIONS": b"TS2=730508;"}},
)
@ -1042,3 +1141,562 @@ def test_global_slot_blocks_foreign_vterm_during_active_rx() -> None:
slot, bytes_3(71442), panama_stream, now + 0.1, 0.0, is_vterm=True,
)
assert not hbp_slot_blocks_group_voice(slot, _TG_A, chile_stream, now + 0.1, 0.0)
def test_single0_listen_vterm_hangtime_blocks_obp_and_dmra() -> None:
"""SINGLE=0 listen-only: post-VTERM GROUP_HANGTIME blocks OBP voice and TA on another TG."""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_accepts_dmra,
peer_slot_blocks_downlink,
track_peer_group_dmrd,
)
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039269)
peer = {"OPTIONS": b"TS2=730501,730504;"}
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={
1: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM},
2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM},
},
connected_count=5,
)
now = 1_000_000.0
j39jq_stream = bytes_4(0x11111111)
obp_stream = bytes_4(0x22222222)
vhead = b"".join([
b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), j39jq_stream,
] + [b"\x00"] * 33)
vterm = b"".join([
b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), j39jq_stream,
] + [b"\x00"] * 33)
track_peer_group_dmrd(ctx, hs, vhead, peer, pkt_time=now)
track_peer_group_dmrd(ctx, hs, vterm, peer, pkt_time=now + 4.5)
assert ctx.peer_voice_hangtime[hs][2] == (730502, now + 4.5)
obp_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), obp_stream,
] + [b"\x00"] * 33)
overlap = now + 9.0
assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=overlap)
assert not peer_accepts_dmra(ctx, hs, 2, 730504, pkt_time=overlap)
assert not peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 11.0)
def test_live_listen_session_blocks_foreign_tg_downlink() -> None:
"""Mid-QSO downlink listen on TG A blocks foreign TG B on same slot (730502 vs 730504)."""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_accepts_dmra,
peer_slot_blocks_downlink,
touch_peer_voice_slot,
)
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039269)
peer = {"OPTIONS": b"TS2=730501,730504;"}
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}},
connected_count=5,
)
now = 1_000_000.0
j39jq_stream = bytes_4(0x11111111)
touch_peer_voice_slot(ctx, hs, 2, j39jq_stream, bytes_3(730502), pkt_time=now)
obp_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222),
] + [b"\x00"] * 33)
assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 0.5)
assert not peer_accepts_dmra(ctx, hs, 2, 730504, pkt_time=now + 0.5)
def test_idle_static_bridges_do_not_block_options_fanout() -> None:
"""Static ACTIVE bridge rows must not drop OPTIONS fan-out for another TG on the slot."""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
)
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
"SYSTEM": system,
"TS": ts,
"TGID": bytes_3(tgid),
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "OFF",
"ON": [bytes_3(tgid)],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": True}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039252)
peer = {"OPTIONS": b"TS2=730500,730504;"}
store = InMemorySubscriptionStore()
store.replace_all(
subscriptions_from_routing_table(
{
"730501": [_row(system="MASTER-A", ts=2, tgid=730501)],
"730502": [_row(system="MASTER-A", ts=2, tgid=730502)],
}
)
)
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}},
connected_count=6,
subscription_store=store,
)
ref_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039251), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (1 << 4) | HBPF_SLT_VHEAD]), bytes_4(0x33333333),
] + [b"\x00"] * 33)
assert not peer_slot_blocks_downlink(ctx, hs, peer, ref_vhead, pkt_time=1_000_002.0)
def test_foreign_obp_blocked_when_master_slot_carries_other_tg() -> None:
"""OBP (non-hotspot rf_src) blocked when this hotspot still holds the prior TG on slot."""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
)
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039269)
peer = {"OPTIONS": b"TS2=730501,730504;"}
now = 1_000_000.0
j39jq_stream = bytes_4(0x11111111)
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={
2: {
"RX_PEER": bytes_4(730039270),
"RX_TGID": bytes_3(730502),
"RX_STREAM_ID": j39jq_stream,
"RX_TIME": now,
"RX_TYPE": HBPF_SLT_VTERM,
"TX_TYPE": HBPF_SLT_VTERM,
},
},
peer_voice_slots={
hs: {
2: {
"stream_id": b"",
"tgid": 730502,
"time": now,
"bridge_hold": True,
},
},
},
connected_count=5,
)
obp_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222),
] + [b"\x00"] * 33)
hbp_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039270), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x33333333),
] + [b"\x00"] * 33)
assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 0.5)
assert peer_slot_blocks_downlink(ctx, hs, peer, hbp_vhead, pkt_time=now + 0.5)
def test_single0_ingress_vterm_bridge_hold_blocks_within_hangtime() -> None:
"""SINGLE=0 ingress TX: bridge_hold blocks a foreign TG only within GROUP_HANGTIME.
Matches legacy bridge_master contention on TX_TIME: after hangtime expires a fresh
PTT on a different TG is delivered.
"""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
track_peer_group_dmrd,
)
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
"SYSTEM": system,
"TS": ts,
"TGID": bytes_3(tgid),
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "OFF",
"ON": [bytes_3(tgid)],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039270)
peer = {"OPTIONS": b"TS2=730500,730508;"}
store = InMemorySubscriptionStore()
store.replace_all(
subscriptions_from_routing_table(
{"730502": [_row(system="MASTER-A", ts=2, tgid=730502)]},
)
)
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}},
connected_count=5,
subscription_store=store,
)
now = 1_000_000.0
stream = bytes_4(0x11111111)
vhead = b"".join([
b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), stream,
] + [b"\x00"] * 33)
vterm = b"".join([
b"DMRD", b"\x00", bytes_3(3520001), bytes_3(730502), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), stream,
] + [b"\x00"] * 33)
track_peer_group_dmrd(ctx, hs, vhead, peer, pkt_time=now, from_ingress=True, voice_slot=2)
track_peer_group_dmrd(ctx, hs, vterm, peer, pkt_time=now + 3, from_ingress=True, voice_slot=2)
assert ctx.peer_voice_slots[hs][2].get("bridge_hold") is True
obp_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(7140023), bytes_3(730504), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0x22222222),
] + [b"\x00"] * 33)
# Within GROUP_HANGTIME of the VTERM (age < 5s): blocked
assert peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 4.0)
# After GROUP_HANGTIME (age > 5s): fresh PTT on a different TG is delivered
assert not peer_slot_blocks_downlink(ctx, hs, peer, obp_vhead, pkt_time=now + 9.0)
def test_single0_listener_bridge_hold_expires_after_hangtime_allows_fresh_ptt() -> None:
"""SINGLE=0 listener (downlink RX, not ingress TX): bridge_hold must expire after
GROUP_HANGTIME so a fresh PTT on a different TG is delivered.
Reproduces the hs1/hs2/hs3 rule: hs1 receives 730500; after it ends and hangtime
clears, a NEW call on 730501 must reach hs1 from the start (no stale bridge_hold).
"""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
track_peer_group_dmrd,
)
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
"SYSTEM": system,
"TS": ts,
"TGID": bytes_3(tgid),
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "OFF",
"ON": [bytes_3(tgid)],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039251)
peer = {"OPTIONS": b"TS2=730500,730501;"}
store = InMemorySubscriptionStore()
store.replace_all(
subscriptions_from_routing_table(
{
"730500": [_row(system="MASTER-A", ts=2, tgid=730500)],
"730501": [_row(system="MASTER-A", ts=2, tgid=730501)],
}
)
)
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}},
connected_count=5,
subscription_store=store,
)
now = 1_000_000.0
ref_stream = bytes_4(0xAAAAAAAA)
ref_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream,
] + [b"\x00"] * 33)
ref_vterm = b"".join([
b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), ref_stream,
] + [b"\x00"] * 33)
track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, from_ingress=False, voice_slot=2)
track_peer_group_dmrd(ctx, hs, ref_vterm, peer, pkt_time=now + 8.0, from_ingress=False, voice_slot=2)
fresh_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0xBBBBBBBB),
] + [b"\x00"] * 33)
# During hangtime: fresh PTT on different TG must be blocked (listening on 730500).
assert peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 9.0)
# After hangtime clears: fresh PTT on different TG must be delivered.
assert not peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 20.0)
def test_single0_listener_fresh_ptt_after_blocked_second_stream() -> None:
"""Full hs1/hs2/hs3 flow: hs_a receives 730500; a second stream 730501 transits
the server (delivered to witness, blocked for hs_a); after 730500 ends and
hangtime clears, a fresh PTT on 730501 must reach hs_a from the start.
This reproduces the live scenario where a concurrent second stream leaves
residual slot state that blocked the fresh PTT in production.
"""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
track_peer_group_dmrd,
)
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
"SYSTEM": system,
"TS": ts,
"TGID": bytes_3(tgid),
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "OFF",
"ON": [bytes_3(tgid)],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039251)
peer = {"OPTIONS": b"TS2=730500,730501;"}
store = InMemorySubscriptionStore()
store.replace_all(
subscriptions_from_routing_table(
{
"730500": [_row(system="MASTER-A", ts=2, tgid=730500)],
"730501": [_row(system="MASTER-A", ts=2, tgid=730501)],
}
)
)
now = 1_000_000.0
ref_stream = bytes_4(0xAAAAAAAA)
second_stream = bytes_4(0xCCCCCCCC)
fresh_stream = bytes_4(0xDDDDDDDD)
status = {2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status=status,
connected_count=5,
subscription_store=store,
)
ref_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream,
] + [b"\x00"] * 33)
ref_vterm = b"".join([
b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), ref_stream,
] + [b"\x00"] * 33)
second_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), second_stream,
] + [b"\x00"] * 33)
second_vterm = b"".join([
b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VTERM]), second_stream,
] + [b"\x00"] * 33)
fresh_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), fresh_stream,
] + [b"\x00"] * 33)
# t=0: hs_a starts receiving 730500.
track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, from_ingress=False, voice_slot=2)
# t=3: second stream 730501 arrives — blocked for hs_a (busy on 730500).
assert peer_slot_blocks_downlink(ctx, hs, peer, second_vhead, pkt_time=now + 3.0)
# While the second stream is active on the wire, STATUS reflects the second
# peer as RX owner of slot 2 (the ingress updates RX_PEER/TGID/TIME/TYPE).
status[2] = {
"RX_PEER": bytes_4(730039252),
"RX_TGID": bytes_3(730501),
"RX_TIME": now + 5.0,
"RX_TYPE": HBPF_SLT_VHEAD,
"RX_STREAM_ID": second_stream,
"TX_TYPE": HBPF_SLT_VTERM,
"TX_TIME": 0.0,
}
# t=8: 730500 ends for hs_a.
track_peer_group_dmrd(ctx, hs, ref_vterm, peer, pkt_time=now + 8.0, from_ingress=False, voice_slot=2)
# t=12: second stream 730501 still in progress — no mid-join (hangtime clear path).
assert peer_slot_blocks_downlink(ctx, hs, peer, second_vhead, pkt_time=now + 12.0)
# t=17: second stream ends. STATUS global still carries the second peer as RX_OWNER.
track_peer_group_dmrd(ctx, hs, second_vterm, peer, pkt_time=now + 17.0, from_ingress=False, voice_slot=2)
status[2]["RX_TYPE"] = HBPF_SLT_VTERM
status[2]["RX_TIME"] = now + 17.0
# t=20: fresh PTT on 730501 — must reach hs_a (all prior hangtime cleared).
blocked = peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 20.0)
assert not blocked, f"fresh 730501 PTT blocked after all streams ended: slots={ctx.peer_voice_slots.get(hs)}, status={status.get(2)}"
def test_peer_stale_session_different_tg_expires_after_stream_to() -> None:
"""A stale per-peer voice session (no VTERM seen) on TG A must not block a
fresh stream on TG B once ``STREAM_TO`` has elapsed.
Reproduces the production case where a listener's downlink session on 730500
was left in ``peer_voice_slots`` with a non-empty stream_id after the stream
ended without VTERM, and blocked a subsequent 730501 PTT indefinitely.
"""
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
)
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
config = {"SYSTEMS": {"MASTER-A": sys_cfg}}
hs = bytes_4(730039251)
peer = {"OPTIONS": b"TS2=730500,730501;"}
now = 1_000_000.0
status = {2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}}
ctx = DownlinkContext(
config=config,
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status=status,
connected_count=5,
subscription_store=None,
)
ctx.peer_voice_slots[hs] = {
2: {"stream_id": bytes_4(0x952AFD04), "tgid": 730500, "time": now},
}
fresh_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039252), bytes_3(730501), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), bytes_4(0xDDDDDDDD),
] + [b"\x00"] * 33)
# While within STREAM_TO window: different TG is blocked (active QSO).
assert peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 0.1)
# After the stale-session timeout elapses: dead session expires, fresh PTT on 730501 passes.
assert not peer_slot_blocks_downlink(ctx, hs, peer, fresh_vhead, pkt_time=now + 6.0)
def test_single1_duplex_listen_lock_does_not_block_other_rf_slot() -> None:
"""SINGLE=1 duplex hotspot: a listen lock on TS1 must not block TS2.
Reproduces the intermittent ``single1-overlap-second-longer`` failure where
``hs_a`` (TS1=730500, TS2=730501, SINGLE=1) receives 730500 on TS1, then a
concurrent 730501 stream on TS2 was incorrectly blocked because
``peer_single_blocks_group_voice`` iterated both RF slots instead of
scoping the lock to the peer's listen slot for the incoming TG.
Duplex hotspots have independent RF timeslots; a SINGLE listen lock on one
timeslot must not deny voice on the other.
"""
from adn_server.application.routing.downlink import (
DownlinkContext,
track_peer_group_dmrd,
)
from adn_server.application.routing.helpers import (
peer_should_receive_group_voice,
)
from adn_server.application.subscription.routing_table_import import (
subscriptions_from_routing_table,
)
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row(*, system: str, ts: int, tgid: int) -> dict:
return {
"SYSTEM": system,
"TS": ts,
"TGID": bytes_3(tgid),
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "OFF",
"ON": [bytes_3(tgid)],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
sys_cfg = {"GROUP_HANGTIME": 5.0, "MODE": "MASTER", "MAX_PEERS": 8, "SINGLE_MODE": False}
hs = bytes_4(730039251)
peer = {"OPTIONS": b"TS1=730500;TS2=730501;SINGLE=1;TIMER=60;"}
store = InMemorySubscriptionStore()
store.replace_all(
subscriptions_from_routing_table(
{
"730500": [_row(system="MASTER-A", ts=1, tgid=730500)],
"730501": [_row(system="MASTER-A", ts=2, tgid=730501)],
}
)
)
ctx = DownlinkContext(
config={"SYSTEMS": {"MASTER-A": sys_cfg}},
system_name="MASTER-A",
sys_cfg=sys_cfg,
peers={hs: peer},
status={
1: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM},
2: {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM},
},
connected_count=5,
subscription_store=store,
)
now = 1_000_000.0
ref_stream = bytes_4(0xAAAAAAAA)
ref_vhead = b"".join([
b"DMRD", b"\x00", bytes_3(730039253), bytes_3(730500), b"\x00\x00\x00\x00",
bytes([0x80 | (HBPF_DATA_SYNC << 4) | HBPF_SLT_VHEAD]), ref_stream,
] + [b"\x00"] * 33)
# hs_a receives 730500 on TS1 — registers a SINGLE listen lock on slot 1.
track_peer_group_dmrd(ctx, hs, ref_vhead, peer, pkt_time=now, voice_slot=1)
assert peer["_UA_SESSION"][1]["tgid"] == 730500
# A 730501 stream arrives on TS2 — different RF slot, must not be blocked.
assert peer_should_receive_group_voice(
peer, 2, 730501, peer_id=hs, system="MASTER-A",
subscription_store=store, connected_count=5, sys_cfg=sys_cfg, now=now + 9.0,
)

@ -51,7 +51,13 @@ def test_hbp_source_timeout_drops_after_180_seconds() -> None:
assert scenario.protocols["MASTER-A"].STATUS[2].get("LOOPLOG") is True
def test_hbp_stream_collision_drops_conflicting_new_stream() -> None:
def test_hbp_stream_collision_silent_activation_on_busy_tg() -> None:
"""Spec §3 divergence: TX onto a TG with an active QSO is not rejected.
The stream passes the ingress gate (silent activation), but uplink audio is
suppressed — no forwarding to other systems. The downlink of the active QSO
continues to reach the peer independently.
"""
bridges = active_routing_table(91, (("MASTER-A", 2), ("MASTER-B", 2)))
scenario = DeterministicScenario(routing_table=bridges)
t0 = scenario.clock.time()
@ -74,8 +80,11 @@ def test_hbp_stream_collision_drops_conflicting_new_stream() -> None:
ingress_pkt_time=t0 + 0.1,
)
assert not ok
assert_inject_ok(ok)
# Uplink is suppressed: no packets forwarded to MASTER-B
assert len(scenario.capture.for_system("MASTER-B")) == 0
# Silent activation marker is set on the source slot
assert scenario.protocols["MASTER-A"].STATUS[2].get("_suppress_uplink") is True
@pytest.mark.behavior

Loading…
Cancel
Save

Powered by TurnKey Linux.