fix(talker-alias): passthrough from radio, dedupe relay and logs

Buffer DMRA before VHEAD, wait for hotspot TA in both mode, relay
passthrough DMRA once per stream, and dedupe DEBUG policy/embed lines.
pull/1/head
Rodrigo Pérez 4 months ago
parent 67710558e8
commit 84ccbc8f5c

@ -37,6 +37,10 @@ class BridgeLcTaMixin:
"OPENBRIDGE",
)
def _master_ta_wait_source(self, source_system: str) -> bool:
"""``both`` mode waits for hotspot TA only on HBP MASTER sources (not OBP)."""
return self._config.get("SYSTEMS", {}).get(source_system, {}).get("MODE") == "MASTER"
def _cancel_both_ta_wait(self, source_system: str, stream_id: bytes) -> None:
wait = self._both_ta_wait.pop(self._both_ta_key(source_system, stream_id), None)
if wait and wait.get("timer") and getattr(wait["timer"], "cancel", None):
@ -109,6 +113,8 @@ class BridgeLcTaMixin:
continue
if st.get("TX_TA_ON"):
continue
if st.get("_ta_embed_kind") == "passthrough" and not force_inject:
continue
if st.get("TX_STREAM_ID") == stream_id and st.get("TX_RFS") == rf_src:
target = st.get("_ta_target_system", source_system)
self._init_talker_alias_embed(
@ -137,9 +143,10 @@ class BridgeLcTaMixin:
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
if not blocks or not passthrough_complete(blocks):
return
key = self._both_ta_key(source_system, stream_id)
wait = self._both_ta_wait.get(key)
if wait:
relay_key = self._both_ta_key(source_system, stream_id)
already_relayed = relay_key in self._passthrough_relayed
wait = self._both_ta_wait.get(relay_key)
if wait and not already_relayed:
wait["peer"] = peer_id
self._cancel_both_ta_wait(source_system, stream_id)
for target_system in wait["targets"]:
@ -147,7 +154,9 @@ class BridgeLcTaMixin:
source_system, target_system, rf_src, stream_id, peer_id, force=True,
)
self._apply_both_ta_embed(source_system, rf_src, stream_id)
self._relay_passthrough_dmra(source_system, peer_id, rf_src, stream_id)
if not already_relayed:
self._relay_passthrough_dmra(source_system, peer_id, rf_src, stream_id)
self._passthrough_relayed.add(relay_key)
def _relay_passthrough_dmra(
self,
@ -179,9 +188,39 @@ class BridgeLcTaMixin:
)
src_cfg = self._config.get("SYSTEMS", {}).get(source_system, {})
if src_cfg.get("MODE") == "MASTER" and src_cfg.get("REPEAT", True):
self._talker_alias.clear_stream(source_system, stream_id)
self.send_talker_alias_local_repeat(source_system, peer_id, rf_src, stream_id)
def _dispatch_talker_alias_on_bridge_open(
self,
target_st: dict[str, Any],
source_system: str,
target_system: str,
rf_src: bytes,
stream_id: bytes,
source_peer: bytes,
) -> None:
"""Prepare TA on a new bridged HBP leg (VHEAD or first burst after hangtime)."""
settings = talker_alias_settings(self._config, source_system)
if not settings["enabled"]:
return
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
have_source_ta = bool(blocks and passthrough_complete(blocks))
if settings["mode"] == "both" and self._master_ta_wait_source(source_system) and not have_source_ta:
self._register_both_ta_wait(
source_system, target_system, rf_src, stream_id, source_peer,
)
return
self._init_talker_alias_embed(
target_st,
source_system,
target_system,
rf_src,
stream_id,
)
self._send_talker_alias_to_target(
source_system, target_system, rf_src, stream_id, source_peer,
)
def _send_talker_alias_to_target(
self,
source_system: str,
@ -201,8 +240,12 @@ class BridgeLcTaMixin:
return
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
have_passthrough = bool(blocks and passthrough_complete(blocks))
if have_passthrough and force:
self._talker_alias.clear_stream(target_system, stream_id)
if have_passthrough:
if force:
if not self._talker_alias.should_resend_passthrough_dmra(target_system, stream_id):
return
elif not self._talker_alias.should_send_on_vhead(target_system, stream_id):
return
elif not self._talker_alias.should_send_on_vhead(target_system, stream_id):
return
if not have_passthrough:
@ -224,7 +267,11 @@ class BridgeLcTaMixin:
except Exception as e:
logger.warning("(ROUTER) send_dmra_to_system %s failed: %s", target_system, e)
return
self._talker_alias.mark_dmra_sent(target_system, stream_id)
self._talker_alias.mark_dmra_sent(
target_system,
stream_id,
kind="passthrough" if have_passthrough else "inject",
)
sid = int_id(stream_id)
if peer_count:
logger.debug(
@ -311,6 +358,7 @@ class BridgeLcTaMixin:
def clear_talker_alias_stream(self, system_name: str, stream_id: bytes) -> None:
"""Release per-stream TA dedupe state after VTERM."""
self._cancel_both_ta_wait(system_name, stream_id)
self._passthrough_relayed.discard(self._both_ta_key(system_name, stream_id))
self._talker_alias.clear_stream(system_name, stream_id)
if not self._get_protocols:
return
@ -375,6 +423,11 @@ class BridgeLcTaMixin:
st["TX_TA_PHASE"] = 0
# First B1–B4 cycle carries group LC; TA on the next cycle.
st["TX_TA_ON"] = False
blocks = self._get_stream_dmra_blocks(source_system, stream_id)
if blocks and passthrough_complete(blocks) and not force_inject:
st["_ta_embed_kind"] = "passthrough"
else:
st["_ta_embed_kind"] = "inject"
def _rewrite_embed_lc(
self,
@ -412,4 +465,5 @@ class BridgeLcTaMixin:
st.pop("TX_TA_PHASE", None)
st.pop("TX_TA_BLOCK_COUNT", None)
st.pop("TX_TA_ON", None)
st.pop("_ta_embed_kind", None)

@ -89,6 +89,8 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi
self._talker_alias = TalkerAliasUseCases(config, ta_emblc_encoder=ta_emblc_encoder)
# (source_system, stream_id) -> {rf_src, peer, targets, timer}
self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {}
# Passthrough DMRA/embed relay already applied for this source stream.
self._passthrough_relayed: set[tuple[str, bytes]] = set()
def _send_bridge_event(self, event: str | bytes) -> bool:
"""Send BRDG_EVENT via ReportingUseCases when REPORT is enabled."""
@ -562,12 +564,13 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi
_ts_st["TX_H_LC"] = bptc.encode_header_lc(dst_lc)
_ts_st["TX_T_LC"] = bptc.encode_terminator_lc(dst_lc)
_ts_st["TX_EMB_LC"] = self._encode_emblc(dst_lc)
self._init_talker_alias_embed(
self._dispatch_talker_alias_on_bridge_open(
_ts_st,
system_name,
entry["SYSTEM"],
rf_src,
stream_id,
peer_id,
)
logger.info(
"(%s) Conference Bridge: %s, Call Bridged to HBP System: %s TS: %s, TGID: %s",
@ -583,10 +586,6 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi
int_id(entry_tgid_b),
)
)
# First successful forward to HBP (may not be VHEAD if earlier frames were hangtime-blocked).
self._send_talker_alias_to_target(
system_name, entry["SYSTEM"], rf_src, stream_id, peer_id,
)
_ts_st["TX_TIME"] = pkt_time
_ts_st["TX_TYPE"] = dtype_vseq
# Slot bit rewrite (legacy bridge.py 457-460 / 770-773)

@ -133,17 +133,60 @@ class TalkerAliasUseCases:
self._config = config
self._ta_emblc = ta_emblc_encoder
self._sent_streams: set[tuple[str, bytes]] = set()
self._embed_logged: set[tuple[str, bytes]] = set()
self._sent_kind: dict[tuple[str, bytes], str] = {}
self._embed_logged: set[tuple[str, str, bytes]] = set()
self._policy_logged: set[tuple[str, str, bytes, str]] = set()
def clear_dmra_sent(self, system_name: str, stream_id: bytes) -> None:
key = (system_name, stream_id)
self._sent_streams.discard(key)
self._sent_kind.pop(key, None)
def clear_stream(self, system_name: str, stream_id: bytes) -> None:
self._sent_streams.discard((system_name, stream_id))
self._embed_logged.discard((system_name, stream_id))
self.clear_dmra_sent(system_name, stream_id)
self._embed_logged = {k for k in self._embed_logged if k[2] != stream_id}
self._policy_logged = {k for k in self._policy_logged if k[2] != stream_id}
def should_send_on_vhead(self, target_system: str, stream_id: bytes) -> bool:
return (target_system, stream_id) not in self._sent_streams
def mark_dmra_sent(self, target_system: str, stream_id: bytes) -> None:
self._sent_streams.add((target_system, stream_id))
def should_resend_passthrough_dmra(self, target_system: str, stream_id: bytes) -> bool:
"""True when passthrough DMRA should replace a prior inject on this leg."""
key = (target_system, stream_id)
if key not in self._sent_streams:
return True
return self._sent_kind.get(key) == "inject"
def mark_dmra_sent(self, target_system: str, stream_id: bytes, *, kind: str) -> None:
key = (target_system, stream_id)
self._sent_streams.add(key)
self._sent_kind[key] = kind
def _log_policy_once(
self,
source_system: str,
target_system: str,
stream_id: bytes,
kind: str,
text: str,
*,
suffix: str = "",
via: str,
) -> None:
log_key = (source_system, target_system, stream_id, kind)
if log_key in self._policy_logged:
return
self._policy_logged.add(log_key)
if kind == "passthrough":
logger.debug(
"(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s",
source_system, text, via, target_system, int_id(stream_id),
)
else:
logger.debug(
"(%s) *TALKER ALIAS* inject '%s'%s via %s -> %s stream %s",
source_system, text, suffix, via, target_system, int_id(stream_id),
)
def packets_for_stream(
self,
@ -168,16 +211,14 @@ class TalkerAliasUseCases:
if not have_passthrough:
return None
text = decode_ta_from_blocks(blocks)
logger.debug(
"(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
self._log_policy_once(
source_system, target, stream_id, "passthrough", text, via=via,
)
return passthrough_packets_from_blocks(rf_src, blocks)
if mode == "inject":
text = format_talker_alias_text(self._config, rf_src)
logger.debug(
"(%s) *TALKER ALIAS* inject '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
self._log_policy_once(
source_system, target, stream_id, "inject", text, via=via,
)
return build_dmra_packets(rf_src, text, settings["text_formats"][0])
# both: prefer the source's own TA. If a valid MMDVM DMRA buffer arrived,
@ -185,17 +226,15 @@ class TalkerAliasUseCases:
# through unchanged in the DMRD voice, so do NOT inject a template here.
if have_passthrough:
text = decode_ta_from_blocks(blocks)
logger.debug(
"(%s) *TALKER ALIAS* passthrough '%s' via %s -> %s stream %s",
source_system, text, via, target, int_id(stream_id),
self._log_policy_once(
source_system, target, stream_id, "passthrough", text, via=via,
)
return passthrough_packets_from_blocks(rf_src, blocks)
# Legacy resolve_ta (both): inject template when no DMRA buffer yet (VHEAD).
text = format_talker_alias_text(self._config, rf_src)
suffix = " (no source TA yet)" if fallback_inject else ""
logger.debug(
"(%s) *TALKER ALIAS* inject '%s'%s via %s -> %s stream %s",
source_system, text, suffix, via, target, int_id(stream_id),
self._log_policy_once(
source_system, target, stream_id, "inject", text, suffix=suffix, via=via,
)
return build_dmra_packets(rf_src, text, settings["text_formats"][0])
@ -225,7 +264,7 @@ class TalkerAliasUseCases:
mode = settings["mode"]
target = target_system or source_system
via = "repeat" if source_system == target else "bridge"
log_key = (source_system, stream_id)
log_key = (source_system, target, stream_id)
def _log(text: str, suffix: str) -> None:
if log_key in self._embed_logged:

@ -318,6 +318,30 @@ class HBPProtocol(DatagramProtocol):
if not self._ta_buffer_enabled() or not stream_id:
return
self._dmra_rf_stream[(peer_id, rf_src)] = stream_id
self._promote_dmra_provisional_key(rf_src, stream_id)
def _promote_dmra_provisional_key(self, provisional: bytes, stream_id: bytes) -> None:
"""Move TA blocks buffered under ``rf_src`` (pre-VHEAD) to the real stream id."""
if provisional == stream_id:
return
entry = self._dmra_by_stream.pop(provisional, None)
if not entry:
return
target = self._dmra_by_stream.setdefault(
stream_id,
{
"blocks": {},
"rf_src": entry.get("rf_src", b""),
"peer": entry.get("peer", b""),
"last": entry.get("last", time.time()),
},
)
target["blocks"].update(entry.get("blocks", {}))
target["last"] = max(float(target.get("last", 0)), float(entry.get("last", 0)))
if entry.get("rf_src"):
target["rf_src"] = entry["rf_src"]
if entry.get("peer"):
target["peer"] = entry["peer"]
def store_dmra_packet(self, peer_id: bytes, data: bytes) -> None:
"""Buffer one DMRA block from a hotspot (MASTER receive path)."""
@ -327,7 +351,8 @@ class HBPProtocol(DatagramProtocol):
rf_src, block_id, payload = parsed
stream_id = self._dmra_rf_stream.get((peer_id, rf_src))
if not stream_id:
return
# Legacy hblink: DMRA may arrive before the first DMRD; key by rf_src until then.
stream_id = rf_src
now = time.time()
entry = self._dmra_by_stream.setdefault(
stream_id,

@ -83,7 +83,8 @@ def test_both_mode_without_ta_still_rewrites_group_embedded_lc() -> None:
)
ts_st = scenario.protocols["MASTER-B"].STATUS[2]
assert ts_st.get("TX_TA_EMB") is not None # both mode injects template at VHEAD
# both + MASTER source: TA overlay deferred until source TA or 2s fallback (see docs).
assert ts_st.get("TX_TA_EMB") is None
bursts = [p for p in scenario.capture.for_system("MASTER-B") if p.fields["dtype_vseq"] == 1]
assert bursts, "voice burst B should be forwarded to MASTER-B"
bits = bitarray(endian="big")

@ -0,0 +1,65 @@
"""DMRA may arrive before DMRD VHEAD — legacy buffers under rf_src until stream note."""
from __future__ import annotations
from adn_server.domain import bytes_3, bytes_4
from adn_server.domain.talker_alias import build_dmra_packet, decode_ta_from_blocks
from adn_server.infrastructure.hbp_constants import DMRA
from tests.support.hbp_repeat_stack import build_hbp_repeat_stack
from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.talker_alias.test_mmdvm_wire import mmdvm_wire_blocks
_PEER_TX = bytes_4(730039210)
_PEER_RX = bytes_4(730039101)
_ADDR_TX = ("10.0.0.1", 62001)
_ADDR_RX = ("10.0.0.2", 62002)
def test_dmra_before_vhead_passthrough_on_repeat() -> None:
"""Regression: early DMRA was dropped when stream_id mapping did not exist yet."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.config["GLOBAL"]["TALKER_ALIAS_MODE"] = "passthrough"
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
rf_src = bytes_3(3120001)
text = "CE5RPY Radio"
for block_id, payload in mmdvm_wire_blocks(text).items():
stack.inject(build_dmra_packet(rf_src, block_id, payload), _ADDR_TX)
base = PacketSpec(peer_id=730039210, rf_src=3120001, dst_id=7304, slot=2, stream_id=0xA1B2C3D4)
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
assert stack.dmra_capture, "passthrough DMRA should be sent on VHEAD"
packets, exclude = stack.dmra_capture[0]
assert exclude == _PEER_TX
assert all(p[:4] == DMRA for p in packets)
rebuilt = {i: packets[i][8:15] for i in range(len(packets))}
assert decode_ta_from_blocks(rebuilt) == text
dmra_to_rx = [p for p, _ in stack.transport.sent if p[:4] == DMRA and _ == _ADDR_RX]
assert dmra_to_rx, "listening peer should receive passthrough DMRA"
def test_promote_provisional_dmra_buffer() -> None:
stack = build_hbp_repeat_stack(talker_alias=True)
hbp = stack.hbp
peer = _PEER_TX
rf_src = bytes_3(3120001)
stream_id = bytes_4(0x12345678)
blocks = mmdvm_wire_blocks("TA early")
for block_id, payload in blocks.items():
hbp.store_dmra_packet(peer, build_dmra_packet(rf_src, block_id, payload))
assert hbp.get_dmra_blocks(rf_src) is not None
assert hbp.get_dmra_blocks(stream_id) is None
hbp.note_dmrd_stream(peer, rf_src, stream_id)
promoted = hbp.get_dmra_blocks(stream_id)
assert promoted is not None
assert decode_ta_from_blocks(promoted) == "TA early"
assert hbp.get_dmra_blocks(rf_src) is None

@ -21,9 +21,18 @@ def test_talker_alias_clear_stream_allows_resend() -> None:
key = ("MASTER-B", stream_id)
assert ta.should_send_on_vhead("MASTER-B", stream_id) is True
ta.mark_dmra_sent("MASTER-B", stream_id)
ta.mark_dmra_sent("MASTER-B", stream_id, kind="inject")
assert ta.should_send_on_vhead("MASTER-B", stream_id) is False
ta.clear_stream("MASTER-B", stream_id)
assert key not in ta._sent_streams
assert ta.should_send_on_vhead("MASTER-B", stream_id) is True
def test_clear_dmra_sent_preserves_embed_log_dedupe() -> None:
ta = make_talker_alias_use_cases(talker_alias_config())
stream_id = PacketSpec(stream_id=0x93939393).data()[16:20]
ta._embed_logged.add(("MASTER-A", "MASTER-B", stream_id))
ta.clear_dmra_sent("MASTER-B", stream_id)
assert ("MASTER-A", "MASTER-B", stream_id) in ta._embed_logged

@ -0,0 +1,55 @@
"""Talker Alias DMRA relay and log deduplication."""
from __future__ import annotations
from adn_server.application.bridge_use_cases import BridgeUseCases
from adn_server.domain import bytes_3, bytes_4
from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter
from adn_server.infrastructure.talker_alias_emblc import default_ta_emblc_encoder
from adn_server.domain.dmr.bptc import encode_emblc
from tests.harness.scenarios import make_talker_alias_use_cases, talker_alias_config
from tests.talker_alias.test_mmdvm_wire import mmdvm_wire_blocks
def test_should_resend_passthrough_only_after_inject() -> None:
ta = make_talker_alias_use_cases(talker_alias_config())
stream_id = bytes_4(0xA1B2C3D4)
assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is True
ta.mark_dmra_sent("SYSTEM", stream_id, kind="inject")
assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is True
ta.mark_dmra_sent("SYSTEM", stream_id, kind="passthrough")
assert ta.should_resend_passthrough_dmra("SYSTEM", stream_id) is False
def test_on_dmra_fragment_stored_relays_passthrough_once() -> None:
"""Second completion callback must not re-send passthrough DMRA on the same stream."""
config = talker_alias_config()
config["GLOBAL"]["TALKER_ALIAS_MODE"] = "both"
config["SYSTEMS"]["MASTER-A"]["REPEAT"] = True
blocks = mmdvm_wire_blocks("CE5RPY")
sent: list[str] = []
def _send_dmra(target_system: str, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
sent.append(target_system)
return 1
bridge = BridgeUseCases(
InMemoryBridgeRouter(),
config,
send_dmra_to_system=_send_dmra,
get_dmra_blocks=lambda _sys, _sid: blocks,
encode_emblc=encode_emblc,
ta_emblc_encoder=default_ta_emblc_encoder,
)
peer = bytes_4(1001)
rf_src = bytes_3(3120001)
stream_id = bytes_4(0xB1B2B3B4)
bridge.on_dmra_fragment_stored("MASTER-A", peer, rf_src, stream_id)
first_count = len(sent)
assert first_count == 1
bridge.on_dmra_fragment_stored("MASTER-A", peer, rf_src, stream_id)
assert len(sent) == first_count
Loading…
Cancel
Save

Powered by TurnKey Linux.