From 6c28754ce4fc470d05e98b45fe9e893b10a2a20a Mon Sep 17 00:00:00 2001 From: ce5rpy <169016246+ce5rpy@users.noreply.github.com> Date: Tue, 28 Jul 2026 15:50:26 -0400 Subject: [PATCH 1/2] fix: dual-slot delivery for mixed static+dynamic TG subscriptions (#66) * fix: deliver TG on both slots when static on one, dynamic on the other register_peer_ua_multi_tg silently dropped dynamic tracking on a slot whenever the same TG was already static on the other slot, and the REPEAT path never expanded to both slots at all. * fix: strip NUL-padded CALLSIGN instead of showing raw bytes Some peers NUL-pad instead of space-pad; reuse the existing normalize_fixed_width_ascii helper instead of a plain .strip(). --- src/adn_server/application/ident_use_cases.py | 4 +- .../application/routing/downlink.py | 76 ++++++-- src/adn_server/application/routing/helpers.py | 67 ++++++- .../application/routing/subscription_table.py | 6 +- .../twisted_adapters/udp_hbp.py | 58 +++++- .../application/test_peer_single_downlink.py | 26 +++ .../test_hbp_repeat_group_dual_slot.py | 180 ++++++++++++++++++ 7 files changed, 384 insertions(+), 33 deletions(-) create mode 100644 tests/infrastructure/test_hbp_repeat_group_dual_slot.py diff --git a/src/adn_server/application/ident_use_cases.py b/src/adn_server/application/ident_use_cases.py index 043e404..73d5866 100644 --- a/src/adn_server/application/ident_use_cases.py +++ b/src/adn_server/application/ident_use_cases.py @@ -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) diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index e8e90fc..0fd9e84 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -41,6 +41,8 @@ from .helpers import ( 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, @@ -165,23 +167,28 @@ def peer_hangtime_voice_slots( return {int(rf_slot), int(wire_slot), 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 +225,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 +248,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, @@ -263,8 +270,20 @@ def peer_slot_blocks_downlink( voice_slot=voice_slot, sys_cfg=ctx.sys_cfg, ): - return True - return False + return f"slot {voice_slot} busy or in GROUP_HANGTIME" + 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 +675,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)] diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 8fb6c8c..cda5025 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -1208,6 +1208,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 +1242,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 +1280,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, diff --git a/src/adn_server/application/routing/subscription_table.py b/src/adn_server/application/routing/subscription_table.py index 6eb4b64..2c76473 100644 --- a/src/adn_server/application/routing/subscription_table.py +++ b/src/adn_server/application/routing/subscription_table.py @@ -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]) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index d9e23ce..6a1e196 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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, ) @@ -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( @@ -752,16 +753,28 @@ class HBPProtocol(DatagramProtocol): stream_id = route_pkt[16:20] if len(route_pkt) >= 20 else b"" return pk, stream_id - def _log_downlink_drop_once(self, peer_id: bytes, route_pkt: bytes) -> None: + def _log_downlink_drop_once( + self, peer_id: bytes, route_pkt: bytes, peer: dict[str, Any] | None, + ) -> None: key = self._downlink_drop_key(peer_id, route_pkt) if key in self._downlink_drop_logged: return self._downlink_drop_logged.add(key) + # 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" 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 +783,34 @@ 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,8 +1394,23 @@ 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: diff --git a/tests/application/test_peer_single_downlink.py b/tests/application/test_peer_single_downlink.py index 654e591..7ae0534 100644 --- a/tests/application/test_peer_single_downlink.py +++ b/tests/application/test_peer_single_downlink.py @@ -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;"} diff --git a/tests/infrastructure/test_hbp_repeat_group_dual_slot.py b/tests/infrastructure/test_hbp_repeat_group_dual_slot.py new file mode 100644 index 0000000..b172a0d --- /dev/null +++ b/tests/infrastructure/test_hbp_repeat_group_dual_slot.py @@ -0,0 +1,180 @@ +# ADN DMR Peer Server - tests infrastructure hbp repeat group dual slot +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# 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_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] From c7883822c61b56317e57cc18b68eb724f1d85c25 Mon Sep 17 00:00:00 2001 From: ce5rpy <169016246+ce5rpy@users.noreply.github.com> Date: Wed, 29 Jul 2026 00:53:30 -0400 Subject: [PATCH 2/2] fix: deliver a hotspot's own group call back to its other slot when subscribed there too (#67) Adds self-echo (a hotspot with the same TG on both slots, static or dynamic, hears its own TX on the other slot) and fixes four occurrences of the same slot-resolution bug that blocked or truncated it: peer_downlink_voice_slot always preferred an unambiguous static/SINGLE=1-locked slot over the caller's own wire slot, corrupting the busy-check and the parrot/echo anti-loopback guard alike. Also makes the "Downlink dropped" log reason specific instead of generic, and throttles repeated identical drop logs. --- .../application/routing/downlink.py | 32 +++- src/adn_server/application/routing/helpers.py | 156 ++++++++++++++---- .../twisted_adapters/udp_hbp.py | 55 ++++-- tests/application/test_peer_rf_mode.py | 42 +++++ tests/application/test_slot_contention.py | 50 ++++++ .../test_hbp_repeat_group_dual_slot.py | 148 +++++++++++++++++ 6 files changed, 434 insertions(+), 49 deletions(-) diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index 0fd9e84..1bc6f5e 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -34,7 +34,7 @@ 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, @@ -159,12 +159,27 @@ 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_block_reason( @@ -256,7 +271,7 @@ def peer_slot_block_reason( 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, @@ -269,8 +284,9 @@ def peer_slot_block_reason( peer_hang_row=hang_row, voice_slot=voice_slot, sys_cfg=ctx.sys_cfg, - ): - return f"slot {voice_slot} busy or in GROUP_HANGTIME" + ) + if busy_reason is not None: + return busy_reason return None diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index cda5025..3256058 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -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 ) @@ -1793,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)) @@ -1806,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: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 6a1e196..ca181aa 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -257,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() @@ -672,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"): @@ -748,18 +754,11 @@ 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, peer: dict[str, Any] | None, ) -> None: - key = self._downlink_drop_key(peer_id, route_pkt) - if key in self._downlink_drop_logged: - return - self._downlink_drop_logged.add(key) # 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 @@ -769,6 +768,17 @@ class HBPProtocol(DatagramProtocol): 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[key] = now logger.info( "(%s) Downlink dropped for peer %s TG %s (%s)", self._system, @@ -797,7 +807,6 @@ class HBPProtocol(DatagramProtocol): ): 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): @@ -1415,6 +1424,26 @@ class HBPProtocol(DatagramProtocol): 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, diff --git a/tests/application/test_peer_rf_mode.py b/tests/application/test_peer_rf_mode.py index b424142..b9b8d67 100644 --- a/tests/application/test_peer_rf_mode.py +++ b/tests/application/test_peer_rf_mode.py @@ -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) diff --git a/tests/application/test_slot_contention.py b/tests/application/test_slot_contention.py index 6c2d758..6e2bf38 100644 --- a/tests/application/test_slot_contention.py +++ b/tests/application/test_slot_contention.py @@ -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 diff --git a/tests/infrastructure/test_hbp_repeat_group_dual_slot.py b/tests/infrastructure/test_hbp_repeat_group_dual_slot.py index b172a0d..d401040 100644 --- a/tests/infrastructure/test_hbp_repeat_group_dual_slot.py +++ b/tests/infrastructure/test_hbp_repeat_group_dual_slot.py @@ -152,6 +152,88 @@ def test_group_call_static_on_one_slot_dynamic_on_other_repeats_to_both() -> Non 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 @@ -178,3 +260,69 @@ def test_group_call_on_dynamic_tg_single_mode_stays_exclusive_to_one_slot() -> N 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}" + )