diff --git a/docs/en/server/user-guide/plugins.md b/docs/en/server/user-guide/plugins.md index d9e0a24..e9ca5f1 100644 --- a/docs/en/server/user-guide/plugins.md +++ b/docs/en/server/user-guide/plugins.md @@ -45,7 +45,7 @@ PLUGINS: | `directory` | Absolute path or relative to project root | | `master_kill` | Emergency disable — no plugins loaded | | `overrides` | Per-plugin config patches without editing `plugins//config.yaml` | -| `send` | Per-plugin permission to send unit data — see [Sending unit data](#sending-unit-data-opt-in) | +| `send` | Per-plugin permission to send unit data and group voice — see [Sending](#sending-unit-data-and-group-voice-opt-in) | On **SIGHUP**, `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed. @@ -94,15 +94,15 @@ Passed to `on_load` as `server_ctx` (`application/plugins/application/context.py | `defer_to_thread(fn, *args)` | Run blocking I/O off the reactor | | `call_from_reactor(fn, *args)` | Schedule callback on reactor thread | | `call_later(delay_s, fn, *args)` | Reactor timer | -| `send_dmrd(pkt) -> bool` | Send one DMRD frame — **only** for a plugin granted it in `PLUGINS.send`, otherwise `None` ([Sending](#sending-unit-data-opt-in)) | +| `send_dmrd(pkt) -> bool` | Send one DMRD frame — **only** for a plugin granted it in `PLUGINS.send`, otherwise `None` ([Sending](#sending-unit-data-and-group-voice-opt-in)) | **Pattern:** do only dispatch in `on_event`; call `defer_to_thread` for file writes, HTTP, heavy CPU. --- -## Sending unit data (opt-in) +## Sending unit data and group voice (opt-in) -A plugin can send **unit data** (data header, rate 1/2 and 3/4 blocks, CSBK: ARS, LRRP, SMS…) with `server_ctx.send_dmrd(pkt)`, one complete HBP `DMRD` frame per call. Enabling a plugin never lets it transmit: the sysop grants it per plugin in `adn-server.yaml`: +A plugin can send **unit data** (data header, rate 1/2 and 3/4 blocks, CSBK: ARS, LRRP, SMS…) and **group voice** on granted talkgroups (voice beacons, announcements) with `server_ctx.send_dmrd(pkt)`, one complete HBP `DMRD` frame per call. Enabling a plugin never lets it transmit: the sysop grants it per plugin in `adn-server.yaml`: ```yaml PLUGINS: @@ -110,12 +110,15 @@ PLUGINS: d-aprs: allowed_src_ids: [900999] # rf_src the plugin may send as; required max_frames_per_s: 40 # per-plugin token bucket (default 40) + beacon: + allowed_src_ids: [2130035] + group_voice_tgs: [213] # group voice only on these TGs; none: unit data only ``` - `send_dmrd` is `None` unless the plugin has an entry with at least one source ID. - Each frame is checked against the **current** config: removing the entry (SIGHUP) or `master_kill` stops sending at once. Granting it to a plugin already loaded needs that plugin reloaded. - Rejected frames (not unit data, source not allowed, over the rate) return `False` and are logged and counted; the first frame of each stream is logged at INFO. -- Safe from any thread: frames are handed to the reactor. +- Called on the reactor thread (`on_event`, `call_later`), the frame is routed at once and the result is whether the server **accepted** it; a plugin sending voice must stop when it gets `False`. From another thread the frame is queued to the reactor and `True` only means the guards passed. How the server treats them: @@ -127,6 +130,8 @@ How the server treats them: | OpenBridge / `DATA-GATEWAY` | No fan-out: plugin frames stay on this server | | Events | Reach plugins with `is_synthetic=True`, so a plugin can ignore its own frames | +Group voice is routed **like a scheduled announcement** (synthetic PTT on the same MASTER): through the TG's bridges, OpenBridge legs included, and to the hotspots of that MASTER. While a plugin stream plays it holds that MASTER slot, so routed calls find it busy; a radio or another stream on the slot makes the next frame fail. The terminator frees the slot. Private voice can't be sent. + Pacing (about 60 ms per burst) is the plugin's job, e.g. with `call_later`. --- diff --git a/docs/es/server/user-guide/plugins.md b/docs/es/server/user-guide/plugins.md index e0c7be7..ddae1a5 100644 --- a/docs/es/server/user-guide/plugins.md +++ b/docs/es/server/user-guide/plugins.md @@ -45,7 +45,7 @@ PLUGINS: | `directory` | Ruta absoluta o relativa al project root | | `master_kill` | Desactivación de emergencia — no carga plugins | | `overrides` | Parches por plugin sin editar `plugins//config.yaml` | -| `send` | Permiso por plugin para enviar datos — ver [Envío de datos](#envío-de-datos-opcional) | +| `send` | Permiso por plugin para enviar datos y voz de grupo — ver [Envío](#envío-de-datos-y-voz-de-grupo-opcional) | Con **SIGHUP**, `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`. @@ -94,15 +94,15 @@ Se pasa a `on_load` como `server_ctx` (`application/plugins/application/context. | `defer_to_thread(fn, *args)` | Ejecutar I/O bloqueante fuera del reactor | | `call_from_reactor(fn, *args)` | Programar callback en el hilo del reactor | | `call_later(delay_s, fn, *args)` | Temporizador del reactor | -| `send_dmrd(pkt) -> bool` | Enviar una trama DMRD — **solo** para un plugin autorizado en `PLUGINS.send`; si no, `None` ([Envío](#envío-de-datos-opcional)) | +| `send_dmrd(pkt) -> bool` | Enviar una trama DMRD — **solo** para un plugin autorizado en `PLUGINS.send`; si no, `None` ([Envío](#envío-de-datos-y-voz-de-grupo-opcional)) | **Patrón:** solo despachar en `on_event`; usar `defer_to_thread` para archivos, HTTP o CPU intensiva. --- -## Envío de datos (opcional) +## Envío de datos y voz de grupo (opcional) -Un plugin puede enviar **datos** (cabecera, bloques de 1/2 y 3/4, CSBK: ARS, LRRP, SMS…) con `server_ctx.send_dmrd(pkt)`, una trama HBP `DMRD` completa por llamada. Activar un plugin nunca le permite transmitir: el sysop lo autoriza por plugin en `adn-server.yaml`: +Un plugin puede enviar **datos** (cabecera, bloques de 1/2 y 3/4, CSBK: ARS, LRRP, SMS…) y **voz de grupo** en los TGs autorizados (balizas de voz, anuncios) con `server_ctx.send_dmrd(pkt)`, una trama HBP `DMRD` completa por llamada. Activar un plugin nunca le permite transmitir: el sysop lo autoriza por plugin en `adn-server.yaml`: ```yaml PLUGINS: @@ -110,12 +110,15 @@ PLUGINS: d-aprs: allowed_src_ids: [900999] # rf_src con el que puede enviar; obligatorio max_frames_per_s: 40 # cubo de tokens por plugin (por defecto 40) + baliza: + allowed_src_ids: [2130035] + group_voice_tgs: [213] # voz de grupo solo en estos TGs; sin ella, solo datos ``` - `send_dmrd` es `None` salvo que el plugin tenga una entrada con al menos un ID de origen. - Cada trama se comprueba contra la configuración **actual**: quitar la entrada (SIGHUP) o `master_kill` corta el envío al instante. Autorizar a un plugin ya cargado exige recargar ese plugin. - Las tramas rechazadas (no son datos, origen no permitido, exceso de ritmo) devuelven `False` y se registran y cuentan; la primera trama de cada stream se registra en INFO. -- Se puede llamar desde cualquier hilo: las tramas se entregan al reactor. +- Llamado desde el hilo del reactor (`on_event`, `call_later`), la trama se enruta en el acto y el resultado indica si el servidor la **aceptó**; un plugin que emite voz debe parar cuando recibe `False`. Desde otro hilo la trama se encola al reactor y `True` solo significa que pasó las salvaguardas. Cómo las trata el servidor: @@ -127,6 +130,8 @@ Cómo las trata el servidor: | OpenBridge / `DATA-GATEWAY` | Sin reparto: las tramas del plugin se quedan en este servidor | | Eventos | Llegan a los plugins con `is_synthetic=True`, para que un plugin ignore sus propias tramas | +La voz de grupo se enruta **como un anuncio programado** (PTT sintético en el mismo MASTER): por los puentes del TG, OpenBridge incluidos, y a los hotspots de ese MASTER. Mientras suena, el stream del plugin ocupa ese slot del MASTER, así que las llamadas enrutadas lo encuentran ocupado; una radio u otro stream en el slot hace fallar la trama siguiente. El terminador libera el slot. No se puede enviar voz privada. + El ritmo (unos 60 ms por ráfaga) lo marca el plugin, por ejemplo con `call_later`. --- diff --git a/src/adn_server/application/plugins/application/ingress.py b/src/adn_server/application/plugins/application/ingress.py new file mode 100644 index 0000000..fc89361 --- /dev/null +++ b/src/adn_server/application/plugins/application/ingress.py @@ -0,0 +1,112 @@ +# ADN DMR Peer Server - plugin frame ingress +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Where a plugin's frames enter the server: the announcement MASTER, as a synthetic ingress.""" + +from __future__ import annotations + +import logging +import time +from typing import Any, Callable + +from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_4, int_id +from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd +from ...routing.helpers import slot_voice_held_by_other_stream +from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind + +logger = logging.getLogger(__name__) + + +class PluginIngress: + """Routes plugin frames on the reactor thread and reports whether each was accepted. + + Unit data goes through the unit data path only (see ``dmrd_received(plugin_origin=)``). + + Group voice is routed like a scheduled announcement: through the bridges (OpenBridge + legs included, as for any local ingress), and out to the hotspots of the MASTER it + enters on. While a plugin's stream plays it holds that MASTER slot (TX_TYPE=VHEAD, + TX_STREAM_ID, TX_RFS), so routed voice finds it busy; a radio or another stream on the + slot makes the frame fail, which tells the plugin to stop. The terminator frees it. + """ + + def __init__( + self, + routing: Any, + config: dict[str, Any], + get_protocols: Callable[[], dict[str, Any]], + send_local: Callable[[str, bytes], None], + clock: Callable[[], float] = time.time, + ) -> None: + self._routing = routing + self._config = config + self._get_protocols = get_protocols + self._send_local = send_local + self._clock = clock + + def set_routing(self, routing: Any) -> None: + """Bootstrap builds the plugin manager before routing; routing is set before any load.""" + self._routing = routing + + 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 + master = announcement_ptt_system(self._config) + if kind is None or not master: + if not master: + logger.warning("(PLUGIN) %s: no MASTER to send from, frame dropped", plugin) + return False + server_id = self._server_id() + now = self._clock() + if kind == UNIT_DATA: + return self._route(master, pkt, now, server_id, plugin) is not False + return self._group_voice(master, header, pkt[:11] + server_id + pkt[15:], now, server_id, plugin) + + def _group_voice(self, master: str, header: Any, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool: + proto = self._get_protocols().get(master) + 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: + self._release(slot, header.stream_id) + 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 + slot["TX_STREAM_ID"] = header.stream_id + slot["TX_RFS"] = header.rf_src + slot["TX_TGID"] = header.dst_id + slot["TX_TIME"] = now + self._send_local(master, pkt) + return True + + 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) + + @staticmethod + def _release(slot: dict[str, Any], stream_id: bytes) -> None: + if slot.get("TX_STREAM_ID") == stream_id: + slot["TX_TYPE"] = HBPF_SLT_VTERM + + def _server_id(self) -> bytes: + server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") + if isinstance(server_id, bytes) and len(server_id) >= 4: + return server_id[:4] + return bytes_4(int(int_id(server_id) if isinstance(server_id, bytes) else server_id or 0) & 0xFFFFFFFF) diff --git a/src/adn_server/application/plugins/application/sender.py b/src/adn_server/application/plugins/application/sender.py index f6217ba..6d7e715 100644 --- a/src/adn_server/application/plugins/application/sender.py +++ b/src/adn_server/application/plugins/application/sender.py @@ -29,7 +29,7 @@ import time from typing import Any, Callable from ....domain import int_id -from ..domain.send import SendPermission, is_plugin_sendable, parse_dmrd_header, send_permission +from ..domain.send import GROUP_VOICE, SendPermission, parse_dmrd_header, plugin_frame_kind, send_permission logger = logging.getLogger(__name__) @@ -38,8 +38,12 @@ class PluginDmrdSender: """Checks each frame against the plugin's live permission, then hands it to the reactor. The permission is read from the server config on every call, so a SIGHUP that - removes or narrows ``PLUGINS.send.`` takes effect at once. Safe to call - from any thread; delivery always happens on the reactor. + removes or narrows ``PLUGINS.send.`` takes effect at once. + + Called on the reactor thread (``on_event``, ``call_later``), the frame is routed at + once and the result is whether the server accepted it: a plugin sending voice stops + when it gets False. From another thread the frame is queued to the reactor and the + result only says the guards passed. """ def __init__( @@ -49,12 +53,14 @@ class PluginDmrdSender: deliver: Callable[[bytes, str], Any], call_from_reactor: Callable[..., Any], clock: Callable[[], float] = time.monotonic, + in_reactor_thread: Callable[[], bool] = lambda: False, ) -> None: self._plugin = plugin self._config = server_config self._deliver = deliver self._call_from_reactor = call_from_reactor self._clock = clock + self._in_reactor_thread = in_reactor_thread self._lock = threading.Lock() self._tokens = math.inf # starts full: clamped to the rate on first use self._refilled = clock() @@ -69,8 +75,11 @@ class PluginDmrdSender: header = parse_dmrd_header(bytes(pkt)) if isinstance(pkt, (bytes, bytearray)) else None if header is None: return self._reject("not a DMRD frame") - if not is_plugin_sendable(header.call_type, header.frame_type, header.dtype_vseq): - return self._reject("only unit data may be sent") + kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq) + if kind is None: + return self._reject("only unit data and group voice may be sent") + if kind == GROUP_VOICE and int_id(header.dst_id) not in permission.group_voice_tgs: + return self._reject(f"TG {int_id(header.dst_id)} not in group_voice_tgs") rf_src = int_id(header.rf_src) stream_id, dst_id = header.stream_id, header.dst_id if rf_src not in permission.allowed_src_ids: @@ -87,6 +96,8 @@ class PluginDmrdSender: logger.info( "(PLUGIN) %s sent stream %s src %s -> dst %s", self._plugin, int_id(stream_id), rf_src, int_id(dst_id) ) + if self._in_reactor_thread(): + return bool(self._deliver(bytes(pkt), self._plugin)) self._call_from_reactor(self._deliver, bytes(pkt), self._plugin) return True diff --git a/src/adn_server/application/plugins/domain/send.py b/src/adn_server/application/plugins/domain/send.py index be3362a..066c4fe 100644 --- a/src/adn_server/application/plugins/domain/send.py +++ b/src/adn_server/application/plugins/domain/send.py @@ -25,10 +25,12 @@ from __future__ import annotations from dataclasses import dataclass from typing import Any, NamedTuple -from ....domain import HBPF_DATA_SYNC +from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, HBPF_VOICE_SYNC # CSBK, data header, rate 1/2 and rate 3/4 data blocks: ARS, LRRP, SMS and the like. PLUGIN_SENDABLE_DTYPES = frozenset({3, 6, 7, 8}) +UNIT_DATA = "unit_data" +GROUP_VOICE = "group_voice" DEFAULT_MAX_FRAMES_PER_S = 40.0 @@ -68,15 +70,28 @@ def parse_dmrd_header(pkt: bytes) -> DmrdHeader | None: ) +def plugin_frame_kind(call_type: str, frame_type: int, dtype_vseq: int) -> str | None: + """UNIT_DATA, GROUP_VOICE, or None for anything a plugin may not send (e.g. private voice).""" + if call_type == "unit": + return UNIT_DATA if frame_type == HBPF_DATA_SYNC and dtype_vseq in PLUGIN_SENDABLE_DTYPES else None + if call_type == "group": + if frame_type in (HBPF_VOICE, HBPF_VOICE_SYNC): + return GROUP_VOICE + if frame_type == HBPF_DATA_SYNC and dtype_vseq in (HBPF_SLT_VHEAD, HBPF_SLT_VTERM): + return GROUP_VOICE + return None + + def is_plugin_sendable(call_type: str, frame_type: int, dtype_vseq: int) -> bool: - """Plugins send unit data only (no voice) in this version.""" - return call_type == "unit" and frame_type == HBPF_DATA_SYNC and dtype_vseq in PLUGIN_SENDABLE_DTYPES + return plugin_frame_kind(call_type, frame_type, dtype_vseq) is not None @dataclass(frozen=True) class SendPermission: allowed_src_ids: frozenset[int] max_frames_per_s: float + # Talkgroups the plugin may speak on (group voice); empty: unit data only. + group_voice_tgs: frozenset[int] = frozenset() def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermission | None: @@ -94,8 +109,9 @@ def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermissio try: ids = frozenset(int(i) for i in entry.get("allowed_src_ids") or ()) rate = float(entry.get("max_frames_per_s", DEFAULT_MAX_FRAMES_PER_S)) + tgs = frozenset(int(t) for t in entry.get("group_voice_tgs") or ()) except (TypeError, ValueError): return None if not ids or rate <= 0: return None - return SendPermission(ids, rate) + return SendPermission(ids, rate, tgs) diff --git a/src/adn_server/application/routing/announcement_ptt_inject.py b/src/adn_server/application/routing/announcement_ptt_inject.py index bd83e1f..cb02518 100644 --- a/src/adn_server/application/routing/announcement_ptt_inject.py +++ b/src/adn_server/application/routing/announcement_ptt_inject.py @@ -113,6 +113,7 @@ def inject_plugin_dmrd( """Feed one plugin frame through ``dmrd_received`` on the same MASTER as announcements. The peer field is always the SERVER_ID: a plugin never speaks as a connected peer. + Group voice is marked ``synthetic_announcement`` exactly like announcements are. """ header = parse_dmrd_header(pkt) if header is None: @@ -130,5 +131,6 @@ def inject_plugin_dmrd( header.stream_id, pkt, ingress_pkt_time=pkt_time, + synthetic_announcement=header.call_type == "group", plugin_origin=plugin, ) diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 7e98c2b..f5590e6 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -222,10 +222,12 @@ class RoutingUseCases( dynamic TG) to that peer. ``plugin_origin`` names the plugin that sent the frame through - ``ServerContext.send_dmrd``. Such frames are unit data only, are delivered by - the unit data path alone (SUB_MAP / hotspot ID, never the private call path, - so SUB_MAP never learns their source), stay on this server (no OpenBridge or - DATA-GATEWAY fan-out), and reach plugins as synthetic events. + ``ServerContext.send_dmrd``. Unit data is delivered by the unit data path + alone (SUB_MAP / hotspot ID, never the private call path, so SUB_MAP never + learns its source), stays on this server (no OpenBridge or DATA-GATEWAY + fan-out), and reaches plugins as synthetic events. Group voice comes with + ``synthetic_announcement=True`` and is routed exactly like an announcement. + Anything else (private voice) is dropped. """ if not self._send_to_system: return @@ -240,7 +242,7 @@ class RoutingUseCases( if call_type == "unit" and int_id(dst_id) == 4000: return if plugin_origin is not None and not is_plugin_sendable(call_type, frame_type, dtype_vseq): - logger.warning("(PLUGIN) %s: only unit data may be sent, dropped", plugin_origin) + logger.warning("(PLUGIN) %s: only unit data and group voice may be sent, dropped", plugin_origin) return False if call_type == "unit": _int_dst = int_id(dst_id) diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 6d93bbc..5d8e591 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -29,6 +29,7 @@ import time from typing import Any from twisted.internet import reactor, task, threads +from twisted.python.threadable import isInIOThread from adn_server.application import ( IdentUseCases, @@ -41,6 +42,7 @@ from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases from adn_server.application.plugins.application.bus import PluginBus from adn_server.application.plugins.application.context import ServerContext from adn_server.application.plugins.application.data_bridge import DataPluginBridge +from adn_server.application.plugins.application.ingress import PluginIngress from adn_server.application.plugins.application.manager import PluginManager from adn_server.application.plugins.application.sender import PluginDmrdSender from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge @@ -448,27 +450,20 @@ def run_peer_server( call_from_reactor=reactor.callFromThread, call_later=reactor.callLater, ) - def _deliver_plugin_dmrd(pkt: bytes, plugin: str) -> None: - # Reactor thread (PluginDmrdSender schedules it there). Same ingress MASTER as announcements. - from adn_server.application.routing.announcement_ptt_inject import ( - announcement_ptt_system, - inject_plugin_dmrd, - ) + def _send_local(system: str, pkt: bytes) -> None: + send_system = getattr(protocols.get(system), "send_system", None) + if callable(send_system): + send_system(pkt) - master = announcement_ptt_system(config) - if not master: - logger.warning("(PLUGIN) %s: no MASTER to send from, frame dropped", plugin) - return - server_id = config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") - if not isinstance(server_id, bytes): - server_id = bytes_4(int(server_id or 0) & 0xFFFFFFFF) - inject_plugin_dmrd(routing_use_cases, master, pkt, pkt_time=time.time(), server_id=server_id, plugin=plugin) + plugin_ingress = PluginIngress(None, config, lambda: protocols, _send_local) # routing set below plugin_manager = PluginManager( plugin_bus, server_ctx, project_root, - sender_factory=lambda name: PluginDmrdSender(name, config, _deliver_plugin_dmrd, reactor.callFromThread), + sender_factory=lambda name: PluginDmrdSender( + name, config, plugin_ingress.deliver, reactor.callFromThread, in_reactor_thread=isInIOThread, + ), ) voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks) data_plugin_bridge = DataPluginBridge(plugin_bus, config) @@ -490,6 +485,7 @@ def run_peer_server( voice_plugin_bridge=voice_plugin_bridge, data_plugin_bridge=data_plugin_bridge, ) + plugin_ingress.set_routing(routing_use_cases) plugin_manager.discover_and_load(config) reactor.addSystemEventTrigger("before", "shutdown", plugin_manager.shutdown_all) routing_use_cases.apply_startup_subscriptions() diff --git a/tests/application/test_plugin_sender.py b/tests/application/test_plugin_sender.py index 6d1a1c3..51122ab 100644 --- a/tests/application/test_plugin_sender.py +++ b/tests/application/test_plugin_sender.py @@ -81,14 +81,40 @@ def test_a_plugin_cannot_send_as_a_radio() -> None: assert sender(_frame(src=7140023)) is False and delivered == [] -def test_voice_and_group_frames_are_refused() -> None: - sender, delivered = _sender(_config()) - assert sender(_frame(call_type="group")) is False - assert sender(_frame(frame_type=HBPF_VOICE, dtype=1)) is False +def test_private_voice_and_non_dmrd_are_refused() -> None: + sender, delivered = _sender(_config(group_voice_tgs=[213])) + assert sender(_frame(frame_type=HBPF_VOICE, dtype=1)) is False # unit call type, voice frame assert sender(b"not a dmrd frame") is False assert delivered == [] +def test_group_voice_only_on_the_granted_talkgroups() -> None: + sender, delivered = _sender(_config(group_voice_tgs=[213])) + voice = PacketSpec(rf_src=GATEWAY_ID, dst_id=213, call_type="group", frame_type=HBPF_VOICE, dtype_vseq=1).data() + other = PacketSpec(rf_src=GATEWAY_ID, dst_id=214, call_type="group", frame_type=HBPF_VOICE, dtype_vseq=1).data() + assert sender(voice) is True + assert sender(other) is False + assert [pkt for pkt, _ in delivered] == [voice] + + +def test_group_voice_needs_a_grant_even_with_a_source_id() -> None: + sender, delivered = _sender(_config()) + voice = PacketSpec(rf_src=GATEWAY_ID, dst_id=213, call_type="group", frame_type=HBPF_VOICE, dtype_vseq=1).data() + assert sender(voice) is False and delivered == [] + + +def test_on_the_reactor_thread_the_result_is_the_routing_result() -> None: + results = iter([True, False]) + queued: list = [] + sender = PluginDmrdSender( + "d-aprs", _config(), lambda pkt, name: next(results), + call_from_reactor=lambda *a: queued.append(a), in_reactor_thread=lambda: True, + ) + assert sender(_frame()) is True + assert sender(_frame()) is False # e.g. the slot was taken: the plugin should stop + assert queued == [] + + def test_rate_limit_per_plugin() -> None: clock = _Clock() sender, delivered = _sender(_config(max_frames_per_s=5), clock) diff --git a/tests/routing/test_plugin_send_routing.py b/tests/routing/test_plugin_send_routing.py index 8e6a04d..bc9b9cb 100644 --- a/tests/routing/test_plugin_send_routing.py +++ b/tests/routing/test_plugin_send_routing.py @@ -25,6 +25,7 @@ from __future__ import annotations from tests.harness.deterministic import ( DeterministicScenario, PacketSpec, + active_routing_table, add_openbridge_system, minimal_config, patch_routing_wall_time, @@ -33,9 +34,10 @@ from tests.routing.unit_data_helpers import idle_hbp_slot from adn_server.application.plugins.application.bus import PluginBus from adn_server.application.plugins.application.data_bridge import DataPluginBridge +from adn_server.application.plugins.application.ingress import PluginIngress from adn_server.application.plugins.domain.events import UnitDataFrame from adn_server.application.routing.announcement_ptt_inject import inject_plugin_dmrd -from adn_server.domain import HBPF_DATA_SYNC, bytes_3, bytes_4 +from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_3, bytes_4 GATEWAY_ID = 900999 RADIO_ID = 7140023 @@ -98,10 +100,10 @@ def test_the_plugin_source_is_never_learned_in_sub_map() -> None: assert bytes_3(GATEWAY_ID) not in sc.config["_SUB_MAP"] -def test_group_traffic_from_a_plugin_is_refused() -> None: +def test_private_voice_from_a_plugin_is_refused() -> None: sc = _scenario() - assert _send(sc, _ars_ack(dst=213, call_type="group", frame_type=0)) is False - assert sc.capture.packets == [] + assert _send(sc, _ars_ack(frame_type=HBPF_VOICE)) is False + assert sc.capture.packets == [] and sc.protocols["SYSTEM-B"].sent_to_peer == [] def test_plugins_see_their_own_frames_as_synthetic() -> None: @@ -113,3 +115,75 @@ def test_plugins_see_their_own_frames_as_synthetic() -> None: _send(sc, _ars_ack()) frames = [e for e in events if isinstance(e, UnitDataFrame)] assert frames and all(f.context.is_synthetic for f in frames) + + +# --- group voice (voice beacons and announcements as plugins) --- + +TG = 213 +BEACON_ID = 2130035 + + +def _voice_scenario(): + config = minimal_config(("SYSTEM", "SYSTEM-B")) + config["SYSTEMS"]["SYSTEM"]["PEERS"] = {b"\x00\x00\x03\xe9": {"CALLSIGN": "HOTSPOT"}} + table = active_routing_table(TG, (("SYSTEM", 2), ("SYSTEM-B", 2)), timeout_minutes=10**6) + sc = DeterministicScenario(config=config, routing_table=table) + sc.routing.apply_startup_subscriptions() + for name in ("SYSTEM", "SYSTEM-B"): + sc.protocols[name].STATUS[2] = idle_hbp_slot() + local: list[bytes] = [] + ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda system, pkt: local.append(pkt), clock=sc.clock.time) + return sc, ingress, local + + +def _beacon(stream: int = 0x0B0B0B0B) -> list[bytes]: + base = PacketSpec(rf_src=BEACON_ID, dst_id=TG, peer_id=7, slot=2, stream_id=stream) + frames = [DeterministicScenario.voice_head_spec(base)] + frames += [DeterministicScenario.voice_burst_spec(base, seq=n + 1, dtype_vseq=n % 6) for n in range(6)] + frames.append(DeterministicScenario.voice_term_spec(base, seq=7)) + return [f.data() for f in frames] + + +def _play(sc, ingress, frames) -> list[bool]: + out = [] + for pkt in frames: + sc.clock.advance(0.06) + with patch_routing_wall_time(sc.clock): + out.append(ingress.deliver(pkt, "beacon")) + return out + + +def test_a_voice_beacon_is_routed_like_an_announcement_and_heard_locally() -> None: + sc, ingress, local = _voice_scenario() + frames = _beacon() + assert _play(sc, ingress, frames) == [True] * len(frames) + to_b = sc.capture.for_system("SYSTEM-B") + assert len(to_b) == len(frames) + assert all(p.packet[5:8] == bytes_3(BEACON_ID) for p in to_b) + assert len(local) == len(frames) # hotspots of the MASTER it enters on + assert {p[11:15] for p in local} == {ingress._server_id()} # never the peer the plugin wrote + assert ingress._server_id() != bytes_4(7) + + +def test_the_master_slot_is_held_while_it_plays_and_freed_by_the_terminator() -> None: + sc, ingress, _ = _voice_scenario() + frames = _beacon() + _play(sc, ingress, frames[:3]) + slot = sc.protocols["SYSTEM"].STATUS[2] + assert slot["TX_TYPE"] == HBPF_SLT_VHEAD and slot["TX_RFS"] == bytes_3(BEACON_ID) + _play(sc, ingress, frames[3:]) + assert slot["TX_TYPE"] == HBPF_SLT_VTERM + + +def test_a_radio_talking_on_the_slot_stops_the_beacon() -> None: + sc, ingress, local = _voice_scenario() + slot = sc.protocols["SYSTEM"].STATUS[2] + slot.update(RX_TYPE=HBPF_SLT_VHEAD, RX_TIME=sc.clock.time() + 0.06, RX_STREAM_ID=b"\x01\x01\x01\x01") + assert _play(sc, ingress, _beacon()[:1]) == [False] + assert sc.capture.for_system("SYSTEM-B") == [] and local == [] + + +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]