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