From d76de25cbc8a26be1c77a62c7f8f282d13a454ab Mon Sep 17 00:00:00 2001 From: yo Date: Thu, 24 Sep 2026 23:43:40 +0200 Subject: [PATCH] feat(plugins): voice_slot_for_tg and monitor TX for plugin voice What a voice plugin needs from the core that it can't see itself: - ServerContext.voice_slot_for_tg(tg): the MASTER slot to speak a TG on now, or None while every slot is busy. Same choice announcements make today (dynamic UA session or active bridge leg of the TG on that MASTER first, then TS2, then TS1), ported to PluginIngress so the announcement code can leave the core. Only for granted talkgroups. - PluginIngress reports GROUP VOICE START/END,TX on the MASTER a plugin stream plays on, the line announcements send today: its hotspots hear the stream but no bridge leg reports it. A stream cut by a radio, or never terminated, still gets its END. Co-Authored-By: Claude Opus 5.5 --- .../plugins/application/context.py | 2 + .../plugins/application/ingress.py | 83 ++++++++++++++++++- .../plugins/application/manager.py | 12 ++- .../application/plugins/application/sender.py | 9 ++ .../infrastructure/bootstrap/peer_server.py | 7 +- tests/application/test_plugin_sender.py | 9 ++ tests/routing/test_plugin_send_routing.py | 48 +++++++++++ 7 files changed, 162 insertions(+), 8 deletions(-) diff --git a/src/adn_server/application/plugins/application/context.py b/src/adn_server/application/plugins/application/context.py index 28693f7..e384dc0 100644 --- a/src/adn_server/application/plugins/application/context.py +++ b/src/adn_server/application/plugins/application/context.py @@ -35,3 +35,5 @@ class ServerContext: call_later: Callable[..., Any] # Present only for a plugin listed in PLUGINS.send (see domain/send.py). send_dmrd: Callable[[bytes], bool] | None = None + # With group voice granted: the MASTER slot to speak a TG on now, None while all are busy. + voice_slot_for_tg: Callable[[int], int | None] | None = None diff --git a/src/adn_server/application/plugins/application/ingress.py b/src/adn_server/application/plugins/application/ingress.py index d3254ad..01878fe 100644 --- a/src/adn_server/application/plugins/application/ingress.py +++ b/src/adn_server/application/plugins/application/ingress.py @@ -29,7 +29,7 @@ from typing import Any, Callable from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM from ....domain.mesh_engine import server_id_bytes from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd -from ...routing.helpers import slot_voice_held_by_other_stream +from ...routing.helpers import master_dynamic_tg_slots, slot_voice_held_by_other_stream from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind logger = logging.getLogger(__name__) @@ -54,17 +54,60 @@ class PluginIngress: get_protocols: Callable[[], dict[str, Any]], send_local: Callable[[str, bytes], None], clock: Callable[[], float] = time.time, + send_routing_event: Callable[[str], None] | None = None, ) -> None: self._routing = routing self._config = config self._get_protocols = get_protocols self._send_local = send_local self._clock = clock + self._send_routing_event = send_routing_event + # Plugin voice streams on air: stream_id -> (master, slot, tg, rf_src, start). + self._on_air: dict[bytes, tuple[str, int, int, int, float]] = {} def set_routing(self, routing: Any) -> None: """Bootstrap builds the plugin manager before routing; routing is set before any load.""" self._routing = routing + def voice_slot_for_tg(self, tg: int) -> int | None: + """The MASTER slot a plugin should speak ``tg`` on now, or None while every slot is busy. + + Same choice as scheduled announcements: the slot where ``tg`` is a dynamic UA + session or has an active bridge leg on that MASTER first, then TS2, then TS1. + """ + master = announcement_ptt_system(self._config) + proto = self._get_protocols().get(master) if master else None + status = getattr(proto, "STATUS", None) + sys_cfg = self._config.get("SYSTEMS", {}).get(master or "", {}) + if not status or sys_cfg.get("MODE") != "MASTER": + return None + dynamic = master_dynamic_tg_slots(sys_cfg, int(tg)) + bridged = self._active_bridge_slots(int(tg), master) + now = self._clock() + for ts in dict.fromkeys([*sorted(dynamic, reverse=True), *sorted(bridged, reverse=True), 2, 1]): + slot = status.get(ts) + if slot and not self._slot_busy(slot, ts in dynamic or ts in bridged, now): + return ts + return None + + @staticmethod + def _slot_busy(slot: dict[str, Any], tg_lives_here: bool, now: float) -> bool: + if slot_voice_held_by_other_stream(slot, b"", now): + return True # live voice on the slot, in or out + if slot.get("RX_TYPE") != HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM: + return True # an outside QSO still in its hang time + if tg_lives_here: + return False + return not (slot.get("RX_TYPE") == HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM) + + def _active_bridge_slots(self, tg: int, master: str) -> set[int]: + table = self._routing.routing_table_for_report() if self._routing is not None else {} + return { + int(be["TS"]) + for be in table.get(str(tg), []) + if isinstance(be, dict) and be.get("SYSTEM") == master and be.get("ACTIVE") and be.get("TS") is not None + } + def deliver(self, pkt: bytes, plugin: str) -> bool: header = parse_dmrd_header(pkt) kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq) if header else None @@ -84,10 +127,11 @@ class PluginIngress: slot = getattr(proto, "STATUS", {}).get(header.slot) if proto is not None else None if slot is None: return False - if slot_voice_held_by_other_stream(slot, header.stream_id, now): - return False - if self._route(master, pkt, now, server_id, plugin) is not True: + if slot_voice_held_by_other_stream(slot, header.stream_id, now) or ( + self._route(master, pkt, now, server_id, plugin) is not True + ): self._release(slot, header.stream_id) + self._off_air(header.stream_id, now) return False is_term = header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VTERM slot["TX_TYPE"] = HBPF_SLT_VTERM if is_term else HBPF_SLT_VHEAD @@ -96,8 +140,39 @@ class PluginIngress: slot["TX_TGID"] = header.dst_id slot["TX_TIME"] = now self._send_local(master, pkt) + if header.stream_id not in self._on_air: + self._on_air_start(master, header, now) + if is_term: + self._off_air(header.stream_id, now) return True + def _on_air_start(self, master: str, header: Any, now: float) -> None: + """Monitor TX on the MASTER itself: its hotspots hear the stream, which no bridge leg reports.""" + if len(self._on_air) >= 64: # plugins that never sent a terminator + for stream_id in list(self._on_air)[:32]: + self._off_air(stream_id, now) + tg, rf_src = int_id(header.dst_id), int_id(header.rf_src) + self._on_air[header.stream_id] = (master, header.slot, tg, rf_src, now) + self._report("START", master, header.stream_id, header.slot, tg, rf_src) + + def _off_air(self, stream_id: bytes, now: float) -> None: + entry = self._on_air.pop(stream_id, None) + if entry is not None: + master, slot, tg, rf_src, start = entry + self._report("END", master, stream_id, slot, tg, rf_src, now - start) + + def _report( + self, action: str, master: str, stream_id: bytes, slot: int, tg: int, rf_src: int, duration: float | None = None + ) -> None: + # Same line announcements send (VoiceUseCases._emit_announcement_voice_event). + if self._send_routing_event is None: + return + parts = ["GROUP VOICE", action, "TX", master, str(int_id(stream_id)), str(rf_src), str(rf_src), str(slot), str(tg)] + if duration is not None: + parts.append(f"{duration:.2f}") + parts.append("1") + self._send_routing_event(",".join(parts)) + def _route(self, master: str, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool | None: return inject_plugin_dmrd(self._routing, master, pkt, pkt_time=now, server_id=server_id, plugin=plugin) diff --git a/src/adn_server/application/plugins/application/manager.py b/src/adn_server/application/plugins/application/manager.py index be69cbb..d503189 100644 --- a/src/adn_server/application/plugins/application/manager.py +++ b/src/adn_server/application/plugins/application/manager.py @@ -136,10 +136,18 @@ class PluginManager: The sender re-reads the permission on every frame, so revoking it on SIGHUP is immediate; granting it to an already loaded plugin needs that plugin reloaded. """ - if self._sender_factory is None or send_permission(server_config, name) is None: + permission = send_permission(server_config, name) if self._sender_factory is not None else None + if permission is None: return self._ctx + sender = self._sender_factory(name) + if permission.group_voice_tgs: + logger.info( + "(PLUGIN-MANAGER) %s may send unit data and group voice on TG %s (PLUGINS.send)", + name, sorted(permission.group_voice_tgs), + ) + return replace(self._ctx, send_dmrd=sender, voice_slot_for_tg=getattr(sender, "voice_slot_for_tg", None)) logger.info("(PLUGIN-MANAGER) %s may send unit data (PLUGINS.send)", name) - return replace(self._ctx, send_dmrd=self._sender_factory(name)) + return replace(self._ctx, send_dmrd=sender) def shutdown_all(self) -> None: for name in list(self._loaded): diff --git a/src/adn_server/application/plugins/application/sender.py b/src/adn_server/application/plugins/application/sender.py index 6d7e715..80a8b7e 100644 --- a/src/adn_server/application/plugins/application/sender.py +++ b/src/adn_server/application/plugins/application/sender.py @@ -54,6 +54,7 @@ class PluginDmrdSender: call_from_reactor: Callable[..., Any], clock: Callable[[], float] = time.monotonic, in_reactor_thread: Callable[[], bool] = lambda: False, + slot_for_tg: Callable[[int], int | None] | None = None, ) -> None: self._plugin = plugin self._config = server_config @@ -61,6 +62,7 @@ class PluginDmrdSender: self._call_from_reactor = call_from_reactor self._clock = clock self._in_reactor_thread = in_reactor_thread + self._slot_for_tg = slot_for_tg self._lock = threading.Lock() self._tokens = math.inf # starts full: clamped to the rate on first use self._refilled = clock() @@ -101,6 +103,13 @@ class PluginDmrdSender: self._call_from_reactor(self._deliver, bytes(pkt), self._plugin) return True + def voice_slot_for_tg(self, tg: int) -> int | None: + """``ServerContext.voice_slot_for_tg``: only for granted talkgroups; reactor thread.""" + permission = send_permission(self._config, self._plugin) + if permission is None or int(tg) not in permission.group_voice_tgs or self._slot_for_tg is None: + return None + return self._slot_for_tg(int(tg)) + def _take_token(self, permission: SendPermission) -> bool: with self._lock: now = self._clock() diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 5d8e591..b170fa9 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -455,14 +455,17 @@ def run_peer_server( if callable(send_system): send_system(pkt) - plugin_ingress = PluginIngress(None, config, lambda: protocols, _send_local) # routing set below + plugin_ingress = PluginIngress( # routing set below + None, config, lambda: protocols, _send_local, send_routing_event=reporting_use_cases.send_routing_event, + ) plugin_manager = PluginManager( plugin_bus, server_ctx, project_root, sender_factory=lambda name: PluginDmrdSender( - name, config, plugin_ingress.deliver, reactor.callFromThread, in_reactor_thread=isInIOThread, + name, config, plugin_ingress.deliver, reactor.callFromThread, + in_reactor_thread=isInIOThread, slot_for_tg=plugin_ingress.voice_slot_for_tg, ), ) voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks) diff --git a/tests/application/test_plugin_sender.py b/tests/application/test_plugin_sender.py index 18888fc..011c544 100644 --- a/tests/application/test_plugin_sender.py +++ b/tests/application/test_plugin_sender.py @@ -250,3 +250,12 @@ def test_a_rate_that_is_not_a_positive_number_grants_nothing(rate) -> None: from adn_server.application.plugins.domain.send import send_permission assert send_permission(_config(max_frames_per_s=rate), "d-aprs") is None + + +def test_voice_slot_query_only_for_granted_talkgroups() -> None: + sender = PluginDmrdSender( + "d-aprs", _config(group_voice_tgs=[213]), lambda *a: True, call_from_reactor=lambda *a: None, + slot_for_tg=lambda tg: 2, + ) + assert sender.voice_slot_for_tg(213) == 2 + assert sender.voice_slot_for_tg(214) is None diff --git a/tests/routing/test_plugin_send_routing.py b/tests/routing/test_plugin_send_routing.py index bc9b9cb..e751acf 100644 --- a/tests/routing/test_plugin_send_routing.py +++ b/tests/routing/test_plugin_send_routing.py @@ -187,3 +187,51 @@ def test_a_second_plugin_stream_cannot_talk_over_the_first() -> None: sc, ingress, _ = _voice_scenario() _play(sc, ingress, _beacon(stream=1)[:2]) assert _play(sc, ingress, _beacon(stream=2)[:1]) == [False] + + +def test_voice_slot_prefers_the_slot_where_the_tg_is_bridged_on_the_master() -> None: + config = minimal_config(("SYSTEM", "SYSTEM-B")) + config["SYSTEMS"]["SYSTEM"]["PEERS"] = {b"\x00\x00\x03\xe9": {"CALLSIGN": "HOTSPOT"}} + table = active_routing_table(TG, (("SYSTEM", 1), ("SYSTEM-B", 2)), timeout_minutes=10**6) + sc = DeterministicScenario(config=config, routing_table=table) + sc.routing.apply_startup_subscriptions() + for ts in (1, 2): + sc.protocols["SYSTEM"].STATUS[ts] = idle_hbp_slot() | {"RX_TYPE": HBPF_SLT_VTERM} + ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None, clock=sc.clock.time) + assert ingress.voice_slot_for_tg(TG) == 1 # TG 213 lives on TS1 of this MASTER + assert ingress.voice_slot_for_tg(9999) == 2 # nowhere yet: TS2 first + + +def test_voice_slot_is_none_while_every_slot_is_busy() -> None: + sc, ingress, _ = _voice_scenario() + for ts in (1, 2): + sc.protocols["SYSTEM"].STATUS[ts] = { + "RX_TYPE": HBPF_SLT_VHEAD, "RX_TIME": sc.clock.time(), "TX_TYPE": HBPF_SLT_VTERM, + "RX_STREAM_ID": b"\x05\x05\x05\x05", + } + assert ingress.voice_slot_for_tg(TG) is None + + +def test_the_monitor_sees_the_beacon_on_the_master_it_plays_on() -> None: + sc, _, _ = _voice_scenario() + events: list[str] = [] + ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None, + clock=sc.clock.time, send_routing_event=events.append) + _play(sc, ingress, _beacon(stream=0x0C0C0C0C)) + starts = [e for e in events if e.startswith("GROUP VOICE,START,TX,SYSTEM,")] + ends = [e for e in events if e.startswith("GROUP VOICE,END,TX,SYSTEM,")] + assert len(starts) == 1 and len(ends) == 1 + assert starts[0].split(",")[4:9] == [str(0x0C0C0C0C), str(BEACON_ID), str(BEACON_ID), "2", str(TG)] + + +def test_a_beacon_cut_by_a_radio_still_ends_on_the_monitor() -> None: + sc, _, _ = _voice_scenario() + events: list[str] = [] + ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None, + clock=sc.clock.time, send_routing_event=events.append) + frames = _beacon() + _play(sc, ingress, frames[:3]) + sc.protocols["SYSTEM"].STATUS[2].update(RX_TYPE=HBPF_SLT_VHEAD, RX_TIME=sc.clock.time() + 0.06, RX_STREAM_ID=b"\x09" * 4) + assert _play(sc, ingress, frames[3:4]) == [False] + assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1 + assert sc.protocols["SYSTEM"].STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM