fix: cross-slot downlink and REPEAT monitor activity

Cross-slot TG downlink, DMRA index parity, and REPEAT monitor START/TX reporting.
pull/22/head
ce5rpy 3 months ago committed by GitHub
parent e05dc00318
commit f561b2685c
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -338,7 +338,7 @@ def peer_owns_multi_dynamic_ua(
*,
peer_id: bytes | None = None,
) -> bool:
"""True when SINGLE=0 peer has keyed this non-static dynamic TG on ``slot``."""
"""True when SINGLE=0 peer has keyed this non-static dynamic TG (either slot)."""
if not sys_cfg or peer_single_mode(peer, sys_cfg):
return False
if peer_id is None:
@ -352,8 +352,12 @@ def peer_owns_multi_dynamic_ua(
per_peer = store.get(pk)
if not isinstance(per_peer, dict):
return False
slot_set = per_peer.get(int(slot))
return isinstance(slot_set, set) and int(tgid) in slot_set
tgid_i = int(tgid)
for voice_slot in (1, 2):
slot_set = per_peer.get(voice_slot)
if isinstance(slot_set, set) and tgid_i in slot_set:
return True
return False
def register_peer_ua_session(
@ -582,17 +586,21 @@ def peer_single_blocks_group_voice(
peer_id: bytes | None = None,
now: float | None = None,
) -> bool:
"""True when SINGLE=1 peer must not receive downlink for ``tgid`` on ``slot``.
"""True when SINGLE=1 peer must not receive downlink for ``tgid``.
With an active session on TG X, every other TG (static or dynamic) is blocked
until the TIMER expires or a new local TX replaces the session.
"""
del slot
if not sys_cfg:
return False
locked = peer_single_exclusive_tgid(peer, slot, sys_cfg, peer_id=peer_id, now=now)
if locked is None:
for voice_slot in (1, 2):
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now,
)
if locked is not None and int(tgid) != locked:
return True
return False
return int(tgid) != locked
def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool:
@ -602,7 +610,11 @@ def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool
slot while self-service lists the TG on the other.
"""
del slot
return peer_options_static_tg_slot(peer, tgid) is not None
from adn_server.application.report.payloads import parse_peer_options_static
ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS"))
tg = str(tgid)
return tg in ts1 or tg in ts2
def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None:
@ -620,6 +632,56 @@ def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None:
return None
def synthetic_group_dmrd_route_packet(slot: int, tgid: int) -> bytes:
"""Minimal DMRD for inject-only downlink fan-out lookup (slot + TG only)."""
bits = 0x80 if int(slot) == 2 else 0
return b"DMRD" + b"\x00" * 4 + bytes_3(tgid) + b"\x00" * 4 + bytes([bits]) + b"\x00" * 38
def repeat_downlink_report_slot(
wire_slot: int,
tgid: int,
peers: dict[Any, Any],
downlink_peer_ids: tuple[bytes, ...],
sys_cfg: dict[str, Any] | None,
) -> int:
"""Monitor timeslot for REPEAT downlink START/TX (OBP bridge uses target TS, not wire slot).
When every downlink peer maps the TG to the same OPTIONS or UA slot, use that slot so
CTABLE chips match RF (cross-slot static/dynamic). Otherwise fall back to wire slot.
"""
display_slots: set[int] = set()
tgid_i = int(tgid)
for peer_id in downlink_peer_ids:
peer = peers.get(peer_id)
if not isinstance(peer, dict):
continue
static_slot = peer_options_static_tg_slot(peer, tgid_i)
if static_slot is not None:
display_slots.add(static_slot)
continue
if not sys_cfg:
continue
pk = bytes_4(int_id(peer_id))
for voice_slot in (1, 2):
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id,
)
if locked is not None and int(locked) == tgid_i:
display_slots.add(voice_slot)
store = sys_cfg.get("_PEER_UA_MULTI_TGS")
if isinstance(store, dict):
per_peer = store.get(pk)
if isinstance(per_peer, dict):
for voice_slot in (1, 2):
slot_set = per_peer.get(voice_slot)
if isinstance(slot_set, set) and tgid_i in slot_set:
display_slots.add(voice_slot)
if len(display_slots) == 1:
return display_slots.pop()
return int(wire_slot)
def peer_single_blocks_uplink(
peer: dict[str, Any],
peer_id: bytes,
@ -651,8 +713,14 @@ def _peer_owns_dynamic_ua(
return False
if not sys_cfg:
return False
locked = peer_single_exclusive_tgid(peer, slot, sys_cfg, peer_id=peer_id, now=now)
return locked is not None and int(tgid) == locked
tgid_i = int(tgid)
for voice_slot in (1, 2):
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now,
)
if locked is not None and tgid_i == locked:
return True
return False
def peer_should_receive_group_voice(

@ -268,27 +268,7 @@ class RoutingUseCases(
# legs; re-arm OPTIONS static destinations before forward (inject-only merged lists).
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
self.apply_static_tg_to_bridge(dst_int)
has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
)
if not has_source and systems_cfg.get(system_name, {}).get("MODE") == "MASTER":
self.options_config_for_system(system_name)
has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
)
# Do not call ensure_dynamic_relay here for 9990–9999 when BRIDGES["9990"] already exists:
# ensure_dynamic_relay replaces the whole table and only sets the source MASTER ACTIVE; every
# other system (including ECHO with TO_TYPE NONE) becomes ACTIVE False — to_target then has
# no active target for the echo path (legacy _seed_echo_routing_table / make_bridges keeps ECHO
# ACTIVE True; activation of the source row is via in-band ON on VTERM, bridge_master ~3465).
if not has_source:
logger.debug(
"(ROUTER) No matching source rule for TG %s from %s slot %s (ACTIVE), not forwarding",
relay_table_key, system_name, bridge_match_slot,
)
return True
# Legacy bridge.py: BRDG_EVENT (OBP group/vcsbk START/END handled in _obp_group_voice_router_obp / post-forward VTERM)
# Legacy bridge_master routerHBP: BRDG_EVENT on VHEAD/VTERM before bridge source scan.
_obp_grp = source_is_obp and call_type in ("group", "vcsbk")
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD:
if not _obp_grp:
@ -342,6 +322,26 @@ class RoutingUseCases(
duration,
)
)
has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
)
if not has_source and systems_cfg.get(system_name, {}).get("MODE") == "MASTER":
self.options_config_for_system(system_name)
has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
)
# Do not call ensure_dynamic_relay here for 9990–9999 when BRIDGES["9990"] already exists:
# ensure_dynamic_relay replaces the whole table and only sets the source MASTER ACTIVE; every
# other system (including ECHO with TO_TYPE NONE) becomes ACTIVE False — to_target then has
# no active target for the echo path (legacy _seed_echo_routing_table / make_bridges keeps ECHO
# ACTIVE True; activation of the source row is via in-band ON on VTERM, bridge_master ~3465).
if not has_source:
logger.debug(
"(ROUTER) No matching source rule for TG %s from %s slot %s (ACTIVE), not forwarding",
relay_table_key, system_name, bridge_match_slot,
)
return True
# ── Exact port of legacy bridge.py routerOBP/routerHBP forwarding to targets ──
pkt_time = time.time()
dmrpkt = data[20:53] if len(data) >= 53 else b""

@ -51,7 +51,9 @@ from ...application.routing.helpers import (
peer_single_exclusive_tgid,
register_peer_ua_session,
resolve_voice_peer_id,
repeat_downlink_report_slot,
seed_peer_ua_session_from_status,
synthetic_group_dmrd_route_packet,
tg4000_reset_on_vhead,
)
from ...application.routing.peer_downlink_index import (
@ -233,6 +235,7 @@ class HBPProtocol(DatagramProtocol):
self._connected_peer_count = 0
self._config_push_delayed = None
self._config_push_throttle = ConfigPushThrottle()
self._repeat_downlink_report_start: dict[bytes, float] = {}
self._refresh_connected_peer_count()
else:
self._peers = {}
@ -428,6 +431,57 @@ class HBPProtocol(DatagramProtocol):
)
)
def _emit_repeat_downlink_group_voice_report(
self,
action: str,
tx_peer: bytes,
rf_src: bytes,
dst_id: bytes,
wire_slot: int,
stream_id: bytes,
downlink_peers: tuple[bytes, ...],
*,
pkt_time: float,
duration: float = 0.0,
) -> None:
"""START/END TX for REPEAT downlink peers (monitor CTABLE, same role as OBP bridge TX leg)."""
if not downlink_peers:
return
report = self._report
if report is None or not self._CONFIG.get("REPORTS", {}).get("REPORT", True):
return
if not hasattr(report, "send_routing_event"):
return
tgid = int_id(dst_id)
report_slot = repeat_downlink_report_slot(
wire_slot, tgid, self._peers, downlink_peers, self._config,
)
systems_cfg = self._CONFIG.get("SYSTEMS", {})
report_peer = resolve_voice_peer_id(tx_peer, rf_src, self._system, systems_cfg)
if action == "START":
report.send_routing_event(
"GROUP VOICE,START,TX,{},{},{},{},{},{}".format(
self._system,
int_id(stream_id),
int_id(report_peer),
int_id(rf_src),
report_slot,
tgid,
)
)
elif action == "END":
report.send_routing_event(
"GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format(
self._system,
int_id(stream_id),
int_id(report_peer),
int_id(rf_src),
report_slot,
tgid,
duration,
)
)
def _apply_tg4000_reset(
self,
peer_id: bytes,
@ -672,29 +726,11 @@ class HBPProtocol(DatagramProtocol):
)
return 0
sent = 0
connected = self._cached_connected_peer_count()
store = self._get_subscription_store() if self._get_subscription_store else None
for peer in self._peers:
route_pkt = synthetic_group_dmrd_route_packet(slot, tgid)
for peer in self._iter_downlink_peers(route_pkt):
if exclude_peer and peer == exclude_peer:
continue
if not peer_should_receive_group_voice(
self._peers[peer],
slot,
tgid,
peer_id=peer,
system=self._system,
bridges=None,
subscription_store=store,
connected_count=connected,
sys_cfg=self._config,
):
logger.debug(
"(%s) *TALKER ALIAS* DMRA skip peer %s (not subscribed to TG %s slot %s)",
self._system,
int_id(peer),
tgid,
slot,
)
if not self._peer_should_receive_dmrd(peer, route_pkt):
continue
for pkt in packets:
self.send_peer(peer, pkt)
@ -1059,9 +1095,49 @@ class HBPProtocol(DatagramProtocol):
)
_repeat_tail = b"".join([_data[15:20], _dmrpkt_out, _data[53:]])
_repeat_pkt = b"".join([_data[:11], _peer_id, _repeat_tail])
for _peer in self._iter_downlink_peers(_repeat_pkt):
if _peer != _peer_id:
_downlink_peers = tuple(
_p for _p in self._iter_downlink_peers(_repeat_pkt) if _p != _peer_id
)
for _peer in _downlink_peers:
self.send_peer(_peer, _repeat_pkt)
if _downlink_peers and _slot in self.STATUS:
_slot_st = self.STATUS[_slot]
if (
_frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VHEAD
and _stream_id != _slot_st.get("RX_STREAM_ID")
):
self._repeat_downlink_report_start[_stream_id] = pkt_time
self._emit_repeat_downlink_group_voice_report(
"START",
_peer_id,
_rf_src,
_dst_id,
_slot,
_stream_id,
_downlink_peers,
pkt_time=pkt_time,
)
elif (
_frame_type == HBPF_DATA_SYNC
and _dtype_vseq == HBPF_SLT_VTERM
and _stream_id == _slot_st.get("RX_STREAM_ID")
and _slot_st.get("RX_TYPE") != HBPF_SLT_VTERM
):
_start = self._repeat_downlink_report_start.pop(
_stream_id, _slot_st.get("RX_START", pkt_time),
)
self._emit_repeat_downlink_group_voice_report(
"END",
_peer_id,
_rf_src,
_dst_id,
_slot,
_stream_id,
_downlink_peers,
pkt_time=pkt_time,
duration=pkt_time - float(_start),
)
# 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,

@ -37,6 +37,7 @@ from adn_server.application.routing.helpers import (
peer_single_exclusive_tgid,
register_peer_ua_multi_tg,
register_peer_ua_session,
repeat_downlink_report_slot,
seed_peer_ua_session_from_status,
)
@ -57,6 +58,44 @@ def test_static_tg_on_opposite_slot_receives_group_voice() -> None:
assert peer_should_receive_group_voice(peer, 2, 730170, connected_count=8)
def test_static_tg_listed_on_both_slots_receives_group_voice() -> None:
peer = {"OPTIONS": b"TS1=730444;TS2=730444;"}
assert peer_receives_group_tgid(peer, 1, 730444)
assert peer_receives_group_tgid(peer, 2, 730444)
assert peer_options_static_tg_slot(peer, 730444) is None
assert peer_should_receive_group_voice(peer, 1, 730444, connected_count=8)
def test_dynamic_tg_on_opposite_slot_receives_group_voice() -> None:
"""SINGLE=0 keyed dynamic on TS2 is heard when voice arrives on TS1."""
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
sys_cfg = _sys_cfg()
peer_id = _peer_id()
register_peer_ua_multi_tg(peer, peer_id, 2, 730444, sys_cfg)
assert peer_should_receive_group_voice(
peer, 1, 730444, peer_id=peer_id, connected_count=3, sys_cfg=sys_cfg,
)
def test_repeat_downlink_report_slot_cross_slot_static() -> None:
peers = {
_peer_id(): {"OPTIONS": b"TS2=730444;"},
}
slot = repeat_downlink_report_slot(
1, 730444, peers, (_peer_id(),), _sys_cfg(),
)
assert slot == 2
def test_repeat_downlink_report_slot_dynamic_on_ts2() -> None:
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
peer_id = _peer_id()
sys_cfg = _sys_cfg()
register_peer_ua_multi_tg(peer, peer_id, 2, 730444, sys_cfg)
slot = repeat_downlink_report_slot(1, 730444, {peer_id: peer}, (peer_id,), sys_cfg)
assert slot == 2
def test_single_blocks_other_static_tg_while_session_on_static() -> None:
"""Indigo on 7305 blocks RX on 730 even when both are in OPTIONS."""
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}

@ -72,6 +72,61 @@ def test_repeat_only_reaches_peers_with_matching_options() -> None:
assert other_pkts == []
def test_repeat_cross_slot_static_tg_downlink() -> None:
"""Voice on TS1 reaches hotspot that lists the TG only on TS2 (PR #2 parity)."""
stack = build_hbp_repeat_stack(talker_alias=False, system_name="MASTER-A")
stack.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
stack.config["REPORTS"] = {"REPORT": True}
stack.hbp._CONFIG = stack.config
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS1={_TG};")
stack.register_peer(_PEER_RX_MATCH, _ADDR_MATCH, options=f"TS2={_TG};")
stack.register_peer(_PEER_RX_OTHER, _ADDR_OTHER, options="TS2=91;")
spec = PacketSpec(
peer_id=int.from_bytes(_PEER_TX, "big"),
rf_src=7300444,
dst_id=_TG,
slot=1,
stream_id=0x11223344,
payload=b"\x00" * 33,
)
stack.inject(DeterministicScenario.voice_burst_spec(spec, seq=1, dtype_vseq=1).data(), _ADDR_TX)
match_pkts = [p for p in stack.transport.for_addr(_ADDR_MATCH) if p[:4] == DMRD]
other_pkts = [p for p in stack.transport.for_addr(_ADDR_OTHER) if p[:4] == DMRD]
assert len(match_pkts) == 1
assert other_pkts == []
def test_repeat_cross_slot_emits_downlink_start_tx_report() -> None:
"""REPEAT downlink reports START,TX with OPTIONS slot (monitor CTABLE parity with OBP bridge)."""
stack = build_hbp_repeat_stack(talker_alias=False, system_name="MASTER-A")
stack.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
stack.config["REPORTS"] = {"REPORT": True}
stack.hbp._CONFIG = stack.config
stack.register_peer(_PEER_TX, _ADDR_TX, options=f"TS1={_TG};")
stack.register_peer(_PEER_RX_MATCH, _ADDR_MATCH, options=f"TS2={_TG};")
spec = PacketSpec(
peer_id=int.from_bytes(_PEER_TX, "big"),
rf_src=7300444,
dst_id=_TG,
slot=1,
stream_id=0x22334455,
payload=b"\x00" * 33,
)
stack.inject(DeterministicScenario.voice_head_spec(spec).data(), _ADDR_TX)
tx_starts = [
e for e in stack.report_factory.events
if e.startswith(f"GROUP VOICE,START,TX,{stack.system_name},")
]
assert len(tx_starts) == 1
parts = tx_starts[0].split(",")
assert int(parts[7]) == 2
assert int(parts[8]) == _TG
def test_bridge_downlink_send_peers_filters_by_options() -> None:
stack = _inject_proxy_stack()
burst = _voice_burst()

@ -70,6 +70,7 @@ class HbpRepeatStack:
hbp: HBPProtocol
bridge: RoutingUseCases
transport: RecordingTransport
report_factory: FakeReportFactory
dmra_capture: list[tuple[list[bytes], bytes | None]] = field(default_factory=list)
def register_peer(
@ -120,7 +121,8 @@ def build_hbp_repeat_stack(
sys_cfg["USE_ACL"] = False
transport = RecordingTransport()
hbp = HBPProtocol(system_name, config, router=_AclRouter())
report_factory = FakeReportFactory()
hbp = HBPProtocol(system_name, config, report_factory=report_factory, router=_AclRouter())
hbp.transport = transport # type: ignore[assignment]
protocols: dict[str, Any] = {system_name: hbp}
@ -143,7 +145,6 @@ def build_hbp_repeat_stack(
def _get_dmra_blocks(_system: str, stream_id: bytes) -> dict[int, bytes] | None:
return hbp.get_dmra_blocks(stream_id)
report_factory = FakeReportFactory()
bridge = RoutingUseCases(
InMemoryAclRouter(),
config,
@ -156,6 +157,7 @@ def build_hbp_repeat_stack(
encode_emblc=encode_emblc,
ta_emblc_encoder=default_ta_emblc_encoder,
)
hbp._dmrd_received = bridge.dmrd_received
hbp._on_talker_alias_repeat_prepare = bridge.prepare_talker_alias_local_repeat
hbp._on_talker_alias_repeat_burst = bridge.rewrite_repeat_voice_burst
hbp._on_talker_alias_stream_end = bridge.clear_talker_alias_stream
@ -166,5 +168,6 @@ def build_hbp_repeat_stack(
hbp=hbp,
bridge=bridge,
transport=transport,
report_factory=report_factory,
dmra_capture=dmra_capture,
)

Loading…
Cancel
Save

Powered by TurnKey Linux.