feat(plugins): group voice from plugins, routed like announcements

First step to move scheduled announcements, TTS and voice beacons out of
the core into a plugin: a plugin granted `group_voice_tgs` in
PLUGINS.send can send group voice on those talkgroups.

- PluginIngress (application layer) enters plugin frames on the
  announcement MASTER. Group voice goes through dmrd_received with
  synthetic_announcement=True, exactly the path announcements use (the
  TG's bridges, OpenBridge included), then to that MASTER's hotspots.
- While a plugin stream plays it holds the MASTER slot (TX_TYPE=VHEAD,
  TX_STREAM_ID, TX_RFS, TX_TGID), so routed voice finds it busy; a radio
  or another stream on the slot fails the frame; VTERM frees it.
- send_dmrd called on the reactor thread returns the routing result, so
  a plugin knows when to stop; from a worker thread it is queued.
- Private voice stays refused. Unit data is unchanged (local only).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
pull/104/head
yo 4 days ago
parent 0c3432ebfd
commit ecddd1191d

@ -45,7 +45,7 @@ PLUGINS:
| `directory` | Absolute path or relative to project root | | `directory` | Absolute path or relative to project root |
| `master_kill` | Emergency disable — no plugins loaded | | `master_kill` | Emergency disable — no plugins loaded |
| `overrides` | Per-plugin config patches without editing `plugins/<name>/config.yaml` | | `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) | | `send` | Per-plugin permission to send unit data and group voice — see [Sending](#sending-unit-data-and-group-voice-opt-in) |
On **SIGHUP**, `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed. On **SIGHUP**, `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed.
@ -94,15 +94,15 @@ Passed to `on_load` as `server_ctx` (`application/plugins/application/context.py
| `defer_to_thread(fn, *args)` | Run blocking I/O off the reactor | | `defer_to_thread(fn, *args)` | Run blocking I/O off the reactor |
| `call_from_reactor(fn, *args)` | Schedule callback on reactor thread | | `call_from_reactor(fn, *args)` | Schedule callback on reactor thread |
| `call_later(delay_s, fn, *args)` | Reactor timer | | `call_later(delay_s, fn, *args)` | Reactor timer |
| `send_dmrd(pkt) -> bool` | Send one DMRD frame — **only** for a plugin granted it in `PLUGINS.send`, otherwise `None` ([Sending](#sending-unit-data-opt-in)) | | `send_dmrd(pkt) -> bool` | Send one DMRD frame — **only** for a plugin granted it in `PLUGINS.send`, otherwise `None` ([Sending](#sending-unit-data-and-group-voice-opt-in)) |
**Pattern:** do only dispatch in `on_event`; call `defer_to_thread` for file writes, HTTP, heavy CPU. **Pattern:** do only dispatch in `on_event`; call `defer_to_thread` for file writes, HTTP, heavy CPU.
--- ---
## Sending unit data (opt-in) ## Sending unit data and group voice (opt-in)
A plugin can send **unit data** (data header, rate 1/2 and 3/4 blocks, CSBK: ARS, LRRP, SMS…) with `server_ctx.send_dmrd(pkt)`, one complete HBP `DMRD` frame per call. Enabling a plugin never lets it transmit: the sysop grants it per plugin in `adn-server.yaml`: A plugin can send **unit data** (data header, rate 1/2 and 3/4 blocks, CSBK: ARS, LRRP, SMS…) and **group voice** on granted talkgroups (voice beacons, announcements) with `server_ctx.send_dmrd(pkt)`, one complete HBP `DMRD` frame per call. Enabling a plugin never lets it transmit: the sysop grants it per plugin in `adn-server.yaml`:
```yaml ```yaml
PLUGINS: PLUGINS:
@ -110,12 +110,15 @@ PLUGINS:
d-aprs: d-aprs:
allowed_src_ids: [900999] # rf_src the plugin may send as; required allowed_src_ids: [900999] # rf_src the plugin may send as; required
max_frames_per_s: 40 # per-plugin token bucket (default 40) max_frames_per_s: 40 # per-plugin token bucket (default 40)
beacon:
allowed_src_ids: [2130035]
group_voice_tgs: [213] # group voice only on these TGs; none: unit data only
``` ```
- `send_dmrd` is `None` unless the plugin has an entry with at least one source ID. - `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. - 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. - Rejected frames (not unit data, source not allowed, over the rate) return `False` and are logged and counted; the first frame of each stream is logged at INFO.
- Safe from any thread: frames are handed to the reactor. - Called on the reactor thread (`on_event`, `call_later`), the frame is routed at once and the result is whether the server **accepted** it; a plugin sending voice must stop when it gets `False`. From another thread the frame is queued to the reactor and `True` only means the guards passed.
How the server treats them: How the server treats them:
@ -127,6 +130,8 @@ How the server treats them:
| OpenBridge / `DATA-GATEWAY` | No fan-out: plugin frames stay on this server | | 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 | | Events | Reach plugins with `is_synthetic=True`, so a plugin can ignore its own frames |
Group voice is routed **like a scheduled announcement** (synthetic PTT on the same MASTER): through the TG's bridges, OpenBridge legs included, and to the hotspots of that MASTER. While a plugin stream plays it holds that MASTER slot, so routed calls find it busy; a radio or another stream on the slot makes the next frame fail. The terminator frees the slot. Private voice can't be sent.
Pacing (about 60 ms per burst) is the plugin's job, e.g. with `call_later`. Pacing (about 60 ms per burst) is the plugin's job, e.g. with `call_later`.
--- ---

@ -45,7 +45,7 @@ PLUGINS:
| `directory` | Ruta absoluta o relativa al project root | | `directory` | Ruta absoluta o relativa al project root |
| `master_kill` | Desactivación de emergencia — no carga plugins | | `master_kill` | Desactivación de emergencia — no carga plugins |
| `overrides` | Parches por plugin sin editar `plugins/<name>/config.yaml` | | `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) | | `send` | Permiso por plugin para enviar datos y voz de grupo — ver [Envío](#envío-de-datos-y-voz-de-grupo-opcional) |
Con **SIGHUP**, `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`. Con **SIGHUP**, `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`.
@ -94,15 +94,15 @@ Se pasa a `on_load` como `server_ctx` (`application/plugins/application/context.
| `defer_to_thread(fn, *args)` | Ejecutar I/O bloqueante fuera del reactor | | `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_from_reactor(fn, *args)` | Programar callback en el hilo del reactor |
| `call_later(delay_s, fn, *args)` | Temporizador del reactor | | `call_later(delay_s, fn, *args)` | Temporizador del reactor |
| `send_dmrd(pkt) -> bool` | Enviar una trama DMRD — **solo** para un plugin autorizado en `PLUGINS.send`; si no, `None` ([Envío](#envío-de-datos-opcional)) | | `send_dmrd(pkt) -> bool` | Enviar una trama DMRD — **solo** para un plugin autorizado en `PLUGINS.send`; si no, `None` ([Envío](#envío-de-datos-y-voz-de-grupo-opcional)) |
**Patrón:** solo despachar en `on_event`; usar `defer_to_thread` para archivos, HTTP o CPU intensiva. **Patrón:** solo despachar en `on_event`; usar `defer_to_thread` para archivos, HTTP o CPU intensiva.
--- ---
## Envío de datos (opcional) ## Envío de datos y voz de grupo (opcional)
Un plugin puede enviar **datos** (cabecera, bloques de 1/2 y 3/4, CSBK: ARS, LRRP, SMS…) con `server_ctx.send_dmrd(pkt)`, una trama HBP `DMRD` completa por llamada. Activar un plugin nunca le permite transmitir: el sysop lo autoriza por plugin en `adn-server.yaml`: Un plugin puede enviar **datos** (cabecera, bloques de 1/2 y 3/4, CSBK: ARS, LRRP, SMS…) y **voz de grupo** en los TGs autorizados (balizas de voz, anuncios) con `server_ctx.send_dmrd(pkt)`, una trama HBP `DMRD` completa por llamada. Activar un plugin nunca le permite transmitir: el sysop lo autoriza por plugin en `adn-server.yaml`:
```yaml ```yaml
PLUGINS: PLUGINS:
@ -110,12 +110,15 @@ PLUGINS:
d-aprs: d-aprs:
allowed_src_ids: [900999] # rf_src con el que puede enviar; obligatorio 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) max_frames_per_s: 40 # cubo de tokens por plugin (por defecto 40)
baliza:
allowed_src_ids: [2130035]
group_voice_tgs: [213] # voz de grupo solo en estos TGs; sin ella, solo datos
``` ```
- `send_dmrd` es `None` salvo que el plugin tenga una entrada con al menos un ID de origen. - `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. - 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. - Las tramas rechazadas (no son datos, origen no permitido, exceso de ritmo) devuelven `False` y se registran y cuentan; la primera trama de cada stream se registra en INFO.
- Se puede llamar desde cualquier hilo: las tramas se entregan al reactor. - Llamado desde el hilo del reactor (`on_event`, `call_later`), la trama se enruta en el acto y el resultado indica si el servidor la **aceptó**; un plugin que emite voz debe parar cuando recibe `False`. Desde otro hilo la trama se encola al reactor y `True` solo significa que pasó las salvaguardas.
Cómo las trata el servidor: Cómo las trata el servidor:
@ -127,6 +130,8 @@ Cómo las trata el servidor:
| OpenBridge / `DATA-GATEWAY` | Sin reparto: las tramas del plugin se quedan en este servidor | | 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 | | Eventos | Llegan a los plugins con `is_synthetic=True`, para que un plugin ignore sus propias tramas |
La voz de grupo se enruta **como un anuncio programado** (PTT sintético en el mismo MASTER): por los puentes del TG, OpenBridge incluidos, y a los hotspots de ese MASTER. Mientras suena, el stream del plugin ocupa ese slot del MASTER, así que las llamadas enrutadas lo encuentran ocupado; una radio u otro stream en el slot hace fallar la trama siguiente. El terminador libera el slot. No se puede enviar voz privada.
El ritmo (unos 60 ms por ráfaga) lo marca el plugin, por ejemplo con `call_later`. El ritmo (unos 60 ms por ráfaga) lo marca el plugin, por ejemplo con `call_later`.
--- ---

@ -0,0 +1,112 @@
# 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, bytes_4, int_id
from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd
from ...routing.helpers import slot_voice_held_by_other_stream
from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind
logger = logging.getLogger(__name__)
class PluginIngress:
"""Routes plugin frames on the reactor thread and reports whether each was accepted.
Unit data goes through the unit data path only (see ``dmrd_received(plugin_origin=)``).
Group voice is routed like a scheduled announcement: through the bridges (OpenBridge
legs included, as for any local ingress), and out to the hotspots of the MASTER it
enters on. While a plugin's stream plays it holds that MASTER slot (TX_TYPE=VHEAD,
TX_STREAM_ID, TX_RFS), so routed voice finds it busy; a radio or another stream on the
slot makes the frame fail, which tells the plugin to stop. The terminator frees it.
"""
def __init__(
self,
routing: Any,
config: dict[str, Any],
get_protocols: Callable[[], dict[str, Any]],
send_local: Callable[[str, bytes], None],
clock: Callable[[], float] = time.time,
) -> None:
self._routing = routing
self._config = config
self._get_protocols = get_protocols
self._send_local = send_local
self._clock = clock
def set_routing(self, routing: Any) -> None:
"""Bootstrap builds the plugin manager before routing; routing is set before any load."""
self._routing = routing
def deliver(self, pkt: bytes, plugin: str) -> bool:
header = parse_dmrd_header(pkt)
kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq) if header else None
master = announcement_ptt_system(self._config)
if kind is None or not master:
if not master:
logger.warning("(PLUGIN) %s: no MASTER to send from, frame dropped", plugin)
return False
server_id = self._server_id()
now = self._clock()
if kind == UNIT_DATA:
return self._route(master, pkt, now, server_id, plugin) is not False
return self._group_voice(master, header, pkt[:11] + server_id + pkt[15:], now, server_id, plugin)
def _group_voice(self, master: str, header: Any, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool:
proto = self._get_protocols().get(master)
slot = getattr(proto, "STATUS", {}).get(header.slot) if proto is not None else None
if slot is None:
return False
if slot_voice_held_by_other_stream(slot, header.stream_id, now):
return False
if self._route(master, pkt, now, server_id, plugin) is not True:
self._release(slot, header.stream_id)
return False
is_term = header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VTERM
slot["TX_TYPE"] = HBPF_SLT_VTERM if is_term else HBPF_SLT_VHEAD
slot["TX_STREAM_ID"] = header.stream_id
slot["TX_RFS"] = header.rf_src
slot["TX_TGID"] = header.dst_id
slot["TX_TIME"] = now
self._send_local(master, pkt)
return True
def _route(self, master: str, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool | None:
return inject_plugin_dmrd(self._routing, master, pkt, pkt_time=now, server_id=server_id, plugin=plugin)
@staticmethod
def _release(slot: dict[str, Any], stream_id: bytes) -> None:
if slot.get("TX_STREAM_ID") == stream_id:
slot["TX_TYPE"] = HBPF_SLT_VTERM
def _server_id(self) -> bytes:
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if isinstance(server_id, bytes) and len(server_id) >= 4:
return server_id[:4]
return bytes_4(int(int_id(server_id) if isinstance(server_id, bytes) else server_id or 0) & 0xFFFFFFFF)

@ -29,7 +29,7 @@ import time
from typing import Any, Callable from typing import Any, Callable
from ....domain import int_id from ....domain import int_id
from ..domain.send import SendPermission, is_plugin_sendable, parse_dmrd_header, send_permission from ..domain.send import GROUP_VOICE, SendPermission, parse_dmrd_header, plugin_frame_kind, send_permission
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -38,8 +38,12 @@ class PluginDmrdSender:
"""Checks each frame against the plugin's live permission, then hands it to the reactor. """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 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 removes or narrows ``PLUGINS.send.<plugin>`` takes effect at once.
from any thread; delivery always happens on the reactor.
Called on the reactor thread (``on_event``, ``call_later``), the frame is routed at
once and the result is whether the server accepted it: a plugin sending voice 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__( def __init__(
@ -49,12 +53,14 @@ class PluginDmrdSender:
deliver: Callable[[bytes, str], Any], deliver: Callable[[bytes, str], Any],
call_from_reactor: Callable[..., Any], call_from_reactor: Callable[..., Any],
clock: Callable[[], float] = time.monotonic, clock: Callable[[], float] = time.monotonic,
in_reactor_thread: Callable[[], bool] = lambda: False,
) -> None: ) -> None:
self._plugin = plugin self._plugin = plugin
self._config = server_config self._config = server_config
self._deliver = deliver self._deliver = deliver
self._call_from_reactor = call_from_reactor self._call_from_reactor = call_from_reactor
self._clock = clock self._clock = clock
self._in_reactor_thread = in_reactor_thread
self._lock = threading.Lock() self._lock = threading.Lock()
self._tokens = math.inf # starts full: clamped to the rate on first use self._tokens = math.inf # starts full: clamped to the rate on first use
self._refilled = clock() self._refilled = clock()
@ -69,8 +75,11 @@ class PluginDmrdSender:
header = parse_dmrd_header(bytes(pkt)) if isinstance(pkt, (bytes, bytearray)) else None header = parse_dmrd_header(bytes(pkt)) if isinstance(pkt, (bytes, bytearray)) else None
if header is None: if header is None:
return self._reject("not a DMRD frame") return self._reject("not a DMRD frame")
if not is_plugin_sendable(header.call_type, header.frame_type, header.dtype_vseq): kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq)
return self._reject("only unit data may be sent") 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) rf_src = int_id(header.rf_src)
stream_id, dst_id = header.stream_id, header.dst_id stream_id, dst_id = header.stream_id, header.dst_id
if rf_src not in permission.allowed_src_ids: if rf_src not in permission.allowed_src_ids:
@ -87,6 +96,8 @@ class PluginDmrdSender:
logger.info( logger.info(
"(PLUGIN) %s sent stream %s src %s -> dst %s", self._plugin, int_id(stream_id), rf_src, int_id(dst_id) "(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) self._call_from_reactor(self._deliver, bytes(pkt), self._plugin)
return True return True

@ -25,10 +25,12 @@ from __future__ import annotations
from dataclasses import dataclass from dataclasses import dataclass
from typing import Any, NamedTuple from typing import Any, NamedTuple
from ....domain import HBPF_DATA_SYNC from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, HBPF_VOICE_SYNC
# CSBK, data header, rate 1/2 and rate 3/4 data blocks: ARS, LRRP, SMS and the like. # 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}) PLUGIN_SENDABLE_DTYPES = frozenset({3, 6, 7, 8})
UNIT_DATA = "unit_data"
GROUP_VOICE = "group_voice"
DEFAULT_MAX_FRAMES_PER_S = 40.0 DEFAULT_MAX_FRAMES_PER_S = 40.0
@ -68,15 +70,28 @@ def parse_dmrd_header(pkt: bytes) -> DmrdHeader | None:
) )
def plugin_frame_kind(call_type: str, frame_type: int, dtype_vseq: int) -> str | None:
"""UNIT_DATA, GROUP_VOICE, or None for anything a plugin may not send (e.g. private voice)."""
if call_type == "unit":
return UNIT_DATA if frame_type == HBPF_DATA_SYNC and dtype_vseq in PLUGIN_SENDABLE_DTYPES else None
if call_type == "group":
if frame_type in (HBPF_VOICE, HBPF_VOICE_SYNC):
return GROUP_VOICE
if frame_type == HBPF_DATA_SYNC and dtype_vseq in (HBPF_SLT_VHEAD, HBPF_SLT_VTERM):
return GROUP_VOICE
return None
def is_plugin_sendable(call_type: str, frame_type: int, dtype_vseq: int) -> bool: 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 plugin_frame_kind(call_type, frame_type, dtype_vseq) is not None
return call_type == "unit" and frame_type == HBPF_DATA_SYNC and dtype_vseq in PLUGIN_SENDABLE_DTYPES
@dataclass(frozen=True) @dataclass(frozen=True)
class SendPermission: class SendPermission:
allowed_src_ids: frozenset[int] allowed_src_ids: frozenset[int]
max_frames_per_s: float 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: def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermission | None:
@ -94,8 +109,9 @@ def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermissio
try: try:
ids = frozenset(int(i) for i in entry.get("allowed_src_ids") or ()) 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)) 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): except (TypeError, ValueError):
return None return None
if not ids or rate <= 0: if not ids or rate <= 0:
return None return None
return SendPermission(ids, rate) return SendPermission(ids, rate, tgs)

@ -113,6 +113,7 @@ def inject_plugin_dmrd(
"""Feed one plugin frame through ``dmrd_received`` on the same MASTER as announcements. """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. The peer field is always the SERVER_ID: a plugin never speaks as a connected peer.
Group voice is marked ``synthetic_announcement`` exactly like announcements are.
""" """
header = parse_dmrd_header(pkt) header = parse_dmrd_header(pkt)
if header is None: if header is None:
@ -130,5 +131,6 @@ def inject_plugin_dmrd(
header.stream_id, header.stream_id,
pkt, pkt,
ingress_pkt_time=pkt_time, ingress_pkt_time=pkt_time,
synthetic_announcement=header.call_type == "group",
plugin_origin=plugin, plugin_origin=plugin,
) )

@ -222,10 +222,12 @@ class RoutingUseCases(
dynamic TG) to that peer. dynamic TG) to that peer.
``plugin_origin`` names the plugin that sent the frame through ``plugin_origin`` names the plugin that sent the frame through
``ServerContext.send_dmrd``. Such frames are unit data only, are delivered by ``ServerContext.send_dmrd``. Unit data is delivered by the unit data path
the unit data path alone (SUB_MAP / hotspot ID, never the private call path, alone (SUB_MAP / hotspot ID, never the private call path, so SUB_MAP never
so SUB_MAP never learns their source), stay on this server (no OpenBridge or learns its source), stays on this server (no OpenBridge or DATA-GATEWAY
DATA-GATEWAY fan-out), and reach plugins as synthetic events. fan-out), and reaches plugins as synthetic events. Group voice comes with
``synthetic_announcement=True`` and is routed exactly like an announcement.
Anything else (private voice) is dropped.
""" """
if not self._send_to_system: if not self._send_to_system:
return return
@ -240,7 +242,7 @@ class RoutingUseCases(
if call_type == "unit" and int_id(dst_id) == 4000: if call_type == "unit" and int_id(dst_id) == 4000:
return return
if plugin_origin is not None and not is_plugin_sendable(call_type, frame_type, dtype_vseq): if plugin_origin is not None and not is_plugin_sendable(call_type, frame_type, dtype_vseq):
logger.warning("(PLUGIN) %s: only unit data may be sent, dropped", plugin_origin) logger.warning("(PLUGIN) %s: only unit data and group voice may be sent, dropped", plugin_origin)
return False return False
if call_type == "unit": if call_type == "unit":
_int_dst = int_id(dst_id) _int_dst = int_id(dst_id)

@ -29,6 +29,7 @@ import time
from typing import Any from typing import Any
from twisted.internet import reactor, task, threads from twisted.internet import reactor, task, threads
from twisted.python.threadable import isInIOThread
from adn_server.application import ( from adn_server.application import (
IdentUseCases, IdentUseCases,
@ -41,6 +42,7 @@ from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases
from adn_server.application.plugins.application.bus import PluginBus from adn_server.application.plugins.application.bus import PluginBus
from adn_server.application.plugins.application.context import ServerContext 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.data_bridge import DataPluginBridge
from adn_server.application.plugins.application.ingress import PluginIngress
from adn_server.application.plugins.application.manager import PluginManager from adn_server.application.plugins.application.manager import PluginManager
from adn_server.application.plugins.application.sender import PluginDmrdSender from adn_server.application.plugins.application.sender import PluginDmrdSender
from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge
@ -448,27 +450,20 @@ def run_peer_server(
call_from_reactor=reactor.callFromThread, call_from_reactor=reactor.callFromThread,
call_later=reactor.callLater, call_later=reactor.callLater,
) )
def _deliver_plugin_dmrd(pkt: bytes, plugin: str) -> None: def _send_local(system: str, pkt: bytes) -> None:
# Reactor thread (PluginDmrdSender schedules it there). Same ingress MASTER as announcements. send_system = getattr(protocols.get(system), "send_system", None)
from adn_server.application.routing.announcement_ptt_inject import ( if callable(send_system):
announcement_ptt_system, send_system(pkt)
inject_plugin_dmrd,
)
master = announcement_ptt_system(config) plugin_ingress = PluginIngress(None, config, lambda: protocols, _send_local) # routing set below
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_manager = PluginManager(
plugin_bus, plugin_bus,
server_ctx, server_ctx,
project_root, project_root,
sender_factory=lambda name: PluginDmrdSender(name, config, _deliver_plugin_dmrd, reactor.callFromThread), sender_factory=lambda name: PluginDmrdSender(
name, config, plugin_ingress.deliver, reactor.callFromThread, in_reactor_thread=isInIOThread,
),
) )
voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks) voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks)
data_plugin_bridge = DataPluginBridge(plugin_bus, config) data_plugin_bridge = DataPluginBridge(plugin_bus, config)
@ -490,6 +485,7 @@ def run_peer_server(
voice_plugin_bridge=voice_plugin_bridge, voice_plugin_bridge=voice_plugin_bridge,
data_plugin_bridge=data_plugin_bridge, data_plugin_bridge=data_plugin_bridge,
) )
plugin_ingress.set_routing(routing_use_cases)
plugin_manager.discover_and_load(config) plugin_manager.discover_and_load(config)
reactor.addSystemEventTrigger("before", "shutdown", plugin_manager.shutdown_all) reactor.addSystemEventTrigger("before", "shutdown", plugin_manager.shutdown_all)
routing_use_cases.apply_startup_subscriptions() routing_use_cases.apply_startup_subscriptions()

@ -81,14 +81,40 @@ def test_a_plugin_cannot_send_as_a_radio() -> None:
assert sender(_frame(src=7140023)) is False and delivered == [] assert sender(_frame(src=7140023)) is False and delivered == []
def test_voice_and_group_frames_are_refused() -> None: def test_private_voice_and_non_dmrd_are_refused() -> None:
sender, delivered = _sender(_config()) sender, delivered = _sender(_config(group_voice_tgs=[213]))
assert sender(_frame(call_type="group")) is False assert sender(_frame(frame_type=HBPF_VOICE, dtype=1)) is False # unit call type, voice frame
assert sender(_frame(frame_type=HBPF_VOICE, dtype=1)) is False
assert sender(b"not a dmrd frame") is False assert sender(b"not a dmrd frame") is False
assert delivered == [] 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: def test_rate_limit_per_plugin() -> None:
clock = _Clock() clock = _Clock()
sender, delivered = _sender(_config(max_frames_per_s=5), clock) sender, delivered = _sender(_config(max_frames_per_s=5), clock)

@ -25,6 +25,7 @@ from __future__ import annotations
from tests.harness.deterministic import ( from tests.harness.deterministic import (
DeterministicScenario, DeterministicScenario,
PacketSpec, PacketSpec,
active_routing_table,
add_openbridge_system, add_openbridge_system,
minimal_config, minimal_config,
patch_routing_wall_time, patch_routing_wall_time,
@ -33,9 +34,10 @@ from tests.routing.unit_data_helpers import idle_hbp_slot
from adn_server.application.plugins.application.bus import PluginBus from adn_server.application.plugins.application.bus import PluginBus
from adn_server.application.plugins.application.data_bridge import DataPluginBridge 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.plugins.domain.events import UnitDataFrame
from adn_server.application.routing.announcement_ptt_inject import inject_plugin_dmrd from adn_server.application.routing.announcement_ptt_inject import inject_plugin_dmrd
from adn_server.domain import HBPF_DATA_SYNC, bytes_3, bytes_4 from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_3, bytes_4
GATEWAY_ID = 900999 GATEWAY_ID = 900999
RADIO_ID = 7140023 RADIO_ID = 7140023
@ -98,10 +100,10 @@ def test_the_plugin_source_is_never_learned_in_sub_map() -> None:
assert bytes_3(GATEWAY_ID) not in sc.config["_SUB_MAP"] assert bytes_3(GATEWAY_ID) not in sc.config["_SUB_MAP"]
def test_group_traffic_from_a_plugin_is_refused() -> None: def test_private_voice_from_a_plugin_is_refused() -> None:
sc = _scenario() sc = _scenario()
assert _send(sc, _ars_ack(dst=213, call_type="group", frame_type=0)) is False assert _send(sc, _ars_ack(frame_type=HBPF_VOICE)) is False
assert sc.capture.packets == [] assert sc.capture.packets == [] and sc.protocols["SYSTEM-B"].sent_to_peer == []
def test_plugins_see_their_own_frames_as_synthetic() -> None: def test_plugins_see_their_own_frames_as_synthetic() -> None:
@ -113,3 +115,75 @@ def test_plugins_see_their_own_frames_as_synthetic() -> None:
_send(sc, _ars_ack()) _send(sc, _ars_ack())
frames = [e for e in events if isinstance(e, UnitDataFrame)] frames = [e for e in events if isinstance(e, UnitDataFrame)]
assert frames and all(f.context.is_synthetic for f in frames) 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…
Cancel
Save

Powered by TurnKey Linux.