fix: filter standalone DMRA downlink by TG subscription

Apply peer_should_receive_group_voice to DMRA UDP (same rules as DMRD).
Debug branch for hotspot TA leak on unrelated static TGs.
pull/10/head
Rodrigo Pérez 3 months ago
parent 5f9c5d779d
commit 41f55201dd

@ -259,6 +259,30 @@ class LcTaMixin:
source_system, target_system, rf_src, stream_id, source_peer,
)
def _resolve_dmra_downlink_route(
self,
target_system: str,
stream_id: bytes,
) -> tuple[int, int] | None:
"""Slot and TGID for standalone DMRA downlink on ``target_system``."""
if not self._get_protocols:
return None
proto = self._get_protocols().get(target_system)
status = getattr(proto, "STATUS", None) if proto else None
if not isinstance(status, dict):
return None
for slot in (1, 2):
st = status.get(slot)
if not isinstance(st, dict):
continue
if st.get("TX_STREAM_ID") == stream_id and st.get("TX_TGID"):
return slot, int_id(st["TX_TGID"])
if st.get("REP_STREAM_ID") == stream_id:
tg = st.get("REP_TGID") or st.get("RX_TGID") or st.get("TX_TGID")
if tg:
return slot, int_id(tg)
return None
def _send_talker_alias_to_target(
self,
source_system: str,
@ -300,8 +324,23 @@ class LcTaMixin:
if not packets:
return
exclude = source_peer if target_system == source_system else None
route = self._resolve_dmra_downlink_route(target_system, stream_id)
if route is None:
logger.warning(
"(%s) *TALKER ALIAS* stream %s DMRA not sent: no slot/TG route on target",
target_system,
int_id(stream_id),
)
return
slot, tgid = route
try:
peer_count = self._send_dmra_to_system(target_system, packets, exclude_peer=exclude)
peer_count = self._send_dmra_to_system(
target_system,
packets,
exclude_peer=exclude,
slot=slot,
tgid=tgid,
)
except Exception as e:
logger.warning("(ROUTER) send_dmra_to_system %s failed: %s", target_system, e)
return
@ -313,8 +352,8 @@ class LcTaMixin:
sid = int_id(stream_id)
if peer_count:
logger.debug(
"(%s) *TALKER ALIAS* stream %s sent %d DMRA block(s) to %d peer(s)",
target_system, sid, len(packets), peer_count,
"(%s) *TALKER ALIAS* stream %s sent %d DMRA block(s) to %d peer(s) TG %s slot %s",
target_system, sid, len(packets), peer_count, tgid, slot,
)
elif exclude:
logger.debug(
@ -360,6 +399,7 @@ class LcTaMixin:
if st.get("REP_STREAM_ID") != stream_id:
dst_lc = LC_OPT + dst_id + rf_src
st["REP_STREAM_ID"] = stream_id
st["REP_TGID"] = dst_id
st["REP_EMB_LC"] = self._encode_emblc(dst_lc)
self._init_talker_alias_embed(
st, system_name, system_name, rf_src, stream_id,

@ -282,10 +282,21 @@ def run_peer_server(
system_name: str,
packets: list[bytes],
exclude_peer: bytes | None = None,
*,
slot: int | None = None,
tgid: int | None = None,
) -> int:
p = protocols.get(system_name)
if p is not None and hasattr(p, "send_dmra_system"):
return int(p.send_dmra_system(packets, exclude_peer=exclude_peer) or 0)
return int(
p.send_dmra_system(
packets,
exclude_peer=exclude_peer,
slot=slot,
tgid=tgid,
)
or 0
)
return 0
def get_dmra_blocks(system_name: str, stream_id: bytes) -> dict[int, bytes] | None:

@ -650,23 +650,70 @@ class HBPProtocol(DatagramProtocol):
self._ta_voice_acc.pop(stream_id, None)
self._ta_decoded_logged.discard(stream_id)
def send_dmra_to_peers(self, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
"""Send DMRA packets to logged-in peers (MASTER downlink). Returns peer count."""
def send_dmra_to_peers(
self,
packets: list[bytes],
exclude_peer: bytes | None = None,
*,
slot: int | None = None,
tgid: int | None = None,
) -> int:
"""Send DMRA packets to logged-in peers (MASTER downlink). Returns peer count.
When ``slot`` and ``tgid`` are set, only peers that would receive group voice
on that route get standalone DMRA (same rules as ``_peer_should_receive_dmrd``).
"""
if self._config.get("MODE") != "MASTER":
return 0
if slot is None or tgid is None:
logger.warning(
"(%s) *TALKER ALIAS* DMRA not sent: missing slot/tgid (stream filter)",
self._system,
)
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:
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,
)
continue
for pkt in packets:
self.send_peer(peer, pkt)
sent += 1
return sent
def send_dmra_system(self, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
def send_dmra_system(
self,
packets: list[bytes],
exclude_peer: bytes | None = None,
*,
slot: int | None = None,
tgid: int | None = None,
) -> int:
"""Send DMRA on this system link (MASTER → peers, PEER → upstream master)."""
if self._config.get("MODE") == "MASTER":
return self.send_dmra_to_peers(packets, exclude_peer=exclude_peer)
return self.send_dmra_to_peers(
packets, exclude_peer=exclude_peer, slot=slot, tgid=tgid,
)
if self._config.get("MODE") == "PEER":
for pkt in packets:
self.send_master(pkt)

@ -447,6 +447,7 @@ class DeterministicScenario:
target_system: str,
packets: list[bytes],
exclude_peer: bytes | None = None,
**kwargs: Any,
) -> int:
self.dmra_capture.append(
CapturedDmra(target_system, list(packets), exclude_peer=exclude_peer)

@ -60,7 +60,7 @@ def test_repeat_downlink_embeds_group_lc_when_talker_alias_enabled() -> None:
"""Regression: raw REPEAT copy left embed LC all-zero; MMDVM never saw TA."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
base = _base_spec()
uplink_burst = DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1)
@ -84,7 +84,7 @@ def test_repeat_sends_dmra_and_embed_on_vhead() -> None:
"""MMDVM may ignore DMRA UDP; embed in DMRD must still be prepared on VHEAD."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
@ -101,7 +101,7 @@ def test_repeat_sends_dmra_and_embed_on_vhead() -> None:
def test_repeat_leaves_burst_payload_unchanged_when_talker_alias_disabled() -> None:
stack = build_hbp_repeat_stack(talker_alias=False)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)
stack.transport.clear()
@ -118,7 +118,7 @@ def test_repeat_overlays_ta_on_second_superframe_cycle() -> None:
"""After bursts B–E, the next B should carry TA embed (alternate superframes)."""
stack = build_hbp_repeat_stack(talker_alias=True)
stack.register_peer(_PEER_TX, _ADDR_TX)
stack.register_peer(_PEER_RX, _ADDR_RX)
stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
base = _base_spec()
stack.inject_spec(DeterministicScenario.voice_head_spec(base), _ADDR_TX)

@ -126,10 +126,19 @@ def build_hbp_repeat_stack(
protocols: dict[str, Any] = {system_name: hbp}
dmra_capture: list[tuple[list[bytes], bytes | None]] = []
def _send_dmra(target_system: str, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
def _send_dmra(
target_system: str,
packets: list[bytes],
exclude_peer: bytes | None = None,
*,
slot: int | None = None,
tgid: int | None = None,
) -> int:
proto = protocols[target_system]
dmra_capture.append((list(packets), exclude_peer))
return proto.send_dmra_to_peers(packets, exclude_peer=exclude_peer)
return proto.send_dmra_to_peers(
packets, exclude_peer=exclude_peer, slot=slot, tgid=tgid,
)
def _get_dmra_blocks(_system: str, stream_id: bytes) -> dict[int, bytes] | None:
return hbp.get_dmra_blocks(stream_id)

@ -42,7 +42,7 @@ def test_dmra_before_vhead_passthrough_on_repeat() -> None:
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)
stack.register_peer(_PEER_RX, _ADDR_RX, options="TS2=7304;")
rf_src = bytes_3(3120001)
text = "CE5RPY Radio"

@ -51,23 +51,37 @@ def test_on_dmra_fragment_stored_relays_passthrough_once() -> None:
config["SYSTEMS"]["MASTER-A"]["REPEAT"] = True
blocks = mmdvm_wire_blocks("CE5RPY")
sent: list[str] = []
stream_id = bytes_4(0xB1B2B3B4)
def _send_dmra(target_system: str, packets: list[bytes], exclude_peer: bytes | None = None) -> int:
def _send_dmra(
target_system: str,
packets: list[bytes],
exclude_peer: bytes | None = None,
**kwargs: object,
) -> int:
sent.append(target_system)
return 1
class _Proto:
STATUS = {
2: {
"TX_STREAM_ID": stream_id,
"TX_TGID": bytes_3(52090),
}
}
bridge = RoutingUseCases(
InMemoryAclRouter(),
config,
InMemorySubscriptionStore(),
send_dmra_to_system=_send_dmra,
get_protocols=lambda: {"MASTER-A": _Proto()},
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)

Loading…
Cancel
Save

Powered by TurnKey Linux.