Merge pull request #104 from pyopower/feat/plugin-send-dmrd
feat(plugins): let a plugin send unit data and group voice, opt-in per pluginpull/105/head
commit
044ba970df
@ -0,0 +1,110 @@
|
||||
# ADN DMR Peer Server - plugin frame ingress
|
||||
#
|
||||
# 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
|
||||
###############################################################################
|
||||
|
||||
"""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
|
||||
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
|
||||
|
||||
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:
|
||||
return server_id_bytes(self._config.get("GLOBAL", {}).get("SERVER_ID"))[:4]
|
||||
@ -0,0 +1,121 @@
|
||||
# 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 GROUP_VOICE, SendPermission, parse_dmrd_header, plugin_frame_kind, 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.
|
||||
|
||||
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__(
|
||||
self,
|
||||
plugin: str,
|
||||
server_config: dict[str, Any],
|
||||
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()
|
||||
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")
|
||||
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:
|
||||
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)
|
||||
)
|
||||
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
|
||||
|
||||
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,115 @@
|
||||
# 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
|
||||
|
||||
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})
|
||||
UNIT_DATA = "unit_data"
|
||||
GROUP_VOICE = "group_voice"
|
||||
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
|
||||
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=attrs.slot,
|
||||
call_type=attrs.call_type,
|
||||
frame_type=attrs.frame_type,
|
||||
dtype_vseq=attrs.dtype_vseq,
|
||||
stream_id=pkt[16:20],
|
||||
)
|
||||
|
||||
|
||||
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:
|
||||
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:
|
||||
"""The plugin's entry in ``PLUGINS.send``, or None when it may not send.
|
||||
|
||||
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"):
|
||||
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))
|
||||
tgs = frozenset(int(t) for t in entry.get("group_voice_tgs") or ())
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
if not ids or not math.isfinite(rate) or rate <= 0:
|
||||
return None
|
||||
return SendPermission(ids, rate, tgs)
|
||||
@ -0,0 +1,252 @@
|
||||
# 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
|
||||
|
||||
import pytest
|
||||
|
||||
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_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)
|
||||
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"]
|
||||
|
||||
|
||||
# --- 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
|
||||
@ -0,0 +1,189 @@
|
||||
# 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,
|
||||
active_routing_table,
|
||||
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.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, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, 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_private_voice_from_a_plugin_is_refused() -> None:
|
||||
sc = _scenario()
|
||||
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:
|
||||
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)
|
||||
|
||||
|
||||
# --- 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]
|
||||
Loading…
Reference in new issue