diff --git a/CHANGELOG.md b/CHANGELOG.md index 78649e1..7b2bf22 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,21 @@ All notable changes to **adn-server** are documented here. +## v2.0.1 (2026-06-20) + +### Bug Fixes + +- Filter standalone DMRA downlink by TG subscription + ([#10](https://github.com/ce5rpy/ADN-DMR-Peer-Server/pull/10), + [`0dcdf54`](https://github.com/ce5rpy/ADN-DMR-Peer-Server/commit/0dcdf54a07328f7203b02809cd7e26255987258e)) + +### Chores + +- Run releases on master only, not develop + ([#10](https://github.com/ce5rpy/ADN-DMR-Peer-Server/pull/10), + [`0dcdf54`](https://github.com/ce5rpy/ADN-DMR-Peer-Server/commit/0dcdf54a07328f7203b02809cd7e26255987258e)) + + ## v2.0.0 (2026-06-19) ### Chores diff --git a/README.md b/README.md index 9800300..6e1b5aa 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ # ADN DMR Peer Server -**Version 2.0.0** — pairs with **adn-monitor 2.0.0** (report v2 slim wire + JSON HELLO). +**Version 2.0.1** — pairs with **adn-monitor 2.0.0** (report v2 slim wire + JSON HELLO). ADN DMR conference bridge server. Configuration is YAML; the codebase follows clean architecture (domain, application, infrastructure). v2 adds integrated **PROXY**, **SubscriptionStore** routing, report v2 to the monitor, and a unified **`adn-server.py`** entrypoint (`--echo`, `--doctor`, `--no-proxy`). See [CHANGELOG.md](CHANGELOG.md). diff --git a/pyproject.toml b/pyproject.toml index ef00ce7..cde7675 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -8,7 +8,7 @@ build-backend = "setuptools.build_meta" [project] name = "adn-server" -version = "2.0.0" +version = "2.0.1" description = "ADN DMR Peer Server" readme = "README.md" license = { text = "GPL-3.0-or-later" } diff --git a/src/adn_server/application/routing/lc_ta.py b/src/adn_server/application/routing/lc_ta.py index 739e7e4..c95d682 100644 --- a/src/adn_server/application/routing/lc_ta.py +++ b/src/adn_server/application/routing/lc_ta.py @@ -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, diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index d379b88..bfa2078 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -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: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 4b2753f..177d607 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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) diff --git a/tests/harness/deterministic.py b/tests/harness/deterministic.py index 8f41eef..b8e334f 100644 --- a/tests/harness/deterministic.py +++ b/tests/harness/deterministic.py @@ -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) diff --git a/tests/infrastructure/test_hbp_repeat_talker_alias.py b/tests/infrastructure/test_hbp_repeat_talker_alias.py index 2a9ddf4..bdd5dba 100644 --- a/tests/infrastructure/test_hbp_repeat_talker_alias.py +++ b/tests/infrastructure/test_hbp_repeat_talker_alias.py @@ -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) diff --git a/tests/support/hbp_repeat_stack.py b/tests/support/hbp_repeat_stack.py index c458af8..e46f73e 100644 --- a/tests/support/hbp_repeat_stack.py +++ b/tests/support/hbp_repeat_stack.py @@ -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) diff --git a/tests/talker_alias/test_dmra_early_buffer.py b/tests/talker_alias/test_dmra_early_buffer.py index e0cb7f0..b977707 100644 --- a/tests/talker_alias/test_dmra_early_buffer.py +++ b/tests/talker_alias/test_dmra_early_buffer.py @@ -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" diff --git a/tests/talker_alias/test_relay_dedupe.py b/tests/talker_alias/test_relay_dedupe.py index be462fc..14ee6a9 100644 --- a/tests/talker_alias/test_relay_dedupe.py +++ b/tests/talker_alias/test_relay_dedupe.py @@ -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)