From 0c3432ebfd31b2ae8e702f3cc1e0ce868f837ad0 Mon Sep 17 00:00:00 2001 From: yo Date: Thu, 24 Sep 2026 23:20:07 +0200 Subject: [PATCH 1/3] feat(plugins): let a plugin send unit data, opt-in per plugin (#103) ServerContext.send_dmrd(pkt) hands one DMRD frame to the routing core, through the same synthetic ingress path scheduled announcements use (inject_plugin_dmrd -> dmrd_received on the announcement MASTER, with the SERVER_ID as peer). Only plugins listed in PLUGINS.send get it. Guards (PluginDmrdSender, re-read from the live config on every frame): - allowed_src_ids: a plugin can't send as a radio; required. - max_frames_per_s: per-plugin token bucket, starts full. - unit data only in this version (data header, rate 1/2, 3/4, CSBK). - master_kill or removing the entry stops sending at once. Routing of plugin frames (dmrd_received plugin_origin): - delivered by the unit data path only: SUB_MAP / hotspot peer ID, to the destination's exact hotspot, also on the ingress MASTER itself; - never through the private call path, which would learn the plugin's source in SUB_MAP (spreading replies over every hotspot of that MASTER) and keep call state on its shared slot; - no OpenBridge or DATA-GATEWAY fan-out; - plugin events for them carry is_synthetic=True. Co-Authored-By: Claude Opus 5.5 --- docs/en/server/user-guide/plugins.md | 33 ++++ docs/es/server/user-guide/plugins.md | 33 ++++ .../plugins/application/context.py | 2 + .../plugins/application/data_bridge.py | 13 +- .../plugins/application/manager.py | 21 ++- .../application/plugins/application/sender.py | 110 +++++++++++++ .../application/plugins/domain/send.py | 101 ++++++++++++ .../routing/announcement_ptt_inject.py | 34 +++++ .../application/routing_use_cases.py | 24 ++- .../infrastructure/bootstrap/peer_server.py | 24 ++- tests/application/test_plugin_sender.py | 144 ++++++++++++++++++ tests/routing/test_plugin_send_routing.py | 115 ++++++++++++++ 12 files changed, 641 insertions(+), 13 deletions(-) create mode 100644 src/adn_server/application/plugins/application/sender.py create mode 100644 src/adn_server/application/plugins/domain/send.py create mode 100644 tests/application/test_plugin_sender.py create mode 100644 tests/routing/test_plugin_send_routing.py diff --git a/docs/en/server/user-guide/plugins.md b/docs/en/server/user-guide/plugins.md index 6e53362..d9e0a24 100644 --- a/docs/en/server/user-guide/plugins.md +++ b/docs/en/server/user-guide/plugins.md @@ -45,6 +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) | On **SIGHUP**, `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed. @@ -93,11 +94,43 @@ 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)) | **Pattern:** do only dispatch in `on_event`; call `defer_to_thread` for file writes, HTTP, heavy CPU. --- +## Sending unit data (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`: + +```yaml +PLUGINS: + send: + 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) +``` + +- `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. + +How the server treats them: + +| | | +|---|---| +| Ingress | The same MASTER as scheduled announcements, with the SERVER_ID as peer | +| Delivery | Unit data path only: `SUB_MAP` / hotspot peer ID, to the exact hotspot of the destination — also on the ingress MASTER itself | +| `SUB_MAP` | Never learns the plugin's source ID (so replies to it are not spread over the MASTER's hotspots) | +| 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 | + +Pacing (about 60 ms per burst) is the plugin's job, e.g. with `call_later`. + +--- + ## Event bus `PluginBus` (`application/plugins/application/bus.py`) delivers events to every loaded plugin's `on_event`. diff --git a/docs/es/server/user-guide/plugins.md b/docs/es/server/user-guide/plugins.md index 83b575d..e0c7be7 100644 --- a/docs/es/server/user-guide/plugins.md +++ b/docs/es/server/user-guide/plugins.md @@ -45,6 +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) | Con **SIGHUP**, `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`. @@ -93,11 +94,43 @@ 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)) | **Patrón:** solo despachar en `on_event`; usar `defer_to_thread` para archivos, HTTP o CPU intensiva. --- +## Envío de datos (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`: + +```yaml +PLUGINS: + send: + 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) +``` + +- `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. + +Cómo las trata el servidor: + +| | | +|---|---| +| Entrada | El mismo MASTER que los anuncios programados, con el SERVER_ID como peer | +| Entrega | Solo el camino de datos: `SUB_MAP` / ID de peer del hotspot, al hotspot exacto del destino — también en el propio MASTER de entrada | +| `SUB_MAP` | Nunca aprende el ID de origen del plugin (así las respuestas no se reparten por los hotspots del MASTER) | +| 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 | + +El ritmo (unos 60 ms por ráfaga) lo marca el plugin, por ejemplo con `call_later`. + +--- + ## Bus de eventos `PluginBus` (`application/plugins/application/bus.py`) entrega eventos al `on_event` de cada plugin cargado. diff --git a/src/adn_server/application/plugins/application/context.py b/src/adn_server/application/plugins/application/context.py index d294fe2..28693f7 100644 --- a/src/adn_server/application/plugins/application/context.py +++ b/src/adn_server/application/plugins/application/context.py @@ -33,3 +33,5 @@ class ServerContext: defer_to_thread: Callable[..., Any] call_from_reactor: Callable[..., Any] call_later: Callable[..., Any] + # Present only for a plugin listed in PLUGINS.send (see domain/send.py). + send_dmrd: Callable[[bytes], bool] | None = None diff --git a/src/adn_server/application/plugins/application/data_bridge.py b/src/adn_server/application/plugins/application/data_bridge.py index 8ca505e..9f716dd 100644 --- a/src/adn_server/application/plugins/application/data_bridge.py +++ b/src/adn_server/application/plugins/application/data_bridge.py @@ -39,7 +39,7 @@ class DataPluginBridge: self._bus = bus self._config = config # (origin_system, slot) -> (stream_id, start_pkt_time) - self._active_streams: dict[tuple[str, int], tuple[int, float]] = {} + self._active_streams: dict[tuple[str, int], tuple[int, float, bool]] = {} def update_config(self, config: dict[str, Any]) -> None: self._config = config @@ -66,7 +66,9 @@ class DataPluginBridge: obp_rssi: bytes, obp_source_rptr: bytes, forwarded: list[str], + synthetic: bool = False, ) -> None: + """``synthetic``: the frame was sent by a plugin (``ServerContext.send_dmrd``).""" if not self._bus.has_subscribers(): return systems_cfg = self._config.get("SYSTEMS", {}) @@ -78,10 +80,10 @@ class DataPluginBridge: stream_key = (system_name, slot) prev = self._active_streams.get(stream_key) if prev is not None and prev[0] != sid: - self._emit_end(system_name, slot, pkt_time, prev[0], prev[1]) + self._emit_end(system_name, slot, pkt_time, prev[0], prev[1], prev[2]) is_new_stream = prev is None or prev[0] != sid if is_new_stream: - self._active_streams[stream_key] = (sid, pkt_time) + self._active_streams[stream_key] = (sid, pkt_time, synthetic) ctx = CallLegContext( call_family="DATA", direction="RX", @@ -93,7 +95,7 @@ class DataPluginBridge: slot=slot, stream_id=sid, server_id=server_id, - is_synthetic=False, + is_synthetic=synthetic, is_proxy_ingress=is_proxy_inject_only(self._config, system_name), pkt_time=pkt_time, obp_source_server_id=int_byte(obp_source_server) if source_is_obp else None, @@ -139,6 +141,7 @@ class DataPluginBridge: pkt_time: float, stream_id: int, start_time: float, + synthetic: bool = False, ) -> None: duration = max(0.0, pkt_time - start_time) systems_cfg = self._config.get("SYSTEMS", {}) @@ -155,7 +158,7 @@ class DataPluginBridge: slot=slot, stream_id=stream_id, server_id=server_id, - is_synthetic=False, + is_synthetic=synthetic, is_proxy_ingress=is_proxy_inject_only(self._config, system_name), pkt_time=pkt_time, ) diff --git a/src/adn_server/application/plugins/application/manager.py b/src/adn_server/application/plugins/application/manager.py index 5a54418..be69cbb 100644 --- a/src/adn_server/application/plugins/application/manager.py +++ b/src/adn_server/application/plugins/application/manager.py @@ -23,12 +23,14 @@ from __future__ import annotations import logging +from dataclasses import replace from pathlib import Path -from typing import Any +from typing import Any, Callable from ..application.bus import PluginBus from ..application.context import ServerContext from ..domain.protocol import ServerPlugin +from ..domain.send import send_permission from ..infrastructure.loader import ( load_plugin_from_entry, plugin_config_for_load, @@ -45,9 +47,11 @@ class PluginManager: bus: PluginBus, ctx: ServerContext, project_root: str | Path, + sender_factory: Callable[[str], Callable[[bytes], bool]] | None = None, ) -> None: self._bus = bus self._ctx = ctx + self._sender_factory = sender_factory self._project_root = Path(project_root) self._loaded: dict[str, ServerPlugin] = {} self._configs: dict[str, dict[str, Any]] = {} @@ -79,7 +83,7 @@ class PluginManager: cfg = plugin_config_for_load(entry.config) if overrides.get(entry.name): cfg = {**cfg, **overrides[entry.name]} - plugin.on_load(self._bus, cfg, self._ctx) + plugin.on_load(self._bus, cfg, self._ctx_for(entry.name, server_config)) self._loaded[entry.name] = plugin self._configs[entry.name] = dict(entry.config) loaded.append(plugin) @@ -106,7 +110,7 @@ class PluginManager: cfg = plugin_config_for_load(entry.config) if overrides.get(entry.name): cfg = {**cfg, **overrides[entry.name]} - plugin.on_load(self._bus, cfg, self._ctx) + plugin.on_load(self._bus, cfg, self._ctx_for(entry.name, server_config)) self._loaded[entry.name] = plugin self._configs[entry.name] = dict(entry.config) self._bus.register_plugin(plugin) @@ -126,6 +130,17 @@ class PluginManager: logger.exception("(PLUGIN-MANAGER) reload failed for %s", entry.name) self._configs[entry.name] = dict(new_cfg) + def _ctx_for(self, name: str, server_config: dict[str, Any]) -> ServerContext: + """The shared context, plus ``send_dmrd`` for a plugin granted it in PLUGINS.send. + + 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: + return self._ctx + logger.info("(PLUGIN-MANAGER) %s may send unit data (PLUGINS.send)", name) + return replace(self._ctx, send_dmrd=self._sender_factory(name)) + def shutdown_all(self) -> None: for name in list(self._loaded): self._unload_one(name) diff --git a/src/adn_server/application/plugins/application/sender.py b/src/adn_server/application/plugins/application/sender.py new file mode 100644 index 0000000..f6217ba --- /dev/null +++ b/src/adn_server/application/plugins/application/sender.py @@ -0,0 +1,110 @@ +# ADN DMR Peer Server - plugin DMRD sender +# +# 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 +############################################################################### + +"""``ServerContext.send_dmrd`` for one plugin: guards in front of the routing core.""" + +from __future__ import annotations + +import logging +import math +import threading +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 + +logger = logging.getLogger(__name__) + + +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. + """ + + def __init__( + self, + plugin: str, + server_config: dict[str, Any], + deliver: Callable[[bytes, str], Any], + call_from_reactor: Callable[..., Any], + clock: Callable[[], float] = time.monotonic, + ) -> None: + self._plugin = plugin + self._config = server_config + self._deliver = deliver + self._call_from_reactor = call_from_reactor + self._clock = clock + self._lock = threading.Lock() + self._tokens = math.inf # starts full: clamped to the rate on first use + self._refilled = clock() + self._streams_seen: set[bytes] = set() + self.dropped = 0 + + def __call__(self, pkt: bytes) -> bool: + """Queue one DMRD frame for routing; False (and logged) when a guard rejects it.""" + permission = send_permission(self._config, self._plugin) + if permission is None: + return self._reject("not allowed to send (PLUGINS.send)") + 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") + 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: + return self._reject(f"source {rf_src} not in allowed_src_ids") + if not self._take_token(permission): + return self._reject(f"over {permission.max_frames_per_s:g} frames/s") + with self._lock: + first = stream_id not in self._streams_seen + if first: + if len(self._streams_seen) > 1024: + self._streams_seen.clear() + self._streams_seen.add(stream_id) + if first: + logger.info( + "(PLUGIN) %s sent stream %s src %s -> dst %s", self._plugin, int_id(stream_id), rf_src, int_id(dst_id) + ) + self._call_from_reactor(self._deliver, bytes(pkt), self._plugin) + return True + + def _take_token(self, permission: SendPermission) -> bool: + with self._lock: + now = self._clock() + rate = permission.max_frames_per_s + self._tokens = min(rate, self._tokens + (now - self._refilled) * rate) + self._refilled = now + if self._tokens < 1.0: + return False + self._tokens -= 1.0 + return True + + def _reject(self, reason: str) -> bool: + with self._lock: + self.dropped += 1 + dropped = self.dropped + if dropped == 1 or dropped % 100 == 0: + logger.warning("(PLUGIN) %s: frame dropped, %s (%d dropped so far)", self._plugin, reason, dropped) + return False diff --git a/src/adn_server/application/plugins/domain/send.py b/src/adn_server/application/plugins/domain/send.py new file mode 100644 index 0000000..be3362a --- /dev/null +++ b/src/adn_server/application/plugins/domain/send.py @@ -0,0 +1,101 @@ +# ADN DMR Peer Server - plugin send rules +# +# 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 +############################################################################### + +"""What a plugin may send, and the per-plugin permission read from ``PLUGINS.send``.""" + +from __future__ import annotations + +from dataclasses import dataclass +from typing import Any, NamedTuple + +from ....domain import HBPF_DATA_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}) +DEFAULT_MAX_FRAMES_PER_S = 40.0 + + +class DmrdHeader(NamedTuple): + seq: int + rf_src: bytes + dst_id: bytes + peer_id: bytes + slot: int + call_type: str + frame_type: int + dtype_vseq: int + stream_id: bytes + + +def parse_dmrd_header(pkt: bytes) -> DmrdHeader | None: + """The HBP DMRD header fields of any call type (group, vcsbk or unit).""" + if len(pkt) < 53 or pkt[:4] != b"DMRD": + return None + bits = pkt[15] + if bits & 0x40: + call_type = "unit" + elif (bits & 0x23) == 0x23: + call_type = "vcsbk" + else: + call_type = "group" + return DmrdHeader( + seq=pkt[4], + rf_src=pkt[5:8], + dst_id=pkt[8:11], + peer_id=pkt[11:15], + slot=2 if bits & 0x80 else 1, + call_type=call_type, + frame_type=(bits & 0x30) >> 4, + dtype_vseq=bits & 0xF, + stream_id=pkt[16:20], + ) + + +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 + + +@dataclass(frozen=True) +class SendPermission: + allowed_src_ids: frozenset[int] + max_frames_per_s: float + + +def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermission | None: + """The plugin's entry in ``PLUGINS.send``, or None when it may not send. + + An entry without source IDs grants nothing: the allowlist is what stops a + plugin from sending as a radio. + """ + plugins_cfg = server_config.get("PLUGINS") or {} + if plugins_cfg.get("master_kill"): + return None + entry = (plugins_cfg.get("send") or {}).get(plugin) + if not isinstance(entry, dict): + return None + 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)) + except (TypeError, ValueError): + return None + if not ids or rate <= 0: + return None + return SendPermission(ids, rate) diff --git a/src/adn_server/application/routing/announcement_ptt_inject.py b/src/adn_server/application/routing/announcement_ptt_inject.py index 2096dc6..bd83e1f 100644 --- a/src/adn_server/application/routing/announcement_ptt_inject.py +++ b/src/adn_server/application/routing/announcement_ptt_inject.py @@ -45,6 +45,7 @@ from __future__ import annotations from typing import Any +from ..plugins.domain.send import parse_dmrd_header from ..proxy.deployment import proxy_target_system from .helpers import parse_dmrd_burst_fields @@ -98,3 +99,36 @@ def inject_announcement_ptt( ingress_pkt_time=pkt_time, synthetic_announcement=True, ) + + +def inject_plugin_dmrd( + routing: Any, + master_system: str, + pkt: bytes, + *, + pkt_time: float, + server_id: bytes, + plugin: str, +) -> bool | None: + """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. + """ + header = parse_dmrd_header(pkt) + if header is None: + return False + return routing.dmrd_received( + master_system, + server_id, + header.rf_src, + header.dst_id, + header.seq, + header.slot, + header.call_type, + header.frame_type, + header.dtype_vseq, + header.stream_id, + pkt, + ingress_pkt_time=pkt_time, + plugin_origin=plugin, + ) diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 34935e8..7e98c2b 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -49,6 +49,7 @@ from ..domain import ( ) from ..domain.dmr import bptc from ..domain.mesh_session import ObpBridgeSession, obp_session +from .plugins.domain.send import is_plugin_sendable from .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder from .reporting_use_cases import ReportingUseCases from .routing.hbp_forward import HbpForwardMixin @@ -202,6 +203,7 @@ class RoutingUseCases( obp_rssi: bytes = b"\x00", obp_source_rptr: bytes = b"\x00\x00\x00\x00", synthetic_announcement: bool = False, + plugin_origin: str | None = None, ) -> bool: """Called by UDP when DMRD is received. Forward to other systems in same bridge (to_target). @@ -218,6 +220,12 @@ class RoutingUseCases( fuzzy fallback must be skipped: matching the announcement's spoofed rf_src against a real connected peer's ID would misattribute the call (TX state, 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. """ if not self._send_to_system: return @@ -231,6 +239,9 @@ class RoutingUseCases( # Legacy bridge_master 3080–3085: private call to ID 4000 only disconnects dynamics; do not route as PC. 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) + return False if call_type == "unit": _int_dst = int_id(dst_id) if dtype_vseq in (6, 7, 8) or (dtype_vseq == 3 and not self._is_stream_known(system_name, stream_id, slot)): @@ -240,11 +251,14 @@ class RoutingUseCases( obp_use_parsed=obp_use_parsed, obp_hops=obp_hops, obp_source_server=obp_source_server, obp_ber=obp_ber, obp_rssi=obp_rssi, obp_source_rptr=obp_source_rptr, + plugin_origin=plugin_origin, ) # Legacy routerHBP ~3252-3254: 7-digit routing runs after unit data, not # instead of it. D-APRS ARS/LRRP downlink (dst 7300392) uses pvt_call_received # when SUB_MAP sendDataToHBP is blocked by the strict idle check. - if len(str(_int_dst)) == 7: + # Not for plugin frames: that path learns SUB_MAP from the source and keeps + # per-slot call state on the ingress MASTER, which a plugin is not a radio of. + if len(str(_int_dst)) == 7 and plugin_origin is None: self._pvt_call_received( system_name, peer_id, rf_src, dst_id, seq, slot, frame_type, dtype_vseq, stream_id, data, @@ -1075,6 +1089,7 @@ class RoutingUseCases( obp_ber: bytes = b"\x00", obp_rssi: bytes = b"\x00", obp_source_rptr: bytes = b"\x00\x00\x00\x00", + plugin_origin: str | None = None, ) -> None: """Legacy routerOBP/routerHBP unit data branch: DATA-GATEWAY + OBP fan-out + SUB_MAP/hotspot.""" pkt_time = time.time() @@ -1199,8 +1214,8 @@ class RoutingUseCases( ) ) - # DATA-GATEWAY forwarding (legacy ~2281-2284 / ~3083-3087) - if global_cfg.get("DATA_GATEWAY"): + # DATA-GATEWAY forwarding (legacy ~2281-2284 / ~3083-3087). Plugin frames stay local. + if global_cfg.get("DATA_GATEWAY") and plugin_origin is None: dg_cfg = systems_cfg.get("DATA-GATEWAY", {}) if dg_cfg.get("MODE") == "OPENBRIDGE" and dg_cfg.get("ENABLED"): logger.debug("(%s) DATA packet sent to DATA-GATEWAY", system_name) @@ -1226,7 +1241,7 @@ class RoutingUseCases( and pkt_time - _local_sub[2] < UNIT_DATA_LOCAL_SUB_MAX_AGE ) for sys_name, sys_cfg in systems_cfg.items(): - if _dst_is_fresh_local_sub: + if _dst_is_fresh_local_sub or plugin_origin is not None: break if sys_name == system_name: continue @@ -1331,6 +1346,7 @@ class RoutingUseCases( obp_rssi=_rssi, obp_source_rptr=_source_rptr, forwarded=_forwarded, + synthetic=plugin_origin is not None, ) def _pvt_call_received( diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index cb0610d..6d93bbc 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -42,6 +42,7 @@ 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.manager import PluginManager +from adn_server.application.plugins.application.sender import PluginDmrdSender from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge from adn_server.application.proxy.deployment import ( is_obp_proxy_managed, @@ -447,7 +448,28 @@ def run_peer_server( call_from_reactor=reactor.callFromThread, call_later=reactor.callLater, ) - plugin_manager = PluginManager(plugin_bus, server_ctx, project_root) + 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, + ) + + 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_manager = PluginManager( + plugin_bus, + server_ctx, + project_root, + sender_factory=lambda name: PluginDmrdSender(name, config, _deliver_plugin_dmrd, reactor.callFromThread), + ) voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks) data_plugin_bridge = DataPluginBridge(plugin_bus, config) diff --git a/tests/application/test_plugin_sender.py b/tests/application/test_plugin_sender.py new file mode 100644 index 0000000..6d1a1c3 --- /dev/null +++ b/tests/application/test_plugin_sender.py @@ -0,0 +1,144 @@ +# ADN DMR Peer Server - tests plugin send guards +# +# 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 +############################################################################### + +"""ServerContext.send_dmrd: opt-in, source allowlist, unit data only, rate limit.""" + +from __future__ import annotations + +from pathlib import Path + +from tests.harness.deterministic import PacketSpec + +from adn_server.application.plugins.application.bus import PluginBus +from adn_server.application.plugins.application.context import ServerContext +from adn_server.application.plugins.application.manager import PluginManager +from adn_server.application.plugins.application.sender import PluginDmrdSender +from adn_server.domain import HBPF_DATA_SYNC, HBPF_VOICE + +GATEWAY_ID = 900999 + + +def _config(**send) -> dict: + return {"PLUGINS": {"send": {"d-aprs": {"allowed_src_ids": [GATEWAY_ID], "max_frames_per_s": 5, **send}}}} + + +def _frame(src: int = GATEWAY_ID, call_type: str = "unit", frame_type: int = HBPF_DATA_SYNC, dtype: int = 6) -> bytes: + return PacketSpec(rf_src=src, dst_id=7140023, call_type=call_type, frame_type=frame_type, dtype_vseq=dtype).data() + + +class _Clock: + def __init__(self) -> None: + self.t = 100.0 + + def __call__(self) -> float: + return self.t + + +def _sender(config: dict, clock=None): + delivered: list = [] + sender = PluginDmrdSender( + "d-aprs", config, lambda pkt, name: delivered.append((pkt, name)), + call_from_reactor=lambda fn, *a: fn(*a), clock=clock or _Clock(), + ) + return sender, delivered + + +def test_an_allowed_frame_is_delivered_on_the_reactor_with_the_plugin_name() -> None: + sender, delivered = _sender(_config()) + assert sender(_frame()) is True + assert delivered == [(_frame(), "d-aprs")] + + +def test_nothing_is_sent_without_a_plugins_send_entry() -> None: + sender, delivered = _sender({"PLUGINS": {}}) + assert sender(_frame()) is False and delivered == [] + + +def test_an_entry_without_source_ids_grants_nothing() -> None: + sender, delivered = _sender(_config(allowed_src_ids=[])) + assert sender(_frame()) is False and delivered == [] + + +def test_a_plugin_cannot_send_as_a_radio() -> None: + sender, delivered = _sender(_config()) + 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 + assert sender(b"not a dmrd frame") is False + assert delivered == [] + + +def test_rate_limit_per_plugin() -> None: + clock = _Clock() + sender, delivered = _sender(_config(max_frames_per_s=5), clock) + assert [sender(_frame()) for _ in range(7)] == [True] * 5 + [False] * 2 # starts full + clock.t += 0.2 # one more token + assert sender(_frame()) is True + assert len(delivered) == 6 and sender.dropped == 2 + + +def test_revoking_the_permission_takes_effect_on_the_next_frame() -> None: + config = _config() + sender, delivered = _sender(config) + assert sender(_frame()) is True + del config["PLUGINS"]["send"]["d-aprs"] # what a SIGHUP reload leaves behind + assert sender(_frame()) is False + assert len(delivered) == 1 + + +def test_master_kill_stops_sending() -> None: + config = _config() + sender, delivered = _sender(config) + config["PLUGINS"]["master_kill"] = True + assert sender(_frame()) is False and delivered == [] + + +def _plugin_dir(root: Path, name: str) -> None: + pkg = root / name / "plugin" + pkg.mkdir(parents=True) + (root / name / "config.yaml").write_text("enabled: true\n") + (pkg / "__init__.py").write_text( + "class _P:\n" + f" name = {name!r}\n" + " ctx = None\n" + " def on_load(self, bus, config, ctx):\n" + " type(self).ctx = ctx\n" + " def on_event(self, event): pass\n" + " def on_reload(self, config): pass\n" + " def on_shutdown(self): pass\n" + "def create_plugin():\n" + " return _P()\n" + ) + + +def test_only_the_granted_plugin_gets_send_dmrd(tmp_path) -> None: + _plugin_dir(tmp_path / "plugins", "d-aprs") + _plugin_dir(tmp_path / "plugins", "logger") + ctx = ServerContext(config={}, project_root=str(tmp_path), defer_to_thread=None, call_from_reactor=None, call_later=None) + made: list[str] = [] + manager = PluginManager(PluginBus(), ctx, tmp_path, sender_factory=lambda name: made.append(name) or (lambda pkt: True)) + loaded = {p.name: p for p in manager.discover_and_load(_config())} + assert type(loaded["d-aprs"]).ctx.send_dmrd is not None + assert type(loaded["logger"]).ctx.send_dmrd is None + assert made == ["d-aprs"] diff --git a/tests/routing/test_plugin_send_routing.py b/tests/routing/test_plugin_send_routing.py new file mode 100644 index 0000000..8e6a04d --- /dev/null +++ b/tests/routing/test_plugin_send_routing.py @@ -0,0 +1,115 @@ +# ADN DMR Peer Server - tests plugin send routing +# +# 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 +############################################################################### + +"""Unit data a plugin sends (ServerContext.send_dmrd) is routed locally, and only locally.""" + +from __future__ import annotations + +from tests.harness.deterministic import ( + DeterministicScenario, + PacketSpec, + add_openbridge_system, + minimal_config, + patch_routing_wall_time, +) +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.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 + +GATEWAY_ID = 900999 +RADIO_ID = 7140023 +RADIO_HOTSPOT = bytes_4(714002301) +SERVER_ID = bytes_4(2131) + + +def _scenario(radio_on: str = "SYSTEM-B") -> DeterministicScenario: + config = minimal_config(("SYSTEM", "SYSTEM-B")) + add_openbridge_system(config, "OBP-1") + config["SYSTEMS"]["OBP-1"]["VER"] = 5 + config["GLOBAL"]["DATA_GATEWAY"] = True + add_openbridge_system(config, "DATA-GATEWAY") + config["SYSTEMS"]["DATA-GATEWAY"]["ENABLED"] = True + sc = DeterministicScenario(config=config) + for name in ("SYSTEM", "SYSTEM-B"): + sc.protocols[name].STATUS[2] = idle_hbp_slot() + sc.config["SYSTEMS"][name]["GROUP_HANGTIME"] = 0 + sc.protocols[radio_on]._peers[RADIO_HOTSPOT] = {} + sc.config["_SUB_MAP"] = {bytes_3(RADIO_ID): (radio_on, 2, sc.clock.time(), RADIO_HOTSPOT)} + return sc + + +def _ars_ack(dst: int = RADIO_ID, *, call_type: str = "unit", frame_type: int = HBPF_DATA_SYNC) -> bytes: + return PacketSpec( + rf_src=GATEWAY_ID, dst_id=dst, peer_id=1, slot=2, call_type=call_type, + frame_type=frame_type, dtype_vseq=6, stream_id=0x0A0B0C0D, payload=bytes(range(33)), + ).data() + + +def _send(sc: DeterministicScenario, pkt: bytes): + with patch_routing_wall_time(sc.clock): + return inject_plugin_dmrd(sc.routing, "SYSTEM", pkt, pkt_time=sc.clock.time(), server_id=SERVER_ID, plugin="d-aprs") + + +def test_reaches_the_radio_on_another_system_and_nothing_else() -> None: + sc = _scenario("SYSTEM-B") + _send(sc, _ars_ack()) + assert [peer for peer, _ in sc.protocols["SYSTEM-B"].sent_to_peer] == [RADIO_HOTSPOT] + assert sc.capture.packets == [] # no OpenBridge, no DATA-GATEWAY + + +def test_reaches_a_radio_behind_the_ingress_master_itself() -> None: + # The common case: every hotspot sits on the proxy MASTER the plugin injects on. + sc = _scenario("SYSTEM") + _send(sc, _ars_ack()) + assert [peer for peer, _ in sc.protocols["SYSTEM"].sent_to_peer] == [RADIO_HOTSPOT] + + +def test_unknown_destination_is_not_flooded_to_openbridge() -> None: + sc = _scenario() + _send(sc, _ars_ack(dst=3341234)) + assert sc.capture.for_system("OBP-1") == [] + assert sc.capture.for_system("DATA-GATEWAY") == [] + + +def test_the_plugin_source_is_never_learned_in_sub_map() -> None: + sc = _scenario() + _send(sc, _ars_ack()) + assert bytes_3(GATEWAY_ID) not in sc.config["_SUB_MAP"] + + +def test_group_traffic_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 == [] + + +def test_plugins_see_their_own_frames_as_synthetic() -> None: + sc = _scenario() + events: list[object] = [] + bus = PluginBus() + bus.subscribe(events.append) + sc.routing._data_plugin_bridge = DataPluginBridge(bus, sc.config) + _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) From ecddd1191d17432cf4175b3e1141851501950fb9 Mon Sep 17 00:00:00 2001 From: yo Date: Thu, 24 Sep 2026 23:36:09 +0200 Subject: [PATCH 2/3] feat(plugins): group voice from plugins, routed like announcements First step to move scheduled announcements, TTS and voice beacons out of the core into a plugin: a plugin granted `group_voice_tgs` in PLUGINS.send can send group voice on those talkgroups. - PluginIngress (application layer) enters plugin frames on the announcement MASTER. Group voice goes through dmrd_received with synthetic_announcement=True, exactly the path announcements use (the TG's bridges, OpenBridge included), then to that MASTER's hotspots. - While a plugin stream plays it holds the MASTER slot (TX_TYPE=VHEAD, TX_STREAM_ID, TX_RFS, TX_TGID), so routed voice finds it busy; a radio or another stream on the slot fails the frame; VTERM frees it. - send_dmrd called on the reactor thread returns the routing result, so a plugin knows when to stop; from a worker thread it is queued. - Private voice stays refused. Unit data is unchanged (local only). Co-Authored-By: Claude Opus 5.5 --- docs/en/server/user-guide/plugins.md | 15 ++- docs/es/server/user-guide/plugins.md | 15 ++- .../plugins/application/ingress.py | 112 ++++++++++++++++++ .../application/plugins/application/sender.py | 21 +++- .../application/plugins/domain/send.py | 24 +++- .../routing/announcement_ptt_inject.py | 2 + .../application/routing_use_cases.py | 12 +- .../infrastructure/bootstrap/peer_server.py | 26 ++-- tests/application/test_plugin_sender.py | 34 +++++- tests/routing/test_plugin_send_routing.py | 82 ++++++++++++- 10 files changed, 296 insertions(+), 47 deletions(-) create mode 100644 src/adn_server/application/plugins/application/ingress.py 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] From c418b1a893304fb054617dc8f983b4c129d7bd28 Mon Sep 17 00:00:00 2001 From: yo Date: Fri, 25 Sep 2026 07:53:13 +0200 Subject: [PATCH 3/3] fix(plugins): PLUGINS is read from adn-server.yaml and follows SIGHUP Review of #104: master_kill and revoking PLUGINS.send did not take effect on SIGHUP, because merge_top_level_config never re-read PLUGINS. Going through the real path showed it goes further: YamlConfigLoader builds the config from a fixed list of sections and dropped PLUGINS altogether, so the documented block (directory, master_kill, overrides) was never read, at startup either. Both predate this PR. - YamlConfigLoader keeps PLUGINS when present (like OBP_PROXY). - merge_top_level_config replaces PLUGINS on every reload, absent included: removing the block removes everything in it. - Tests through prepare_reload_config + merge_top_level_config + swap_runtime_config and a ConfigProxy, as the server wires them: entry removed, master_kill, whole section removed, a TG granted; and one with the real loader on a real adn-server.yaml, boot and reload. Also from the review: - parse_dmrd_header uses call_attributes(); PluginIngress uses server_id_bytes(). - A max_frames_per_s that is not a finite positive number grants nothing. - The allowlist is documented as a guard against buggy plugins, not a sandbox: a plugin runs in-process with the live config. Co-Authored-By: Claude Opus 5.5 --- docs/en/server/user-guide/plugins.md | 3 +- docs/es/server/user-guide/plugins.md | 3 +- .../plugins/application/ingress.py | 8 +- .../application/plugins/domain/send.py | 26 +++--- .../infrastructure/config_loader.py | 3 + .../infrastructure/config_reload.py | 5 +- tests/application/test_plugin_sender.py | 82 +++++++++++++++++++ 7 files changed, 108 insertions(+), 22 deletions(-) diff --git a/docs/en/server/user-guide/plugins.md b/docs/en/server/user-guide/plugins.md index e9ca5f1..baa4f0f 100644 --- a/docs/en/server/user-guide/plugins.md +++ b/docs/en/server/user-guide/plugins.md @@ -47,7 +47,7 @@ PLUGINS: | `overrides` | Per-plugin config patches without editing `plugins//config.yaml` | | `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. +On **SIGHUP**, the `PLUGINS` block is re-read from `adn-server.yaml` (removing it counts as removing everything in it) and `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed. ### `config.yaml` reserved keys @@ -116,6 +116,7 @@ PLUGINS: ``` - `send_dmrd` is `None` unless the plugin has an entry with at least one source ID. +- The allowlist and the rate limit guard against a **buggy** plugin (sending as a radio, flooding the mesh). They are no sandbox: a plugin runs in-process with the live config and could rewrite its own entry, so only install plugins you trust. - 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. - 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. diff --git a/docs/es/server/user-guide/plugins.md b/docs/es/server/user-guide/plugins.md index ddae1a5..8cf4120 100644 --- a/docs/es/server/user-guide/plugins.md +++ b/docs/es/server/user-guide/plugins.md @@ -47,7 +47,7 @@ PLUGINS: | `overrides` | Parches por plugin sin editar `plugins//config.yaml` | | `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`. +Con **SIGHUP**, el bloque `PLUGINS` se vuelve a leer de `adn-server.yaml` (quitarlo equivale a quitar todo lo que contenía) y `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`. ### Claves reservadas en `config.yaml` @@ -116,6 +116,7 @@ PLUGINS: ``` - `send_dmrd` es `None` salvo que el plugin tenga una entrada con al menos un ID de origen. +- La lista de IDs permitidos y el límite de ritmo protegen frente a un plugin **con fallos** (que emita como una radio o inunde la red). No son un aislamiento: un plugin corre en el mismo proceso con la configuración viva y podría reescribir su propia entrada, así que instala solo plugins de confianza. - 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. - 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. diff --git a/src/adn_server/application/plugins/application/ingress.py b/src/adn_server/application/plugins/application/ingress.py index fc89361..d3254ad 100644 --- a/src/adn_server/application/plugins/application/ingress.py +++ b/src/adn_server/application/plugins/application/ingress.py @@ -26,7 +26,8 @@ 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 ....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 ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind @@ -106,7 +107,4 @@ class PluginIngress: 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) + return server_id_bytes(self._config.get("GLOBAL", {}).get("SERVER_ID"))[:4] diff --git a/src/adn_server/application/plugins/domain/send.py b/src/adn_server/application/plugins/domain/send.py index 066c4fe..383ab1e 100644 --- a/src/adn_server/application/plugins/domain/send.py +++ b/src/adn_server/application/plugins/domain/send.py @@ -22,10 +22,12 @@ from __future__ import annotations +import math from dataclasses import dataclass from typing import Any, NamedTuple from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, HBPF_VOICE_SYNC +from ....domain.mesh_admission import call_attributes # 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}) @@ -50,22 +52,16 @@ def parse_dmrd_header(pkt: bytes) -> DmrdHeader | None: """The HBP DMRD header fields of any call type (group, vcsbk or unit).""" if len(pkt) < 53 or pkt[:4] != b"DMRD": return None - bits = pkt[15] - if bits & 0x40: - call_type = "unit" - elif (bits & 0x23) == 0x23: - call_type = "vcsbk" - else: - call_type = "group" + attrs = call_attributes(pkt[15]) return DmrdHeader( seq=pkt[4], rf_src=pkt[5:8], dst_id=pkt[8:11], peer_id=pkt[11:15], - slot=2 if bits & 0x80 else 1, - call_type=call_type, - frame_type=(bits & 0x30) >> 4, - dtype_vseq=bits & 0xF, + slot=attrs.slot, + call_type=attrs.call_type, + frame_type=attrs.frame_type, + dtype_vseq=attrs.dtype_vseq, stream_id=pkt[16:20], ) @@ -97,8 +93,10 @@ class SendPermission: def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermission | None: """The plugin's entry in ``PLUGINS.send``, or None when it may not send. - An entry without source IDs grants nothing: the allowlist is what stops a - plugin from sending as a radio. + An entry without source IDs grants nothing. The allowlist and the rate limit + guard against a *buggy* plugin (sending as a radio, flooding the mesh); they are + no sandbox: a plugin runs in-process with the live config and could rewrite its + own entry, so only install plugins you trust. """ plugins_cfg = server_config.get("PLUGINS") or {} if plugins_cfg.get("master_kill"): @@ -112,6 +110,6 @@ def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermissio tgs = frozenset(int(t) for t in entry.get("group_voice_tgs") or ()) except (TypeError, ValueError): return None - if not ids or rate <= 0: + if not ids or not math.isfinite(rate) or rate <= 0: return None return SendPermission(ids, rate, tgs) diff --git a/src/adn_server/infrastructure/config_loader.py b/src/adn_server/infrastructure/config_loader.py index 2b58a87..9ede81f 100644 --- a/src/adn_server/infrastructure/config_loader.py +++ b/src/adn_server/infrastructure/config_loader.py @@ -114,6 +114,9 @@ class YamlConfigLoader: # "absent" (defaults apply) from "present but disabled". if isinstance(data.get("OBP_PROXY"), dict): config["OBP_PROXY"] = data["OBP_PROXY"] + # PLUGINS (directory, master_kill, overrides, send) is read by the plugin manager. + if isinstance(data.get("PLUGINS"), dict): + config["PLUGINS"] = data["PLUGINS"] apply_proxy_env_overrides(config) # Ensure REPORT_CLIENTS is list if "REPORT_CLIENTS" in config["REPORTS"] and isinstance(config["REPORTS"]["REPORT_CLIENTS"], str): diff --git a/src/adn_server/infrastructure/config_reload.py b/src/adn_server/infrastructure/config_reload.py index c8e0478..3aaa633 100644 --- a/src/adn_server/infrastructure/config_reload.py +++ b/src/adn_server/infrastructure/config_reload.py @@ -119,12 +119,15 @@ def merge_system_config(old_cfg: dict[str, Any], new_cfg: dict[str, Any]) -> dic def merge_top_level_config(config: dict[str, Any], incoming: dict[str, Any]) -> None: - """Update GLOBAL / REPORTS / ALIASES / LOGGER in the live config dict.""" + """Update GLOBAL / REPORTS / ALIASES / LOGGER / … / PLUGINS in the live config dict.""" kill_flag = config.get("GLOBAL", {}).get("_KILL_SERVER") for key in ("GLOBAL", "REPORTS", "ALIASES", "LOGGER", "PROXY", "DATABASE", "SELF_SERVICE"): if key not in incoming: continue config[key] = copy.deepcopy(incoming[key]) + # PLUGINS always follows the file, absent included: it carries master_kill and the + # per-plugin send permissions, whose removal must take effect on this reload. + config["PLUGINS"] = copy.deepcopy(incoming.get("PLUGINS") or {}) if kill_flag is not None: config.setdefault("GLOBAL", {})["_KILL_SERVER"] = kill_flag for rk in _RUNTIME_TOP_KEYS: diff --git a/tests/application/test_plugin_sender.py b/tests/application/test_plugin_sender.py index 51122ab..18888fc 100644 --- a/tests/application/test_plugin_sender.py +++ b/tests/application/test_plugin_sender.py @@ -24,6 +24,8 @@ from __future__ import annotations from pathlib import Path +import pytest + from tests.harness.deterministic import PacketSpec from adn_server.application.plugins.application.bus import PluginBus @@ -168,3 +170,83 @@ def test_only_the_granted_plugin_gets_send_dmrd(tmp_path) -> None: assert type(loaded["d-aprs"]).ctx.send_dmrd is not None assert type(loaded["logger"]).ctx.send_dmrd is None assert made == ["d-aprs"] + + +# --- through the real SIGHUP reload path (review of #104) --- + +from adn_server.application.runtime_context import ( # noqa: E402 + ConfigProxy, + RuntimeContext, + RuntimeContextHolder, + prepare_reload_config, + swap_runtime_config, +) +from adn_server.infrastructure.config_reload import merge_top_level_config # noqa: E402 + + +def _reload(holder: RuntimeContextHolder, incoming: dict) -> None: + """What a SIGHUP does with a freshly parsed adn-server.yaml (``incoming``).""" + new_config = prepare_reload_config(holder) + merge_top_level_config(new_config, incoming) + swap_runtime_config(holder, new_config) + + +@pytest.mark.parametrize( + "incoming", + [ + {"GLOBAL": {}, "PLUGINS": {"send": {}}}, # entry removed + {"GLOBAL": {}, "PLUGINS": {"master_kill": True, **_config()["PLUGINS"]}}, # master_kill + {"GLOBAL": {}}, # the whole PLUGINS section removed + ], + ids=["entry-removed", "master-kill", "section-removed"], +) +def test_a_sighup_reload_revokes_sending(incoming) -> None: + holder = RuntimeContextHolder(RuntimeContext(config={"GLOBAL": {}, **_config()})) + sender, delivered = _sender(ConfigProxy(holder)) + assert sender(_frame()) is True + _reload(holder, incoming) + assert sender(_frame()) is False + assert len(delivered) == 1 + + +def test_a_sighup_reload_can_grant_a_new_talkgroup() -> None: + holder = RuntimeContextHolder(RuntimeContext(config={"GLOBAL": {}, **_config()})) + sender, _ = _sender(ConfigProxy(holder)) + 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 + _reload(holder, {"GLOBAL": {}, **_config(group_voice_tgs=[213])}) + assert sender(voice) is True + + +def test_plugins_section_is_read_from_adn_server_yaml_and_followed_on_reload(tmp_path) -> None: + """The whole chain with the real loader: PLUGINS used to be dropped by YamlConfigLoader, + so neither master_kill, overrides nor send permissions ever reached the server.""" + import logging + from pathlib import Path + + from adn_server.infrastructure.config_loader import YamlConfigLoader + from adn_server.infrastructure.config_reload import prepare_incoming_config + + example = Path(__file__).resolve().parents[2] / "adn-server.example.yaml" + base = example.read_text(encoding="utf-8") + path = tmp_path / "adn-server.yaml" + grant = "\nPLUGINS:\n send:\n d-aprs:\n allowed_src_ids: [900999]\n" + path.write_text(base + grant, encoding="utf-8") + log = logging.getLogger("test") + + boot = prepare_incoming_config(YamlConfigLoader(), str(path), log) + assert boot["PLUGINS"]["send"]["d-aprs"]["allowed_src_ids"] == [900999] + holder = RuntimeContextHolder(RuntimeContext(config=boot)) + sender, _ = _sender(ConfigProxy(holder)) + assert sender(_frame()) is True + + path.write_text(base + "\nPLUGINS:\n master_kill: true\n" + grant[len("\nPLUGINS:\n"):], encoding="utf-8") + _reload(holder, prepare_incoming_config(YamlConfigLoader(), str(path), log)) + assert sender(_frame()) is False + + +@pytest.mark.parametrize("rate", [float("inf"), float("nan"), 0, -5]) +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