diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 4a7a555..ab143c9 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -461,6 +461,27 @@ def inject_only_defer_obp_hbp_slot_contention( return source_is_hbp and connected_count > 1 +_SERVER_VOICE_RF_SRC = 5000 + + +def master_slot_holds_server_broadcast(slot_st: dict[str, Any], pkt_time: float) -> bool: + """True when server scheduled/TTS voice holds the flat MASTER slot TX row. + + Announcements stamp TX_TYPE=VHEAD, TX_RFS=5000, and refresh TX_TIME each frame. + Re-apply global slot contention for OBP→MASTER when held so inject-only defer + does not interleave mesh voice (legacy ``bridge_master`` TX_TGID/TX_TIME rules). + """ + if int_id(slot_st.get("TX_RFS", b"\x00\x00\x00")) != _SERVER_VOICE_RF_SRC: + return False + tx_type = slot_st.get("TX_TYPE") + if tx_type is None or tx_type == HBPF_SLT_VTERM: + return False + tx_time = float(slot_st.get("TX_TIME", 0) or 0) + if tx_time <= 0: + return True + return (pkt_time - tx_time) < STREAM_TO + + _OBP_FLAT_TX_KEYS = ( "TX_START", "TX_TGID", diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 7bc1113..33df780 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -46,6 +46,7 @@ from .routing.helpers import ( inject_only_defer_obp_hbp_slot_contention, is_private_subscriber_dst, is_unit_data_ingress, + master_slot_holds_server_broadcast, obp_clear_deferred_bridge_tx_leg, obp_deferred_bridge_tx_leg, obp_flat_bridge_tx_idle, @@ -678,8 +679,12 @@ class RoutingUseCases( ) if not _obp_deferred_bridge: _ts_st["TX_TYPE"] = HBPF_SLT_VTERM - if ( + _apply_master_slot_contention = ( not _defer_slot_contention + or master_slot_holds_server_broadcast(_ts_st, pkt_time) + ) + if ( + _apply_master_slot_contention and hbp_slot_blocks_group_voice( _ts_st, entry_tgid_b, diff --git a/src/adn_server/application/voice_use_cases.py b/src/adn_server/application/voice_use_cases.py index 3fd0c94..f6495d7 100644 --- a/src/adn_server/application/voice_use_cases.py +++ b/src/adn_server/application/voice_use_cases.py @@ -148,11 +148,15 @@ class VoiceUseCases: def _mark_slots_busy(self, targets: list[dict[str, Any]]) -> None: """Mark target slots busy (TX_TYPE=VHEAD) to prevent TS conflict.""" + server_rfs = bytes_3(5000) + now = time.time() for t in targets: try: slot = t.get("slot") if slot is not None: slot["TX_TYPE"] = HBPF_SLT_VHEAD + slot["TX_TIME"] = now + slot["TX_RFS"] = server_rfs except (KeyError, TypeError): pass @@ -305,6 +309,7 @@ class VoiceUseCases: "LAST": now, } slot["TX_TGID"] = dst_id + slot["TX_RFS"] = source_id else: sys_obj.STATUS[stream_id]["LAST"] = now slot["TX_TIME"] = now @@ -607,6 +612,7 @@ class VoiceUseCases: "LAST": now, } slot["TX_TGID"] = dst_id + slot["TX_RFS"] = source_id else: sys_obj.STATUS[stream_id]["LAST"] = now slot["TX_TIME"] = now diff --git a/tests/application/test_master_server_broadcast_contention.py b/tests/application/test_master_server_broadcast_contention.py new file mode 100644 index 0000000..cf1af84 --- /dev/null +++ b/tests/application/test_master_server_broadcast_contention.py @@ -0,0 +1,50 @@ +# ADN DMR Peer Server - server broadcast slot contention +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY + +from __future__ import annotations + +from adn_server.application.routing.helpers import ( + hbp_slot_blocks_group_voice, + master_slot_holds_server_broadcast, +) +from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4 + + +def test_master_slot_holds_server_broadcast_detects_announcement_row() -> None: + slot = { + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": 100.0, + "TX_RFS": bytes_3(5000), + "TX_TGID": bytes_3(91), + } + assert master_slot_holds_server_broadcast(slot, 100.1) is True + assert master_slot_holds_server_broadcast(slot, 100.0 + STREAM_TO + 1) is False + + +def test_master_slot_holds_server_broadcast_ignores_obp_tx() -> None: + slot = { + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": 100.0, + "TX_RFS": bytes_3(3340062), + "TX_TGID": bytes_3(7144), + } + assert master_slot_holds_server_broadcast(slot, 100.1) is False + + +def test_obp_blocked_when_server_broadcast_holds_slot() -> None: + slot = { + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": 200.0, + "TX_RFS": bytes_3(5000), + "TX_TGID": bytes_3(91), + } + blocked = hbp_slot_blocks_group_voice( + slot, + bytes_3(7144), + bytes_4(0x11111111), + 200.05, + group_hangtime=3.0, + ) + assert blocked is True diff --git a/tests/routing/test_obp_announcement_contention.py b/tests/routing/test_obp_announcement_contention.py new file mode 100644 index 0000000..9390982 --- /dev/null +++ b/tests/routing/test_obp_announcement_contention.py @@ -0,0 +1,71 @@ +# ADN DMR Peer Server - OBP vs server announcement slot contention +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY + +from __future__ import annotations + +from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time +from tests.harness.scenarios import obp_bridge_scenario + +from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4 + + +def _seed_server_broadcast_slot(scenario: DeterministicScenario, *, ann_tg: int = 91) -> None: + proto = scenario.protocols["MASTER-A"] + t = scenario.clock.time() + proto.STATUS[2] = { + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VHEAD, + "TX_TIME": t, + "TX_RFS": bytes_3(5000), + "TX_TGID": bytes_3(ann_tg), + "TX_STREAM_ID": bytes_4(0xA0A0A0A0), + } + + +def test_obp_blocked_while_server_broadcast_holds_master_ts2() -> None: + """Legacy parity: OBP must not route onto TS2 while announcement holds TX row.""" + scenario = obp_bridge_scenario("OBP-CL", tg=7144) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + scenario.config["SYSTEMS"]["MASTER-A"]["PEERS"] = { + bytes_4(730002): {"CONNECTION": "YES", "OPTIONS": b"TS2=7144;"}, + } + _seed_server_broadcast_slot(scenario, ann_tg=91) + base = PacketSpec( + peer_id=73010, + rf_src=3340062, + dst_id=7144, + slot=1, + stream_id=0x11111111, + ) + with patch_routing_wall_time(scenario.clock): + scenario.inject_obp("OBP-CL", DeterministicScenario.voice_head_spec(base)) + assert scenario.capture.for_system("MASTER-A") == [] + + +def test_obp_still_forwards_when_no_server_broadcast_hold() -> None: + """Inject-only defer remains when the slot is not held by server voice.""" + scenario = obp_bridge_scenario("OBP-CL", tg=7144) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + scenario.config["SYSTEMS"]["MASTER-A"]["PEERS"] = { + bytes_4(730002): {"CONNECTION": "YES", "OPTIONS": b"TS2=7144;"}, + } + proto = scenario.protocols["MASTER-A"] + t = scenario.clock.time() + proto.STATUS[2] = { + "RX_TYPE": HBPF_SLT_VTERM, + "TX_TYPE": HBPF_SLT_VTERM, + "TX_TIME": t - 60.0, + "TX_RFS": bytes_3(3340062), + "TX_TGID": bytes_3(7144), + } + base = PacketSpec( + peer_id=73010, + rf_src=3340062, + dst_id=7144, + slot=1, + stream_id=0x22222222, + ) + with patch_routing_wall_time(scenario.clock): + scenario.inject_obp("OBP-CL", DeterministicScenario.voice_head_spec(base)) + assert len(scenario.capture.for_system("MASTER-A")) > 0 diff --git a/tests/voice/test_announcement_anticollision.py b/tests/voice/test_announcement_anticollision.py index f49e6a8..6c865b4 100644 --- a/tests/voice/test_announcement_anticollision.py +++ b/tests/voice/test_announcement_anticollision.py @@ -68,3 +68,17 @@ def test_broadcast_aborts_when_qso_starts_mid_transmission() -> None: assert uc._announcement_running[0] is False assert master.sent == [] + + +def test_mark_slots_busy_stamps_server_voice_tx_row() -> None: + scenario, master = voice_master_scenario() + uc = make_voice_uc(scenario, master) + slot = master.STATUS[2] + slot["TX_TYPE"] = HBPF_SLT_VTERM + targets = [{"sys_obj": master, "name": "MASTER-A", "slot": slot, "ts": 2}] + + uc._mark_slots_busy(targets) + + assert slot["TX_TYPE"] == HBPF_SLT_VHEAD + assert int.from_bytes(slot["TX_RFS"], "big") == 5000 + assert slot["TX_TIME"] > 0