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 <noreply@anthropic.com>
pull/104/head
yo 4 days ago
parent 114ad296f5
commit 0c3432ebfd

@ -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/<name>/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`.

@ -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/<name>/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.

@ -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

@ -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,
)

@ -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)

@ -0,0 +1,110 @@
# ADN DMR Peer Server - plugin DMRD sender
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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.<plugin>`` 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

@ -0,0 +1,101 @@
# ADN DMR Peer Server - plugin send rules
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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)

@ -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,
)

@ -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(

@ -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)

@ -0,0 +1,144 @@
# ADN DMR Peer Server - tests plugin send guards
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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"]

@ -0,0 +1,115 @@
# ADN DMR Peer Server - tests plugin send routing
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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)
Loading…
Cancel
Save

Powered by TurnKey Linux.