Merge pull request #68 from ce5rpy/develop

fix: deliver a hotspot's own group call back to its other slot when subscribed there too
pull/69/head
ce5rpy 2 months ago committed by GitHub
commit 1bcbcb8de2
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -31,6 +31,7 @@ import time
from typing import Any, Callable
from ..domain import HBPF_SLT_VTERM, bytes_3, int_id
from ..domain.hbp_protocol import normalize_fixed_width_ascii
from .server_voice import server_voice_rf_src_bytes
logger = logging.getLogger(__name__)
@ -97,8 +98,7 @@ class IdentUseCases:
for _peerid in peers:
peer_cfg = peers.get(_peerid, {})
if isinstance(peer_cfg, dict) and peer_cfg.get("CALLSIGN"):
cs = peer_cfg["CALLSIGN"]
_callsign = cs.decode("utf-8", errors="replace") if isinstance(cs, bytes) else cs
_callsign = normalize_fixed_width_ascii(peer_cfg["CALLSIGN"])
break
if not _callsign:
logger.debug("(IDENT) %s System has no peers or no recorded callsign, skipping", system)

@ -34,13 +34,15 @@ from .helpers import (
_peer_transmit_hangtime_blocks,
_peer_ua_session_entry,
clear_peer_ua_sessions,
hbp_slot_blocks_group_voice_for_peer,
hbp_slot_blocks_group_voice_for_peer_reason,
is_special_tg,
is_ua_session_tgid,
master_per_peer_slot_contention,
parse_dmrd_burst_fields,
parse_dmrd_route_fields,
peer_downlink_voice_slot,
peer_dynamic_tg_active_on_both_slots,
peer_dynamic_tg_active_on_slot,
peer_is_simplex,
peer_options_static_tg_slot,
peer_receives_group_tgid,
@ -157,31 +159,51 @@ def peer_hangtime_voice_slots(
*,
peer_id: bytes | None = None,
) -> set[int]:
"""RF / OPTIONS slots that share transmit hangtime for this downlink."""
rf_slot = normalize_ua_voice_slot(peer, int(wire_slot))
"""RF / OPTIONS slots that share transmit hangtime for this downlink.
``wire_slot`` is already the caller's final, decided delivery slot (the
packet has already been remapped). If this peer is genuinely eligible for
``tgid`` on ``wire_slot`` (static or dynamic, per the same combined logic
``iter_downlink_voice_slots`` uses for the actual delivery decision),
trust it and check only that slot -- falling back to also deriving
``rf_slot``/``listen_slot`` would pull in an unrelated slot (e.g. a static
slot for the same TG on this peer) whenever ``peer_downlink_voice_slot``'s
single-slot answer differs from ``wire_slot``, wrongly reporting that
unrelated slot's own busy/ingress state instead of ``wire_slot``'s.
"""
wire_slot_i = int(wire_slot)
eligible = iter_downlink_voice_slots(peer, wire_slot_i, int(tgid), sys_cfg, peer_id=peer_id)
if wire_slot_i in eligible:
return {wire_slot_i}
rf_slot = normalize_ua_voice_slot(peer, wire_slot_i)
listen_slot = peer_downlink_voice_slot(
peer, int(wire_slot), int(tgid), sys_cfg, peer_id=peer_id,
peer, wire_slot_i, int(tgid), sys_cfg, peer_id=peer_id,
)
return {int(rf_slot), int(wire_slot), int(listen_slot)}
return {int(rf_slot), wire_slot_i, int(listen_slot)}
def peer_slot_blocks_downlink(
def peer_slot_block_reason(
ctx: DownlinkContext,
peer_id: bytes,
peer: dict[str, Any],
packet: bytes,
*,
pkt_time: float | None = None,
) -> bool:
"""P2/P3: True when slot busy or GROUP_HANGTIME blocks this downlink."""
) -> str | None:
"""P2/P3: reason this downlink is slot-busy/hangtime blocked, or None if not.
Same logic as ``peer_slot_blocks_downlink`` (which just checks ``is not
None``) -- kept as one function so the diagnostic reason can never drift
from the actual accept/reject decision.
"""
if packet[:4] != b"DMRD":
return False
return None
parsed = parse_dmrd_route_fields(packet)
if parsed is None:
return False
return None
wire_slot, tgid, call_type = parsed
if call_type not in ("group", "vcsbk"):
return False
return None
pk = bytes_4(int_id(peer_id))
burst = parse_dmrd_burst_fields(packet)
stream_id = burst[3] if burst is not None else b""
@ -218,22 +240,22 @@ def peer_slot_blocks_downlink(
and int_id(slot_st.get("RX_TGID", b"")) == active_tg
)
if not rx_listening:
return True
return f"VTERM: other stream TG {active_tg} still active on this peer"
active = per_slot.get(int(listen_slot))
if not isinstance(active, dict):
for row in per_slot.values():
if not isinstance(row, dict):
continue
if int(row.get("tgid", 0) or 0) == int_id(dst_id):
return False
return False
return None
return None
if active.get("stream_id") == stream_id:
return False
return None
if int(active.get("tgid", 0) or 0) == int_id(dst_id):
return False
return True
return None
return f"VTERM: listen slot busy with TG {int(active.get('tgid', 0) or 0)}"
if not ctx.per_peer_contention():
return False
return None
hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0)
now = time.time() if pkt_time is None else float(pkt_time)
seed_hangtime_for_stale_ingress_voice_slots(ctx, peer_id, pkt_time=now)
@ -241,7 +263,7 @@ def peer_slot_blocks_downlink(
if hang > 0:
for hang_row in ctx.peer_voice_hangtime.get(pk, {}).values():
if _peer_transmit_hangtime_blocks(hang_row, incoming_tgid_b, now, hang):
return True
return "GROUP_HANGTIME: recent local transmit still hanging on this peer"
peer_slots = ctx.peer_voice_slots.get(pk)
voice_slots = peer_hangtime_voice_slots(
peer, wire_slot, tgid, ctx.sys_cfg, peer_id=peer_id,
@ -249,7 +271,7 @@ def peer_slot_blocks_downlink(
for voice_slot in sorted(voice_slots):
hang_row = ctx.peer_voice_hangtime.get(pk, {}).get(voice_slot)
slot_st = ctx.status.get(voice_slot, {})
if hbp_slot_blocks_group_voice_for_peer(
busy_reason = hbp_slot_blocks_group_voice_for_peer_reason(
slot_st,
peer_id,
incoming_tgid_b,
@ -262,9 +284,22 @@ def peer_slot_blocks_downlink(
peer_hang_row=hang_row,
voice_slot=voice_slot,
sys_cfg=ctx.sys_cfg,
):
return True
return False
)
if busy_reason is not None:
return busy_reason
return None
def peer_slot_blocks_downlink(
ctx: DownlinkContext,
peer_id: bytes,
peer: dict[str, Any],
packet: bytes,
*,
pkt_time: float | None = None,
) -> bool:
"""P2/P3: True when slot busy or GROUP_HANGTIME blocks this downlink."""
return peer_slot_block_reason(ctx, peer_id, peer, packet, pkt_time=pkt_time) is not None
def remap_dmrd_for_peer(
@ -656,17 +691,38 @@ def iter_downlink_voice_slots(
peer: dict[str, Any],
wire_slot: int,
tgid: int,
sys_cfg: dict[str, Any] | None = None,
*,
peer_id: bytes | None = None,
) -> list[int]:
"""Voice slot(s) to deliver a group downlink to this peer.
Trusts ``peer_listen_slots`` -- it already collapses simplex peers/bridges
to one slot via ``peer_is_simplex``, and only returns more than one slot
for a peer confirmed duplex-capable with the TG on both TS1 and TS2, which
should genuinely receive on both (see peer_listen_slots docstring).
should genuinely receive on both (see peer_listen_slots docstring). Static
vs dynamic makes no difference here -- a TG independently activated on
both slots (SINGLE=0 keyed, or SINGLE=1 exclusive session per slot) gets
the same dual-slot treatment via peer_dynamic_tg_active_on_both_slots.
Mixed case: static on exactly one slot *and* dynamically active on the
other (e.g. TS1 static, TS2 keyed dynamically) must also deliver to both
-- checked via peer_dynamic_tg_active_on_slot on whichever slot the
static list didn't already claim.
"""
listen = peer_listen_slots(peer, tgid)
if len(listen) > 1:
return listen
if len(listen) == 1 and not peer_is_simplex(peer):
static_slot = listen[0]
other_slot = 2 if static_slot == 1 else 1
if peer_dynamic_tg_active_on_slot(peer, tgid, other_slot, sys_cfg, peer_id=peer_id):
return sorted({static_slot, other_slot})
return listen
if peer_dynamic_tg_active_on_both_slots(peer, tgid, sys_cfg, peer_id=peer_id):
return [1, 2]
if listen:
return listen
if peer_receives_group_tgid(peer, wire_slot, tgid):
return [peer_downlink_voice_slot(peer, wire_slot, tgid)]
return [peer_downlink_voice_slot(peer, wire_slot, tgid, sys_cfg, peer_id=peer_id)]
return [int(wire_slot)]

@ -336,7 +336,7 @@ def peer_single_same_tg_foreign_tx_blocks(
return True
def peer_hotspot_voice_slot_busy(
def peer_hotspot_voice_slot_busy_reason(
peer_id: bytes,
voice_slot: int,
stream_id: bytes,
@ -350,8 +350,12 @@ def peer_hotspot_voice_slot_busy(
peers: dict[Any, Any] | None = None,
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``.
) -> str | None:
"""Reason this hotspot must not receive another group stream on ``voice_slot``, or None.
Same logic as ``peer_hotspot_voice_slot_busy`` (which just checks ``is not
None``) -- kept as one function so the diagnostic reason can never drift
from the actual accept/reject decision.
Hard rules (SINGLE=0 and SINGLE=1):
@ -366,10 +370,10 @@ def peer_hotspot_voice_slot_busy(
if peer is None and peers is not None:
peer = peers.get(peer_id) or peers.get(pk)
if _peer_transmit_hangtime_blocks(hang_row, incoming_tgid_b, pkt_time, group_hangtime):
return True
return f"GROUP_HANGTIME: recent local transmit hanging on slot {voice_slot}"
incoming_tgid = int_id(incoming_tgid_b)
active = (peer_slots or {}).get(int(voice_slot))
if isinstance(active, dict):
incoming_tgid = int_id(incoming_tgid_b)
active_tgid = int(active.get("tgid", 0) or 0)
active_stream = active.get("stream_id")
active_time = float(active.get("time", 0) or 0)
@ -378,10 +382,10 @@ def peer_hotspot_voice_slot_busy(
peer_slots.pop(int(voice_slot), None)
active = None
if isinstance(active, dict) and active.get("ingress"):
return True
return f"peer is transmitting (ingress) on slot {voice_slot}"
if isinstance(active, dict) and active.get("bridge_hold") and active_tgid and incoming_tgid != active_tgid:
if age <= group_hangtime:
return True
return f"bridge hold: TG {active_tgid} active on slot {voice_slot}"
peer_slots.pop(int(voice_slot), None)
active = None
if isinstance(active, dict) and stream_id and active_stream:
@ -391,24 +395,28 @@ def peer_hotspot_voice_slot_busy(
if age >= STREAM_TO:
peer_slots.pop(int(voice_slot), None)
else:
return True
return f"TG {incoming_tgid} already active on slot {voice_slot} with a different stream"
elif active_tgid and incoming_tgid != active_tgid:
return True
return f"slot {voice_slot} busy with different TG {active_tgid}"
else:
return True
return f"slot {voice_slot} busy with an unidentified active stream"
elif isinstance(active, dict):
if active_tgid and active_tgid == incoming_tgid:
pass
else:
return True
return (
f"slot {voice_slot} already occupied by TG {active_tgid}"
if active_tgid
else f"slot {voice_slot} already occupied by another stream"
)
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
return f"SINGLE mode: local UA lock on TG {incoming_tgid} blocks foreign downlink on slot {voice_slot}"
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
return f"SINGLE mode: peer is transmitting a different call on TG {incoming_tgid}"
if bytes_4(int_id(slot_st.get("RX_PEER", b""))) == pk:
rx_active = (
slot_st.get("RX_TYPE") is not None
@ -416,12 +424,47 @@ def peer_hotspot_voice_slot_busy(
and (pkt_time - float(slot_st.get("RX_TIME", 0))) < STREAM_TO
)
if rx_active and stream_id != slot_st.get("RX_STREAM_ID"):
return True
return f"slot {voice_slot} STATUS RX owner with a different active stream"
if _peer_status_rx_hangtime_blocks(
peer_id, slot_st, incoming_tgid_b, pkt_time, group_hangtime,
):
return True
return False
return f"GROUP_HANGTIME: recent STATUS RX hangtime on slot {voice_slot}"
return None
def peer_hotspot_voice_slot_busy(
peer_id: bytes,
voice_slot: int,
stream_id: bytes,
incoming_tgid_b: bytes,
slot_st: dict[str, Any],
peer_slots: PeerVoiceSlotMap | None,
hang_row: tuple[int, float] | None,
pkt_time: float,
group_hangtime: float,
*,
peers: dict[Any, Any] | None = None,
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``."""
return (
peer_hotspot_voice_slot_busy_reason(
peer_id,
voice_slot,
stream_id,
incoming_tgid_b,
slot_st,
peer_slots,
hang_row,
pkt_time,
group_hangtime,
peers=peers,
peer=peer,
sys_cfg=sys_cfg,
)
is not None
)
def master_per_peer_slot_contention(
@ -602,7 +645,7 @@ def obp_clear_deferred_bridge_tx_leg(
obp_clear_flat_bridge_tx(slot_st)
def hbp_slot_blocks_group_voice_for_peer(
def hbp_slot_blocks_group_voice_for_peer_reason(
slot_st: dict[str, Any],
peer_id: bytes,
incoming_tgid_b: bytes,
@ -616,24 +659,21 @@ def hbp_slot_blocks_group_voice_for_peer(
peer_hang_row: tuple[int, float] | None = None,
voice_slot: int | None = None,
sys_cfg: dict[str, Any] | None = None,
) -> bool:
"""Slot contention scoped to one hotspot when ``per_peer`` (inject-only multi-HS).
``STATUS[slot]`` is shared at the MASTER, but each connected hotspot has an
independent RF timeslot. Another peer's active QSO must not block this peer.
) -> str | None:
"""Reason for ``hbp_slot_blocks_group_voice_for_peer``'s block, or None.
Inject-only OBP→HBP defers global contention and stamps bridge ``TX_*`` on the
shared slot row before ``send_peer``. Same-stream exemption must not treat that
bridge TX stamp as the hotspot's own leg while the peer is still on the air (RX).
Same logic as ``hbp_slot_blocks_group_voice_for_peer`` (which just checks
``is not None``) -- kept as one function so the diagnostic reason can
never drift from the actual accept/reject decision.
"""
if per_peer:
if voice_slot is None:
return False
return None
peer = None
if peers is not None:
pk = bytes_4(int_id(peer_id))
peer = peers.get(peer_id) or peers.get(pk)
return peer_hotspot_voice_slot_busy(
return peer_hotspot_voice_slot_busy_reason(
peer_id,
int(voice_slot),
stream_id,
@ -647,8 +687,53 @@ def hbp_slot_blocks_group_voice_for_peer(
peer=peer if isinstance(peer, dict) else None,
sys_cfg=sys_cfg,
)
return hbp_slot_blocks_group_voice(
if hbp_slot_blocks_group_voice(
slot_st, incoming_tgid_b, stream_id, pkt_time, group_hangtime,
):
return "global STATUS slot contention"
return None
def hbp_slot_blocks_group_voice_for_peer(
slot_st: dict[str, Any],
peer_id: bytes,
incoming_tgid_b: bytes,
stream_id: bytes,
pkt_time: float,
group_hangtime: float,
*,
per_peer: bool,
peers: dict[Any, Any] | None = None,
peer_slots: PeerVoiceSlotMap | None = None,
peer_hang_row: tuple[int, float] | None = None,
voice_slot: int | None = None,
sys_cfg: dict[str, Any] | None = None,
) -> bool:
"""Slot contention scoped to one hotspot when ``per_peer`` (inject-only multi-HS).
``STATUS[slot]`` is shared at the MASTER, but each connected hotspot has an
independent RF timeslot. Another peer's active QSO must not block this peer.
Inject-only OBP→HBP defers global contention and stamps bridge ``TX_*`` on the
shared slot row before ``send_peer``. Same-stream exemption must not treat that
bridge TX stamp as the hotspot's own leg while the peer is still on the air (RX).
"""
return (
hbp_slot_blocks_group_voice_for_peer_reason(
slot_st,
peer_id,
incoming_tgid_b,
stream_id,
pkt_time,
group_hangtime,
per_peer=per_peer,
peers=peers,
peer_slots=peer_slots,
peer_hang_row=peer_hang_row,
voice_slot=voice_slot,
sys_cfg=sys_cfg,
)
is not None
)
@ -1208,6 +1293,27 @@ def _peer_ua_multi_store(sys_cfg: dict[str, Any]) -> dict[bytes, dict[int, set[i
return store
def _peer_static_tg_blocks_slot(peer: dict[str, Any], slot: int, tgid: int) -> bool:
"""Does this peer's static OPTIONS already cover ``tgid`` for this exact
slot? Simplex peers have one real RF path regardless of nominal TS1/TS2,
so any static match blocks (matches peer_receives_group_tgid's
either-slot check); duplex peers are checked per-slot, since a static
match on one slot must not block genuinely independent dynamic activity
on the *other* slot (e.g. TG static on TS2, this same peer separately
keying up the same TG on TS1)."""
from adn_server.application.report.payloads import parse_peer_options_static
ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS"))
tg = str(tgid)
if peer_is_simplex(peer):
return tg in ts1 or tg in ts2
if int(slot) == 1:
return tg in ts1
if int(slot) == 2:
return tg in ts2
return False
def register_peer_ua_multi_tg(
peer: dict[str, Any],
peer_id: bytes,
@ -1221,7 +1327,7 @@ def register_peer_ua_multi_tg(
tgid_i = int(tgid)
if not is_ua_session_tgid(tgid_i):
return
if peer_receives_group_tgid(peer, slot, tgid_i):
if _peer_static_tg_blocks_slot(peer, slot, tgid_i):
return
pk = bytes_4(int_id(peer_id))
per_peer = _peer_ua_multi_store(sys_cfg).setdefault(pk, {})
@ -1259,6 +1365,50 @@ def peer_owns_multi_dynamic_ua(
return False
def peer_dynamic_tg_active_on_slot(
peer: dict[str, Any],
tgid: int,
slot: int,
sys_cfg: dict[str, Any] | None,
*,
peer_id: bytes | None = None,
) -> bool:
"""True when a *dynamic* (non-static OPTIONS) TG is active on this
specific slot for this peer right now. Covers SINGLE=1 (independent
exclusive session per slot) and SINGLE=0 (independently keyed multi-TG
set per slot)."""
if not sys_cfg or peer_id is None:
return False
tgid_i = int(tgid)
if peer_single_mode(peer, sys_cfg):
locked = peer_single_exclusive_tgid(peer, slot, sys_cfg, peer_id=peer_id)
return locked is not None and locked == tgid_i
store = sys_cfg.get("_PEER_UA_MULTI_TGS")
if not isinstance(store, dict):
return False
per_peer = store.get(bytes_4(int_id(peer_id)))
if not isinstance(per_peer, dict):
return False
slot_set = per_peer.get(slot)
return isinstance(slot_set, set) and tgid_i in slot_set
def peer_dynamic_tg_active_on_both_slots(
peer: dict[str, Any],
tgid: int,
sys_cfg: dict[str, Any] | None,
*,
peer_id: bytes | None = None,
) -> bool:
"""True when a *dynamic* (non-static OPTIONS) TG is active on both slot 1
and slot 2 for this peer right now -- static-or-dynamic makes no
difference to whether a duplex peer should get the call on both slots."""
return (
peer_dynamic_tg_active_on_slot(peer, tgid, 1, sys_cfg, peer_id=peer_id)
and peer_dynamic_tg_active_on_slot(peer, tgid, 2, sys_cfg, peer_id=peer_id)
)
def register_peer_ua_session(
peer: dict[str, Any],
peer_id: bytes,
@ -1728,6 +1878,17 @@ def peer_downlink_voice_slot(
static = peer_options_static_tg_slot(peer, tgid)
if static is not None:
return static
from adn_server.application.report.payloads import parse_peer_options_static
ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS"))
tg = str(tgid)
if tg in ts1 and tg in ts2:
# Static on both slots (peer_options_static_tg_slot returns None
# because that's ambiguous *as a single answer*, not because the TG
# is unresolved) -- an in-progress SINGLE=1 exclusive lock or SINGLE=0
# UA_MULTI entry from this peer's own current transmission must not
# override a config that already, unambiguously, permits wire_slot.
return int(wire_slot)
if sys_cfg is not None and peer_id is not None:
tgid_i = int(tgid)
pk = bytes_4(int_id(peer_id))
@ -1741,6 +1902,10 @@ def peer_downlink_voice_slot(
if isinstance(store, dict):
per_peer = store.get(pk)
if isinstance(per_peer, dict):
wire_slot_i = int(wire_slot)
wire_slot_set = per_peer.get(wire_slot_i)
if isinstance(wire_slot_set, set) and tgid_i in wire_slot_set:
return wire_slot_i
for voice_slot in (1, 2):
slot_set = per_peer.get(voice_slot)
if isinstance(slot_set, set) and tgid_i in slot_set:

@ -51,6 +51,7 @@ from typing import Any
from ...domain import bytes_3, bytes_4, int_id
from ...domain.config_coerce import coerce_bool, parse_options_single
from ...domain.dynamic_tg import DynamicTgEntry
from ...domain.hbp_protocol import normalize_fixed_width_ascii
from ..proxy.deployment import is_proxy_inject_only
logger = logging.getLogger(__name__)
@ -805,10 +806,7 @@ class SubscriptionTableMixin:
parts.append(f"peers_connected={len(connected)}")
if connected:
def _cs(c):
v = c.get("CALLSIGN") or b""
if isinstance(v, bytes):
return v.decode("utf8", errors="replace").strip() or "?"
return str(v).strip() or "?"
return normalize_fixed_width_ascii(c.get("CALLSIGN")) or "?"
parts.append(
"peers=[%s]"
% ", ".join("%s/%s" % (p.get("RADIO_ID", "?"), _cs(p)) for p in connected[:10])

@ -45,6 +45,7 @@ from ...application.routing.downlink import (
normalize_ua_voice_slot,
peer_accepts_dmra,
peer_accepts_group_dmrd_packet,
peer_slot_block_reason,
remap_dmrd_for_peer,
track_peer_group_dmrd,
)
@ -256,7 +257,7 @@ class HBPProtocol(DatagramProtocol):
self._connected_peer_count = 0
self._peer_voice_slots: dict[bytes, dict[int, dict[str, Any]]] = {}
self._peer_voice_hangtime: dict[bytes, dict[int, tuple[int, float]]] = {}
self._downlink_drop_logged: set[tuple[bytes, bytes]] = set()
self._downlink_drop_logged: dict[tuple[bytes, str], float] = {}
self._config_push_delayed = None
self._config_push_throttle = ConfigPushThrottle()
self._refresh_connected_peer_count()
@ -354,7 +355,7 @@ class HBPProtocol(DatagramProtocol):
if not isinstance(peer, dict) or call_type not in ("group", "vcsbk"):
self.send_peer(_peer, _packet)
continue
voice_slots = iter_downlink_voice_slots(peer, wire_slot, tgid)
voice_slots = iter_downlink_voice_slots(peer, wire_slot, tgid, self._config, peer_id=_peer)
if not voice_slots:
voice_slots = [
peer_downlink_voice_slot(
@ -671,9 +672,15 @@ class HBPProtocol(DatagramProtocol):
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,
)
# `packet` is already in its final, delivery-decided form (self-echo
# calls this via send_peer with the packet pre-remapped to its own
# other slot) -- trust its own wire slot rather than re-deriving one
# via peer_downlink_voice_slot, which prefers an unambiguous static
# slot regardless of which slot this delivery is actually for. That
# would wrongly resolve back to the peer's own currently-transmitting
# static slot and block a self-echo delivery to a different, free
# dynamic slot.
voice_slot = normalize_ua_voice_slot(peer, slot)
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"):
@ -747,21 +754,37 @@ class HBPProtocol(DatagramProtocol):
self._downlink_ctx(), peer_id, peer, packet, routed=routed,
)
def _downlink_drop_key(self, peer_id: bytes, route_pkt: bytes) -> tuple[bytes, bytes]:
pk = bytes_4(int_id(peer_id))
stream_id = route_pkt[16:20] if len(route_pkt) >= 20 else b""
return pk, stream_id
_DOWNLINK_DROP_LOG_COOLDOWN = 5.0
def _log_downlink_drop_once(self, peer_id: bytes, route_pkt: bytes) -> None:
key = self._downlink_drop_key(peer_id, route_pkt)
if key in self._downlink_drop_logged:
def _log_downlink_drop_once(
self, peer_id: bytes, route_pkt: bytes, peer: dict[str, Any] | None,
) -> None:
# By the time send_peer reaches this call, _peer_should_receive_dmrd
# already passed (it returns silently, without logging, on its own
# earlier check) -- so the only thing left that could have rejected
# it here is slot-busy/GROUP_HANGTIME. peer_slot_block_reason shares
# its logic with the actual accept/reject check (peer_slot_blocks_downlink),
# so this can't drift from the real decision.
reason = "unknown"
if isinstance(peer, dict):
reason = peer_slot_block_reason(self._downlink_ctx(), peer_id, peer, route_pkt) or "unknown"
# Keyed on (peer, reason) rather than stream_id: a still-ongoing call
# keeps hitting the same reason on every voice frame, so per-stream
# dedup alone doesn't stop the spam. Cooldown re-logs periodically
# instead of only once, so a long-running block doesn't go silent.
pk = bytes_4(int_id(peer_id))
key = (pk, reason)
now = time.time()
last = self._downlink_drop_logged.get(key)
if last is not None and (now - last) < self._DOWNLINK_DROP_LOG_COOLDOWN:
return
self._downlink_drop_logged.add(key)
self._downlink_drop_logged[key] = now
logger.info(
"(%s) Downlink dropped for peer %s TG %s (OPTIONS filter, slot busy, or GROUP_HANGTIME)",
"(%s) Downlink dropped for peer %s TG %s (%s)",
self._system,
int_id(peer_id),
int_id(route_pkt[8:11]),
reason,
)
def send_peer(self, _peer: bytes, _packet: bytes, *, _skip_dual_expand: bool = False) -> None:
@ -770,20 +793,33 @@ class HBPProtocol(DatagramProtocol):
return
peer = self._peers.get(_peer)
route_pkt = _packet
if peer is not None:
if peer is not None and not _skip_dual_expand:
# Caller already remapped to the intended slot via
# iter_downlink_voice_slots' explicit per-slot loop -- re-running
# the single-slot remap here would collapse a dual-slot delivery
# (e.g. a dynamically-locked TG active on both slots) back onto
# whichever slot that single-slot decision happens to prefer.
route_pkt = remap_dmrd_for_peer(
_packet, peer, self._config, peer_id=_peer,
)
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)
self._log_downlink_drop_once(_peer, route_pkt, peer)
return
self._downlink_drop_logged.discard(self._downlink_drop_key(_peer, route_pkt))
_packet = b"".join([route_pkt[:11], _peer, route_pkt[15:]])
ctx = self._downlink_ctx()
if isinstance(peer, dict):
track_peer_group_dmrd(ctx, _peer, _packet, peer, from_ingress=False)
# Pass the packet's own (already-decided) slot explicitly --
# track_peer_group_dmrd's own voice_slot=None fallback re-derives
# it via peer_downlink_voice_slot, which always prefers slot 1
# when a dynamic TG is locked on both slots, mistracking state
# for whichever delivery actually used slot 2.
_route_parsed = parse_dmrd_route_fields(_packet)
track_peer_group_dmrd(
ctx, _peer, _packet, peer, from_ingress=False,
voice_slot=_route_parsed[0] if _route_parsed else None,
)
self.transport.write(_packet, self._peers[_peer]["SOCKADDR"])
def _ta_buffer_enabled(self) -> bool:
@ -1367,12 +1403,47 @@ class HBPProtocol(DatagramProtocol):
_pvt_targets = None
if _pvt_targets is None:
for _peer in self._iter_downlink_peers(_repeat_pkt):
if _peer != _peer_id:
self.send_peer(_peer, _repeat_pkt)
if _peer == _peer_id:
continue
_repeat_peer_obj = self._peers.get(_peer)
if _call_type in ("group", "vcsbk") and isinstance(_repeat_peer_obj, dict):
_repeat_voice_slots = iter_downlink_voice_slots(
_repeat_peer_obj, _slot, int_id(_dst_id),
self._config, peer_id=_peer,
)
if len(_repeat_voice_slots) > 1:
for _repeat_vs in _repeat_voice_slots:
_repeat_remapped = remap_dmrd_for_peer(
_repeat_pkt, _repeat_peer_obj, self._config,
peer_id=_peer, voice_slot=_repeat_vs,
)
self.send_peer(_peer, _repeat_remapped, _skip_dual_expand=True)
continue
self.send_peer(_peer, _repeat_pkt)
else:
for _peer in _pvt_targets:
if _peer != _peer_id:
self.send_peer(_peer, _repeat_pkt)
# Same repeater, other slot: if this peer is also subscribed
# (static or dynamic) to this TG on the slot it did NOT just
# transmit on, deliver there too -- that slot is a genuinely
# independent RF path also tuned to this TG. Purely additive:
# does not change delivery to any other peer above.
if _call_type in ("group", "vcsbk"):
_src_peer_obj = self._peers.get(_peer_id)
if isinstance(_src_peer_obj, dict):
_src_voice_slots = iter_downlink_voice_slots(
_src_peer_obj, _slot, int_id(_dst_id),
self._config, peer_id=_peer_id,
)
for _src_vs in _src_voice_slots:
if _src_vs == _slot:
continue
_echo_remapped = remap_dmrd_for_peer(
_repeat_pkt, _src_peer_obj, self._config,
peer_id=_peer_id, voice_slot=_src_vs,
)
self.send_peer(_peer_id, _echo_remapped, _skip_dual_expand=True)
# TG 4000: reset after REPEAT so peers see the packet (legacy order)
if self._handle_tg4000_packet(
_peer_id, _slot, _int_dst_id, _call_type, _frame_type, _dtype_vseq,

@ -11,6 +11,7 @@ from adn_server.application.routing.helpers import (
RF_MODE_DUPLEX,
RF_MODE_SIMPLEX,
SIMPLEX_VOICE_SLOT,
bytes_4,
derive_peer_rf_mode,
peer_downlink_voice_slot,
peer_is_simplex,
@ -84,6 +85,47 @@ def test_duplex_keeps_cross_slot_static_remap() -> None:
assert not (remapped[15] & 0x80)
def test_downlink_voice_slot_prefers_wire_slot_when_dynamic_tg_active_on_both() -> None:
# Regression: a SINGLE=0 (UA_MULTI) TG dynamically active on both slots
# used to always resolve to slot 1 regardless of which slot was asked
# about, wrongly reporting the peer's own actively-transmitting slot as
# the "listen slot" for the opposite-direction self-echo delivery.
peer = _duplex_peer()
peer_rf_mode(peer)
peer_id = 714001
pk = bytes_4(peer_id)
sys_cfg = {"_PEER_UA_MULTI_TGS": {pk: {1: {71442}, 2: {71442}}}}
assert peer_downlink_voice_slot(peer, 1, 71442, sys_cfg, peer_id=peer_id) == 1
assert peer_downlink_voice_slot(peer, 2, 71442, sys_cfg, peer_id=peer_id) == 2
def test_downlink_voice_slot_static_both_slots_not_hijacked_by_single_exclusive_lock() -> None:
# Regression: a peer with a TG static on BOTH TS1 and TS2 (SINGLE=1) that
# is currently transmitting registers a SINGLE=1 exclusive session lock
# on its own TX slot. peer_downlink_voice_slot used to consult that lock
# (checking slot 1 before slot 2, unconditionally) whenever the static
# config was "ambiguous" (both slots), hijacking the answer to whichever
# slot the peer's own lock happened to be on -- wrongly reporting the
# peer's own actively-transmitting slot as busy for a self-echo delivery
# actually targeting the *other* slot.
peer = {
"SLOTS": b"3",
"RX_FREQ": b"431612500",
"TX_FREQ": b"438612500",
"OPTIONS": b"TS1=730500;TS2=730500;SINGLE=1;",
}
peer_rf_mode(peer)
peer_id = 730039258
pk = bytes_4(peer_id)
sys_cfg = {
"_PEER_UA_SESSIONS": {
pk: {1: {"tgid": 730500, "expires": 0, "source": "local"}},
},
}
assert peer_downlink_voice_slot(peer, 1, 730500, sys_cfg, peer_id=peer_id) == 1
assert peer_downlink_voice_slot(peer, 2, 730500, sys_cfg, peer_id=peer_id) == 2
def test_topology_peer_row_includes_rf_mode() -> None:
peer = _simplex_peer()
peer_rf_mode(peer)

@ -320,6 +320,32 @@ def test_peer_without_static_tgs_receives_nothing_until_dynamic() -> None:
)
def test_multi_zero_static_on_one_slot_still_tracks_dynamic_on_other() -> None:
"""Real bug: TG static on TS2, peer separately keys the SAME TG on TS1 --
the static match on TS2 must not block persisting TS1 as dynamic too.
register_peer_ua_multi_tg's guard used to be slot-blind (peer_receives_group_tgid
ignores its slot argument and checks either static list), so this dynamic
activation was silently never recorded."""
peer = {"OPTIONS": b"TS2=71442;SINGLE=0;"}
sys_cfg = _sys_cfg()
peer_id = _peer_id()
register_peer_ua_multi_tg(peer, peer_id, 1, 71442, sys_cfg)
pk = bytes_4(int_id(peer_id))
assert 71442 in sys_cfg["_PEER_UA_MULTI_TGS"][pk][1]
def test_multi_zero_simplex_static_on_either_slot_still_blocks_dynamic() -> None:
"""Simplex peer: one real RF path regardless of nominal TS1/TS2, so a
static match on either slot still blocks dynamic tracking (unlike
duplex, which is tracked independently per slot)."""
peer = {"OPTIONS": b"TS2=71442;SINGLE=0;", "RF_MODE": "simplex"}
sys_cfg = _sys_cfg()
peer_id = _peer_id()
register_peer_ua_multi_tg(peer, peer_id, 1, 71442, sys_cfg)
pk = bytes_4(int_id(peer_id))
assert pk not in sys_cfg.get("_PEER_UA_MULTI_TGS", {})
def test_single_zero_dynamic_heard_when_both_peers_keyed() -> None:
"""SINGLE=0: HS1 and HS2 both keyed 7304 → each hears the other's TX on 7304."""
peer_a = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}

@ -11,7 +11,9 @@ from adn_server.application.routing.helpers import (
hbp_master_ingress_repeat_allowed,
hbp_slot_blocks_group_voice,
hbp_slot_blocks_group_voice_for_peer,
hbp_slot_blocks_group_voice_for_peer_reason,
peer_hotspot_voice_slot_busy,
peer_hotspot_voice_slot_busy_reason,
register_peer_ua_session,
slot_has_active_voice,
slot_in_group_hangtime,
@ -333,6 +335,54 @@ def test_peer_hotspot_voice_slot_busy_blocks_downlink_during_local_ingress() ->
)
def test_peer_hotspot_voice_slot_busy_reason_distinguishes_ingress_and_bridge_hold() -> None:
"""Regression: the drop reason must name the actual cause, not a generic string."""
now = 1_000_000.0
hs = bytes_4(730039265)
slot = _active_rx_slot()
slot["RX_PEER"] = bytes_4(730039264)
ingress_slots = {
2: {"stream_id": _STREAM_B, "tgid": 730502, "time": now, "ingress": True},
}
reason = peer_hotspot_voice_slot_busy_reason(
hs, 2, _STREAM_A, _TG_A, slot, ingress_slots, None, now + 0.05, 5.0,
)
assert reason is not None and "transmitting (ingress) on slot 2" in reason
assert peer_hotspot_voice_slot_busy(
hs, 2, _STREAM_A, _TG_A, slot, ingress_slots, None, now + 0.05, 5.0,
) == (reason is not None)
bridge_slots = {
2: {"stream_id": _STREAM_B, "tgid": 71442, "time": now, "bridge_hold": True},
}
reason = peer_hotspot_voice_slot_busy_reason(
hs, 2, _STREAM_A, _TG_A, slot, bridge_slots, None, now + 0.05, 5.0,
)
assert reason is not None and "bridge hold: TG 71442 active on slot 2" in reason
assert peer_hotspot_voice_slot_busy(
hs, 2, _STREAM_A, _TG_A, slot, bridge_slots, None, now + 0.05, 5.0,
) == (reason is not None)
def test_hbp_slot_blocks_group_voice_for_peer_reason_matches_bool() -> None:
now = 1_000_000.0
peer_a = bytes_4(352000133)
peer_b = bytes_4(714002301)
slot = _active_rx_slot()
slot["RX_PEER"] = peer_a
assert hbp_slot_blocks_group_voice_for_peer_reason(
slot, peer_b, _TG_B, _STREAM_B, now + 0.1, 0.0, per_peer=True, voice_slot=2,
) is None
reason = hbp_slot_blocks_group_voice_for_peer_reason(
slot, peer_a, _TG_B, _STREAM_B, now + 0.1, 0.0, per_peer=True, voice_slot=2,
)
assert reason is not None
assert hbp_slot_blocks_group_voice_for_peer(
slot, peer_a, _TG_B, _STREAM_B, now + 0.1, 0.0, per_peer=True, voice_slot=2,
) is True
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

@ -0,0 +1,328 @@
# ADN DMR Peer Server - tests infrastructure hbp repeat group dual slot
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""REPEAT path through real HBPProtocol: peer subscribed to a TG on both
static OPTIONS slots must get the group call on both, regardless of the
source peer's own wire slot.
Regression: the intra-MASTER REPEAT loop (_master_datagram_received) sends
the raw incoming packet unchanged via send_peer(), never going through
iter_downlink_voice_slots()/peer_listen_slots() at all -- unlike send_peers()
(exercised by tests/infrastructure/test_peer_downlink_fanout.py), which
already trusted iter_downlink_voice_slots(). A peer-to-peer call on the same
MASTER (the common case) goes through REPEAT, not send_peers(), so the
TS1+TS2 dual-slot fix landed there without covering the path real traffic
actually uses."""
from __future__ import annotations
from tests.harness.deterministic import DeterministicScenario, PacketSpec, parse_dmr_fields
from tests.support.hbp_repeat_stack import build_hbp_repeat_stack
from adn_server.domain import bytes_4
from adn_server.infrastructure.hbp_constants import DMRD
_TG = 730
_PEER_TX = bytes_4(730039110)
_PEER_RX = bytes_4(730039101)
_ADDR_TX = ("10.0.0.1", 62001)
_ADDR_RX = ("10.0.0.2", 62002)
def _group_spec(slot: int) -> PacketSpec:
return PacketSpec(
peer_id=int.from_bytes(_PEER_TX, "big"),
rf_src=7300391,
dst_id=_TG,
slot=slot,
call_type="group",
stream_id=0xA1B2C3D4,
payload=b"\x00" * 33,
)
def test_group_call_on_both_static_slots_repeats_to_both() -> None:
stack = build_hbp_repeat_stack(talker_alias=False)
# SINGLE_MODE defaults to True in the generic test scenario config, which
# pulls in the (unrelated) dynamic-session machinery -- pin it explicitly
# to False here so this test only exercises the static-OPTIONS path.
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
# No RX_FREQ/TX_FREQ/SLOTS given -> derive_peer_rf_mode defaults to duplex
# (matches a real MMDVM_HS_Dual_Hat hotspot, which reported duplex here).
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX, _ADDR_RX, options=f"TS1={_TG};TS2={_TG};")
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
downlink = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_RX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in downlink)
assert slots == [1, 2], (
"peer with TG on both static slots must be repeated on both, "
f"regardless of the source's own wire slot -- got {slots}"
)
def test_group_call_on_one_static_slot_still_repeats_once() -> None:
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX, _ADDR_RX, options=f"TS1={_TG};")
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
downlink = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_RX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in downlink)
assert slots == [1]
def test_group_call_on_dynamic_tg_keyed_both_slots_repeats_to_both() -> None:
"""SINGLE=0 hotspot that has independently keyed the same dynamic
(non-static) TG on both slots -- static vs dynamic makes no difference to
whether a duplex peer should get the call on both."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX, _ADDR_RX, options="")
rx_pk = bytes_4(730039101)
stack.config["SYSTEMS"][stack.system_name].setdefault("_PEER_UA_MULTI_TGS", {})[rx_pk] = {
1: {_TG}, 2: {_TG},
}
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
downlink = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_RX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in downlink)
assert slots == [1, 2]
def test_group_call_static_on_one_slot_dynamic_on_other_repeats_to_both() -> None:
"""Real-world case: TG static on TS1, independently keyed dynamically
(SINGLE=0) on TS2 (e.g. via peer_dynamic_tgs DB restore) -- must repeat
to both, not just the statically-configured slot."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX, _ADDR_RX, options=f"TS1={_TG};")
rx_pk = bytes_4(730039101)
stack.config["SYSTEMS"][stack.system_name].setdefault("_PEER_UA_MULTI_TGS", {})[rx_pk] = {
2: {_TG},
}
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
downlink = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_RX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in downlink)
assert slots == [1, 2]
def test_group_call_echoes_to_transmitting_peers_own_other_slot() -> None:
"""New behavior: a repeater transmitting a TG on one slot, that is also
subscribed (static or dynamic) to that same TG on its OTHER slot, hears
its own call echoed there too -- that other slot is a genuinely
independent RF path also tuned to the TG."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
# TG static on TS1 only; transmits on slot 2 -- slot 1 is the "other" slot.
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS1={_TG};")
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in echo)
assert slots == [1], f"peer should hear itself echoed on its own other slot -- got {slots}"
def test_group_call_echoes_to_own_dynamically_subscribed_other_slot() -> None:
"""Same self-echo, but the other slot's subscription is dynamic (SINGLE=0
keyed), not static -- static vs dynamic must be treated identically."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options="")
pk = bytes_4(int.from_bytes(_PEER_TX, "big"))
stack.config["SYSTEMS"][stack.system_name].setdefault("_PEER_UA_MULTI_TGS", {})[pk] = {
1: {_TG},
}
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in echo)
assert slots == [1], f"peer should hear itself echoed on its dynamically-subscribed other slot -- got {slots}"
def test_group_call_echoes_symmetrically_on_opposite_slot() -> None:
"""Same behavior, transmitting on the opposite slot: TG static on TS2,
transmits on slot 1 -- echoes back to itself on slot 2."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
base = _group_spec(slot=1)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in echo)
assert slots == [2], f"peer should hear itself echoed on its own other slot -- got {slots}"
def test_group_call_no_self_echo_when_not_subscribed_on_other_slot() -> None:
"""No TG configured at all -- no self-echo, matching existing (unchanged)
single-peer behavior."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options="")
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
assert echo == []
def test_group_call_on_dynamic_tg_single_mode_stays_exclusive_to_one_slot() -> None:
"""SINGLE=1 ("one exclusive dynamic TG per hotspot, either RF slot; new
local TX replaces all others" -- register_peer_ua_session) cannot have
the same TG genuinely active on both slots at once: registering it on
slot 2 clears slot 1's session. This is intentional exclusivity, not a
duplex/simultaneous-dual-slot case like static OPTIONS or SINGLE=0 -- it
must keep collapsing to one slot."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX, _ADDR_RX, options="SINGLE=1;")
rx_pk = bytes_4(730039101)
stack.config["SYSTEMS"][stack.system_name].setdefault("_PEER_UA_SESSIONS", {})[rx_pk] = {
2: {"tgid": _TG, "expires": 0, "source": "local"},
}
base = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
downlink = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_RX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in downlink)
assert slots == [2]
def test_group_call_self_echo_not_blocked_by_own_ingress_when_dynamic_tg_on_both_slots() -> None:
"""Regression: when a TG is dynamically (SINGLE=0) active on both slots,
peer_downlink_voice_slot used to always resolve to slot 1 regardless of
which slot was asked about. That made the busy-check for a slot-1-to-2
self-echo also examine slot 1 -- the peer's own live ingress slot -- so
the echo was wrongly reported busy even though slot 2 was free."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options="")
pk = bytes_4(int.from_bytes(_PEER_TX, "big"))
stack.config["SYSTEMS"][stack.system_name].setdefault("_PEER_UA_MULTI_TGS", {})[pk] = {
1: {_TG}, 2: {_TG},
}
base = _group_spec(slot=1)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in echo)
assert slots == [2], (
"self-echo to slot 2 must not be blocked by the source's own ingress "
f"on slot 1 -- got {slots}"
)
def test_group_call_self_echo_to_dynamic_slot_not_blocked_by_own_static_slot_ingress() -> None:
"""Regression: TG static on TS1 only, dynamically activated on TS2 by an
earlier call. peer_hangtime_voice_slots used to always union in the
static slot (TS1) as a busy-check candidate, via peer_downlink_voice_slot,
even when the caller only cares about TS2 -- so a later TX on TS1 (the
static slot, now busy with the source's own ingress) wrongly blocked its
own self-echo delivery to TS2, even though TS2 was free."""
stack = build_hbp_repeat_stack(talker_alias=False)
stack.config["SYSTEMS"][stack.system_name]["SINGLE_MODE"] = False
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS1={_TG};SINGLE=0;")
# First call: TX on TS2 (not in static OPTIONS) -- dynamically activates
# TS2, and (as a side effect) self-echoes back to the static TS1 slot.
base_a = _group_spec(slot=2)
stack.inject_spec(DeterministicScenario.voice_head_spec(base_a), _ADDR_TX)
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base_a, seq=1, dtype_vseq=1), _ADDR_TX,
)
stack.inject_spec(DeterministicScenario.voice_term_spec(base_a, seq=2), _ADDR_TX)
# Second call: TX on TS1 (the static slot) -- must self-echo to TS2,
# which is now dynamically active from the first call.
stack.transport.clear()
base_b = _group_spec(slot=1)
stack.inject_spec(DeterministicScenario.voice_head_spec(base_b), _ADDR_TX)
stack.inject_spec(
DeterministicScenario.voice_burst_spec(base_b, seq=1, dtype_vseq=1), _ADDR_TX,
)
echo = [pkt for pkt, addr in stack.transport.sent if addr == _ADDR_TX and pkt[:4] == DMRD]
slots = sorted(parse_dmr_fields(pkt)["slot"] for pkt in echo)
assert slots == [2, 2], (
"self-echo to the dynamically-active TS2 must not be blocked by the "
f"source's own busy static TS1 -- got {slots}"
)
Loading…
Cancel
Save

Powered by TurnKey Linux.