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)