fix: block OBP during server announcements on inject-only masters

Re-apply slot contention when a MASTER holds TX_RFS=5000 for VHEAD/TTS
announcements so inject-only OBP deferral cannot interleave and chop audio.
pull/38/head
Rodrigo Pérez 3 months ago
parent ab28cc725e
commit cae5d57d8e

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

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

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

@ -0,0 +1,50 @@
# ADN DMR Peer Server - server broadcast slot contention
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
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

@ -0,0 +1,71 @@
# ADN DMR Peer Server - OBP vs server announcement slot contention
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.