You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
ADN-DMR-Peer-Server/tests/routing/test_obp_concurrent_streams...

720 lines
26 KiB

This file contains ambiguous Unicode characters!

This file contains ambiguous Unicode characters that may be confused with others in your current locale. If your use case is intentional and legitimate, you can safely ignore this warning. Use the Escape button to highlight these characters.

# ADN DMR Peer Server - concurrent OBP streams downlink reproduction
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
"""Reproduce production bug: two concurrent OBP streams on one MASTER TS2 with a
dual-TG hotspot (SINGLE=0) intermittently drops the active listen stream.
Diagnosis-only: no production code changes. Uses DeterministicScenario for OBP
ingress + real HBPProtocol send_peer downlink gate.
"""
from __future__ import annotations
import copy
import logging
from contextlib import contextmanager
from typing import Any
import pytest
from tests.harness.deterministic import (
DeterministicScenario,
FakeClock,
PacketSpec,
add_openbridge_system,
patch_routing_wall_time,
)
from tests.support.hbp_repeat_stack import RecordingTransport
from adn_server.application.routing.downlink import (
DownlinkContext,
peer_slot_blocks_downlink,
)
from adn_server.application.routing.helpers import (
_peer_status_rx_hangtime_blocks,
_peer_transmit_hangtime_blocks,
hbp_slot_blocks_group_voice_for_peer,
peer_hotspot_voice_slot_busy,
peer_single_blocks_foreign_same_tg_downlink,
peer_single_same_tg_foreign_tx_blocks,
slot_has_active_voice,
)
from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id
from adn_server.domain.hbp_protocol import HBPF_VOICE, STREAM_TO
from adn_server.infrastructure.acl_router import InMemoryAclRouter
from adn_server.infrastructure.config_normalizer import ensure_system_runtime_config
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
_TG_A = 3109050
_TG_B = 52090
_STREAM_A = 0xAABBCCDD
_STREAM_B = 0x11223344
_HS_PEER = 730039101
_OBP_PEER = 73010
_RF_A = 3340001
_RF_B = 3340002
_INTERVAL_S = 0.06
_VOICE_BURSTS = 8
def _bridge_row(*, system: str, ts: int, tgid: int) -> dict[str, Any]:
tg_b = bytes_3(tgid)
return {
"SYSTEM": system,
"TS": ts,
"TGID": tg_b,
"ACTIVE": True,
"TIMEOUT": 3600.0,
"TO_TYPE": "ON",
"ON": [tg_b],
"OFF": [],
"RESET": [],
"TIMER": 0.0,
}
def _dual_tg_routing_table() -> dict[str, list[dict[str, Any]]]:
return {
str(_TG_A): [
_bridge_row(system="OBP-CL", ts=1, tgid=_TG_A),
_bridge_row(system="MASTER-A", ts=2, tgid=_TG_A),
],
str(_TG_B): [
_bridge_row(system="OBP-CL", ts=1, tgid=_TG_B),
_bridge_row(system="MASTER-A", ts=2, tgid=_TG_B),
],
}
@contextmanager
def _patch_harness_time(clock: FakeClock):
"""Patch wall clock on routing ingress and downlink gate paths."""
import adn_server.application.routing.downlink as downlink_mod
import adn_server.application.routing_use_cases as routing_mod
orig_routing = routing_mod.time.time
orig_downlink = downlink_mod.time.time
routing_mod.time.time = clock.time
downlink_mod.time.time = clock.time
try:
yield
finally:
routing_mod.time.time = orig_routing
downlink_mod.time.time = orig_downlink
def _build_stack() -> tuple[DeterministicScenario, HBPProtocol, RecordingTransport]:
routing_table = _dual_tg_routing_table()
scenario = DeterministicScenario(
routing_table=routing_table,
enable_reporting=True,
)
config = scenario.config
add_openbridge_system(config, "OBP-CL")
master = config["SYSTEMS"]["MASTER-A"]
master.update(
{
"MAX_PEERS": 8,
"GROUP_HANGTIME": 5.0,
"SINGLE_MODE": False,
"TS2_STATIC": f"{_TG_A},{_TG_B}",
}
)
ensure_system_runtime_config(config)
transport = RecordingTransport()
hbp = HBPProtocol("MASTER-A", config, router=InMemoryAclRouter())
hbp.transport = transport # type: ignore[assignment]
scenario.protocols["MASTER-A"] = hbp
def _send_to_system(target: str, packet: bytes, **kwargs: Any) -> None:
scenario.capture.recorder(target)(
packet,
hops=kwargs.get("_hops", kwargs.get("hops", b"")),
ber=kwargs.get("_ber", kwargs.get("ber", b"\x00")),
rssi=kwargs.get("_rssi", kwargs.get("rssi", b"\x00")),
source_server=kwargs.get(
"_source_server", kwargs.get("source_server", b"\x00\x00\x00\x00"),
),
source_rptr=kwargs.get(
"_source_rptr", kwargs.get("source_rptr", b"\x00\x00\x00\x00"),
),
)
proto = scenario.protocols.get(target)
if proto is not None and hasattr(proto, "send_system"):
proto.send_system(
packet,
_hops=kwargs.get("_hops", b""),
_ber=kwargs.get("_ber", b"\x00"),
_rssi=kwargs.get("_rssi", b"\x00"),
_source_server=kwargs.get("_source_server", b"\x00\x00\x00\x00"),
_source_rptr=kwargs.get("_source_rptr", b"\x00\x00\x00\x00"),
)
scenario.routing._send_to_system = _send_to_system
hs = bytes_4(_HS_PEER)
hbp._peers[hs] = {
"CONNECTION": "YES",
"CONNECTED": scenario.clock.time(),
"LAST_PING": scenario.clock.time(),
"SOCKADDR": ("127.0.0.1", 62031),
"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode(),
}
master.setdefault("PEERS", {})[hs] = hbp._peers[hs]
hbp._refresh_connected_peer_count()
hbp._mark_downlink_index_dirty()
return scenario, hbp, transport
def _base_spec(*, tgid: int, stream_id: int, rf_src: int) -> PacketSpec:
return PacketSpec(
peer_id=_OBP_PEER,
rf_src=rf_src,
dst_id=tgid,
slot=1,
stream_id=stream_id,
)
def _dmrd_packet(
*,
tgid: int,
stream_id: int,
rf_src: int,
frame_type: int,
dtype_vseq: int,
seq: int = 0,
) -> bytes:
base = _base_spec(tgid=tgid, stream_id=stream_id, rf_src=rf_src)
spec = PacketSpec(
peer_id=base.peer_id,
rf_src=base.rf_src,
dst_id=base.dst_id,
slot=base.slot,
stream_id=base.stream_id,
seq=seq,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
payload=base.payload,
)
return spec.data()
def _is_voice_burst(packet: bytes) -> bool:
return ((packet[15] & 0x30) >> 4) == HBPF_VOICE
def _parse_dmrd_args(packet: bytes) -> dict[str, Any]:
bits = packet[15]
return {
"peer_id": packet[11:15],
"rf_src": packet[5:8],
"dst_id": packet[8:11],
"seq": packet[4],
"slot": 2 if bits & 0x80 else 1,
"call_type": "group",
"frame_type": (bits & 0x30) >> 4,
"dtype_vseq": bits & 0xF,
"stream_id": packet[16:20],
"data": packet,
}
def _wrap_send_peer_trace(hbp: HBPProtocol) -> dict[str, Any]:
"""Record per-packet send_peer accept/reject without modifying production code."""
trace: dict[str, Any] = {"accepted": [], "rejected": []}
orig = hbp.send_peer
def _traced(peer_id: bytes, packet: bytes, *, _skip_dual_expand: bool = False) -> None:
route_pkt = packet
peer = hbp._peers.get(peer_id)
if peer is not None and packet[:4] == b"DMRD":
from adn_server.application.routing.downlink import remap_dmrd_for_peer
route_pkt = remap_dmrd_for_peer(packet, peer, hbp._config, peer_id=peer_id)
opts_ok = hbp._peer_should_receive_dmrd(peer_id, packet)
gate_ok = hbp._peer_would_accept_group_dmrd(
peer_id, packet if not _skip_dual_expand else route_pkt, routed=_skip_dual_expand,
)
row = {
"tgid": int_id(route_pkt[8:11]),
"stream": route_pkt[16:20],
"opts_ok": opts_ok,
"gate_ok": gate_ok,
"blocked": not (opts_ok and (gate_ok if packet[:4] == b"DMRD" else True)),
}
if row["blocked"]:
trace["rejected"].append(row)
else:
trace["accepted"].append(row)
orig(peer_id, packet, _skip_dual_expand=_skip_dual_expand)
hbp.send_peer = _traced # type: ignore[method-assign]
return trace
def _inject_obp_at(
scenario: DeterministicScenario,
packet: bytes,
*,
pkt_time: float,
) -> None:
scenario.clock.now = pkt_time
with _patch_harness_time(scenario.clock), patch_routing_wall_time(scenario.clock):
scenario.routing.dmrd_received(
"OBP-CL",
ingress_pkt_time=pkt_time,
obp_use_parsed=True,
obp_hops=b"\x00\x00\x00\x00",
obp_source_server=bytes_4(9990),
**_parse_dmrd_args(packet),
)
def _run_interleaved_qsos(
scenario: DeterministicScenario,
*,
voice_bursts: int = _VOICE_BURSTS,
interval_s: float = _INTERVAL_S,
) -> dict[str, Any]:
t0 = scenario.clock.time()
step = 0
tx_stamp_changes = 0
prev_tx_stream: bytes | None = None
hs_addr = ("127.0.0.1", 62031)
delivered_a: list[bytes] = []
delivered_b: list[bytes] = []
blocked_a_checks: list[dict[str, Any]] = []
def _slot_st() -> dict[str, Any]:
proto = scenario.protocols["MASTER-A"]
assert isinstance(proto, HBPProtocol)
return proto.STATUS.setdefault(2, {})
def _record_delivery(pkt: bytes) -> None:
tgid = int_id(pkt[8:11])
sid = pkt[16:20]
if tgid == _TG_A and sid == bytes_4(_STREAM_A):
delivered_a.append(pkt)
elif tgid == _TG_B and sid == bytes_4(_STREAM_B):
delivered_b.append(pkt)
packets_plan: list[tuple[str, bytes]] = []
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_DATA_SYNC,
dtype_vseq=HBPF_SLT_VHEAD,
),
)
)
for seq in range(1, voice_bursts + 1):
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_VOICE,
dtype_vseq=(seq % 4) or 1,
seq=seq,
),
)
)
for label, tgid, stream, rf in (
("A", _TG_A, _STREAM_A, _RF_A),
("B", _TG_B, _STREAM_B, _RF_B),
):
packets_plan.append(
(
label,
_dmrd_packet(
tgid=tgid,
stream_id=stream,
rf_src=rf,
frame_type=HBPF_DATA_SYNC,
dtype_vseq=HBPF_SLT_VTERM,
seq=99,
),
)
)
scenario.capture.packets.clear()
transport = scenario.protocols["MASTER-A"]
assert isinstance(transport, HBPProtocol)
rec = transport.transport
assert isinstance(rec, RecordingTransport)
rec.clear()
hbp = transport
hs = bytes_4(_HS_PEER)
peer = hbp._peers[hs]
for label, pkt in packets_plan:
pkt_time = t0 + step * interval_s
step += 1
_inject_obp_at(scenario, pkt, pkt_time=pkt_time)
cur = _slot_st().get("TX_STREAM_ID")
if cur != prev_tx_stream:
tx_stamp_changes += 1
prev_tx_stream = cur
bits = pkt[15] | 0x80
route_pkt = pkt[:15] + bytes([bits]) + pkt[16:]
ctx = hbp._downlink_ctx()
if (
label == "A"
and int_id(pkt[8:11]) == _TG_A
and pkt[16:20] == bytes_4(_STREAM_A)
and (pkt[15] & 0xF) != HBPF_SLT_VTERM
):
blocked = peer_slot_blocks_downlink(ctx, hs, peer, route_pkt, pkt_time=pkt_time)
if blocked:
blocked_a_checks.append(
{
"pkt_time": pkt_time,
"step": step - 1,
"slot_st": copy.deepcopy(_slot_st()),
"peer_slots": copy.deepcopy(ctx.peer_voice_slots.get(hs, {})),
"reason": diagnose_downlink_block(ctx, hs, peer, route_pkt, pkt_time),
}
)
for sent_pkt, addr in rec.sent:
if addr == hs_addr:
_record_delivery(sent_pkt)
rec.clear()
return {
"tx_stamp_changes": tx_stamp_changes,
"start_tx_events": len(
[ev for ev in (scenario.report_factory.events if scenario.report_factory else []) if ",START,TX," in ev]
),
"delivered_a": delivered_a,
"delivered_b": delivered_b,
"blocked_a_checks": blocked_a_checks,
"final_slot_st": copy.deepcopy(_slot_st()),
}
def diagnose_downlink_block(
ctx: DownlinkContext,
peer_id: bytes,
peer: dict[str, Any],
route_pkt: bytes,
pkt_time: float,
) -> str:
"""Pin which gate branch would block this downlink (mirrors helpers/downlink)."""
if route_pkt[:4] != b"DMRD":
return "not_dmrd"
stream_id = route_pkt[16:20]
incoming_tgid_b = route_pkt[8:11]
hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0)
peer_slots = ctx.peer_voice_slots.get(bytes_4(int_id(peer_id)))
pk = bytes_4(int_id(peer_id))
for voice_slot in (2,):
hang_row = ctx.peer_voice_hangtime.get(pk, {}).get(voice_slot)
slot_st = ctx.status.get(voice_slot, {})
if _peer_transmit_hangtime_blocks(hang_row, incoming_tgid_b, pkt_time, hang):
return f"helpers.py:_peer_transmit_hangtime_blocks voice_slot={voice_slot}"
active = (peer_slots or {}).get(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)
age = pkt_time - active_time
if active.get("ingress"):
return f"helpers.py:peer_hotspot_voice_slot_busy ingress voice_slot={voice_slot}"
if active.get("bridge_hold") and active_tgid and incoming_tgid != active_tgid:
if age <= hang:
return (
f"helpers.py:peer_hotspot_voice_slot_busy bridge_hold "
f"active_tg={active_tgid} incoming={incoming_tgid}"
)
if active_stream and stream_id:
if active_stream == stream_id:
pass
elif active_tgid and active_tgid == incoming_tgid:
if age < STREAM_TO:
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"same_tg_different_stream active={active_stream!r} incoming={stream_id!r}"
)
elif active_tgid and incoming_tgid != active_tgid:
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"elif active_tgid={active_tgid} incoming_tgid={incoming_tgid}"
)
else:
return "helpers.py:peer_hotspot_voice_slot_busy else branch (stream/tg mismatch)"
elif isinstance(active, dict):
if not (active_tgid and active_tgid == incoming_tgid):
return (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"no_active_stream active_tg={active_tgid} incoming={incoming_tgid}"
)
if peer_single_blocks_foreign_same_tg_downlink(
peer, pk, voice_slot, incoming_tgid_b, peer_slots, ctx.sys_cfg, now=pkt_time,
):
return f"helpers.py:peer_single_blocks_foreign_same_tg_downlink voice_slot={voice_slot}"
if peer_single_same_tg_foreign_tx_blocks(
peer, pk, incoming_tgid_b, stream_id, slot_st, ctx.sys_cfg, pkt_time=pkt_time,
):
return (
f"helpers.py:peer_single_same_tg_foreign_tx_blocks "
f"TX_STREAM_ID={slot_st.get('TX_STREAM_ID')!r} slot_active={slot_has_active_voice(slot_st, pkt_time)}"
)
if bytes_4(int_id(slot_st.get("RX_PEER", b""))) == pk:
rx_active = (
slot_st.get("RX_TYPE") is not None
and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM
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 (
f"helpers.py:peer_hotspot_voice_slot_busy "
f"rx_active stream_mismatch RX_STREAM_ID={slot_st.get('RX_STREAM_ID')!r}"
)
if _peer_status_rx_hangtime_blocks(
peer_id, slot_st, incoming_tgid_b, pkt_time, hang,
):
return f"helpers.py:_peer_status_rx_hangtime_blocks voice_slot={voice_slot}"
if hbp_slot_blocks_group_voice_for_peer(
slot_st,
peer_id,
incoming_tgid_b,
stream_id,
pkt_time,
hang,
per_peer=True,
peers=ctx.peers,
peer_slots=peer_slots,
peer_hang_row=hang_row,
voice_slot=voice_slot,
sys_cfg=ctx.sys_cfg,
):
return f"helpers.py:hbp_slot_blocks_group_voice_for_peer voice_slot={voice_slot}"
if peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt, pkt_time=pkt_time):
return "downlink.py:peer_slot_blocks_downlink (composite)"
return "not_blocked"
def test_flip_flop_slot_st_blocks_active_listen_on_cross_tg() -> None:
"""Pure unit: peer listening on TG A must not RX TG B while session is open."""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_b,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_B),
"RX_TYPE": HBPF_SLT_VTERM,
}
assert peer_hotspot_voice_slot_busy(
hs, 2, stream_b, bytes_3(_TG_B), slot_st, peer_slots, None, now + 0.06, 5.0,
peer=peer, sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False},
)
route_b = _dmrd_packet(
tgid=_TG_B,
stream_id=_STREAM_B,
rf_src=_RF_B,
frame_type=HBPF_VOICE,
dtype_vseq=1,
seq=1,
)
route_b = route_b[:15] + bytes([route_b[15] | 0x80]) + route_b[16:]
ctx = DownlinkContext(
config={"SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}},
system_name="MASTER-A",
sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False, "MODE": "MASTER"},
peers={hs: peer},
status={2: slot_st},
peer_voice_slots={hs: copy.deepcopy(peer_slots)},
connected_count=2,
)
reason = diagnose_downlink_block(ctx, hs, peer, route_b, now + 0.06)
assert "active_tgid=3109050 incoming_tgid=52090" in reason
def test_flip_flop_slot_st_same_stream_not_blocked_despite_foreign_tx_stamp() -> None:
"""Same stream id on active listen must pass even when flat TX row shows other stream."""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_b,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_B),
"RX_TYPE": HBPF_SLT_VTERM,
}
assert not peer_hotspot_voice_slot_busy(
hs, 2, stream_a, bytes_3(_TG_A), slot_st, peer_slots, None, now + 0.06, 5.0,
peer=peer, sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False},
)
def test_rx_active_own_peer_foreign_stream_blocks_active_downlink() -> None:
"""When STATUS RX shows the hotspot mid-TX on another stream, stream-A downlink blocks.
helpers.py peer_hotspot_voice_slot_busy lines 411-418:
``rx_active and stream_id != slot_st.get('RX_STREAM_ID')`` → True.
"""
now = 1_000_000.0
hs = bytes_4(_HS_PEER)
peer = {"OPTIONS": f"TS2={_TG_A},{_TG_B};SINGLE=0;".encode()}
stream_a = bytes_4(_STREAM_A)
stream_b = bytes_4(_STREAM_B)
peer_slots = {
2: {"stream_id": stream_a, "tgid": _TG_A, "time": now, "ingress": False},
}
slot_st = {
"RX_PEER": hs,
"RX_STREAM_ID": stream_b,
"RX_TYPE": HBPF_SLT_VHEAD,
"RX_TIME": now,
"RX_TGID": bytes_3(_TG_B),
"TX_PEER": bytes_4(_OBP_PEER),
"TX_STREAM_ID": stream_a,
"TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": now,
"TX_TGID": bytes_3(_TG_A),
}
route_a = _dmrd_packet(
tgid=_TG_A,
stream_id=_STREAM_A,
rf_src=_RF_A,
frame_type=HBPF_VOICE,
dtype_vseq=1,
seq=1,
)
route_a = route_a[:15] + bytes([route_a[15] | 0x80]) + route_a[16:]
ctx = DownlinkContext(
config={"SYSTEMS": {"MASTER-A": {"MODE": "MASTER"}}},
system_name="MASTER-A",
sys_cfg={"GROUP_HANGTIME": 5.0, "SINGLE_MODE": False, "MODE": "MASTER"},
peers={hs: peer},
status={2: slot_st},
peer_voice_slots={hs: copy.deepcopy(peer_slots)},
connected_count=2,
)
reason = diagnose_downlink_block(ctx, hs, peer, route_a, now + 0.06)
assert "rx_active stream_mismatch" in reason
assert peer_slot_blocks_downlink(ctx, hs, peer, route_a, pkt_time=now + 0.06)
def test_concurrent_obp_tx_row_flip_flop(caplog: pytest.LogCaptureFixture) -> None:
"""BUG: shared STATUS[2] TX_STREAM_ID must change at most twice per stream (VHEAD), not per packet."""
caplog.set_level(logging.INFO)
scenario, _hbp, _transport = _build_stack()
result = _run_interleaved_qsos(scenario, voice_bursts=6)
# Expect: 2 streams × 1 VHEAD stamp each (+ maybe VTERM) — not ~2×packets.
max_expected_stamps = 6
assert result["tx_stamp_changes"] <= max_expected_stamps, (
f"TX_STREAM_ID flip-flop: {result['tx_stamp_changes']} changes "
f"(expected <= {max_expected_stamps}); final STATUS[2]={result['final_slot_st']}"
)
assert result["start_tx_events"] <= max_expected_stamps, (
f"START,TX report spam: {result['start_tx_events']} events"
)
def test_concurrent_obp_active_stream_delivered_without_gaps() -> None:
"""After VHEAD, stream-A voice should not be dropped mid-QSO (minimal harness)."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
result = _run_interleaved_qsos(scenario, voice_bursts=_VOICE_BURSTS)
voice_a = [p for p in result["delivered_a"] if _is_voice_burst(p)]
rejected_a = [r for r in send_trace["rejected"] if r["tgid"] == _TG_A]
expected_voice = _VOICE_BURSTS
# Minimal OBP-only harness delivers all A packets; production drops use RX/status
# corruption paths (see test_rx_active_own_peer_foreign_stream_blocks_active_downlink).
assert not rejected_a, f"unexpected stream-A send_peer rejects: {rejected_a}"
assert len(voice_a) == expected_voice, (
f"stream A voice delivered {len(voice_a)}/{expected_voice}"
)
# Stream B may be dropped for the listening hotspot — not asserted.
def test_concurrent_obp_stream_b_blocked_on_listen_session() -> None:
"""Stream B must be dropped while hotspot listens to stream A (one QSO per slot)."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
_run_interleaved_qsos(scenario, voice_bursts=4)
rejected_b = [r for r in send_trace["rejected"] if r["tgid"] == _TG_B]
assert rejected_b, "stream B should be blocked for listening hotspot"
ctx = hbp._downlink_ctx()
hs = bytes_4(_HS_PEER)
peer = hbp._peers[hs]
sample = _dmrd_packet(
tgid=_TG_B, stream_id=_STREAM_B, rf_src=_RF_B,
frame_type=HBPF_VOICE, dtype_vseq=1, seq=1,
)
sample = sample[:15] + bytes([sample[15] | 0x80]) + sample[16:]
reason = diagnose_downlink_block(ctx, hs, peer, sample, scenario.clock.time())
assert (
"active_tgid=3109050 incoming_tgid=52090" in reason
or "_peer_transmit_hangtime_blocks" in reason
or "peer_hotspot" in reason
or "peer_slot" in reason
), f"stream B blocked but unexpected gate: {reason}"
def test_concurrent_obp_diagnose_active_stream_block_branch() -> None:
"""Identify the exact branch that drops active stream-A packets mid-QSO."""
scenario, hbp, transport = _build_stack()
send_trace = _wrap_send_peer_trace(hbp)
result = _run_interleaved_qsos(scenario, voice_bursts=_VOICE_BURSTS)
rejected_a = [
r for r in send_trace["rejected"]
if r["tgid"] == _TG_A and r["stream"] == bytes_4(_STREAM_A)
]
if not rejected_a and not result["blocked_a_checks"]:
# After per-leg OBP bridge TX wiring, flat TX_STREAM_ID no longer flip-flops per packet.
if result["tx_stamp_changes"] <= 4:
pytest.skip("no stream-A rejects; per-leg bridge TX fix removed flat-row flip-flop")
assert result["tx_stamp_changes"] > 4, "expected TX row flip-flop in concurrent OBP harness"
pytest.skip("no stream-A send_peer rejects — inspect flip-flop collateral (START,TX spam)")
reasons: dict[str, int] = {}
for row in result["blocked_a_checks"]:
reasons[row["reason"]] = reasons.get(row["reason"], 0) + 1
for rej in rejected_a:
if not rej["opts_ok"]:
key = "udp_hbp.py:_peer_should_receive_dmrd OPTIONS/eligibility"
elif not rej["gate_ok"]:
key = "udp_hbp.py:_peer_would_accept_group_dmrd -> peer_slot_blocks_downlink"
else:
key = "unknown"
reasons[key] = reasons.get(key, 0) + 1
top = max(reasons, key=reasons.get)
assert "peer_hotspot" in top or "peer_slot" in top or "send_peer" in top or "udp_hbp" in top, (
f"unexpected block reasons: {reasons}"
)

Powered by TurnKey Linux.