feat(plugins): add drop-in plugin framework for voice and unit-data

Plugin bus, loader, voice/data bridges, post-forward events, circuit
breaker, SIGHUP rescan, OBP/HBP hooks. Reference skeleton plugins/example/
(disabled). Bilingual user guide.
pull/76/head
Rodrigo Pérez 4 weeks ago
parent 834034ea16
commit 0c5f1413e6

12
.gitignore vendored

@ -36,6 +36,18 @@ json/*
# Internal # Internal
docs-priv/ docs-priv/
# Runtime drop-in plugins (see plugins/README.md); example skeleton is versioned
plugins/*
!plugins/README.md
!plugins/.gitkeep
!plugins/example/
!plugins/example/**
plugins/example/example-events/
plugins/example/**/__pycache__/
# Runtime TTS on-demand cache (generated at runtime)
Audio/es_ES/ondemand/
# Dev-only session captures; anonymize before promoting to fixtures/ # Dev-only session captures; anonymize before promoting to fixtures/
tests/fixtures/sessions/_capture/ tests/fixtures/sessions/_capture/

@ -21,6 +21,7 @@ The **ADN DMR Peer Server** is a [GPL-3.0](https://www.gnu.org/licenses/gpl-3.0.
| I want to… | Start here | | I want to… | Start here |
|------------|------------| |------------|------------|
| Run and configure | [Introduction](server/user-guide/introduction.md), [Configuration](server/user-guide/configuration.md) | | Run and configure | [Introduction](server/user-guide/introduction.md), [Configuration](server/user-guide/configuration.md) |
| Write a drop-in plugin | [Plugins](server/user-guide/plugins.md), skeleton `plugins/example/` |
| TG 4000, 999x, echo | [Special numbers](server/user-guide/special-numbers.md) | | TG 4000, 999x, echo | [Special numbers](server/user-guide/special-numbers.md) |
| Private calls | [Private calls](server/user-guide/private-calls.md) | | Private calls | [Private calls](server/user-guide/private-calls.md) |
| Voice / TTS | [Voice, announcements, and TTS](server/user-guide/voice-and-tts.md) | | Voice / TTS | [Voice, announcements, and TTS](server/user-guide/voice-and-tts.md) |

@ -28,6 +28,8 @@ Conceptual comparison with legacy **`BRIDGES`**: [BRIDGES vs Subscriptions](brid
Performance changes in 2.x (indexes, reporting, integrated proxy): [Performance (2.x)](performance.md). Performance changes in 2.x (indexes, reporting, integrated proxy): [Performance (2.x)](performance.md).
**Drop-in plugins** (voice/unit-data bus, `ServerPlugin` contract): [Plugins user guide](../user-guide/plugins.md). Reference skeleton: `plugins/example/` in the repository root.
- **`InMemoryAclRouter`** (`AclRouter` port) — ACL range checks only (`acl_check`). - **`InMemoryAclRouter`** (`AclRouter` port) — ACL range checks only (`acl_check`).
- **`routing_table_for_report()`** — export shim for monitor/report (legacy BRIDGE_SND shape); not used for runtime forwards. - **`routing_table_for_report()`** — export shim for monitor/report (legacy BRIDGE_SND shape); not used for runtime forwards.

@ -0,0 +1,217 @@
# Plugins
Drop-in **plugins** extend the peer server without modifying core code. The framework lives in `src/adn_server/application/plugins/`; each plugin is a directory under `plugins/<name>/` loaded at runtime.
The repository ships one reference skeleton: `plugins/example/` (disabled by default).
---
## Directory layout
```
plugins/<plugin-name>/
config.yaml # enabled: true|false (+ options)
plugin/
__init__.py # must export create_plugin()
domain/ # pure logic, no I/O
application/ # use cases, ServerPlugin adapter
infrastructure/ # files, HTTP, thread pools
tests/
```
Enable a plugin in its `config.yaml` and reload the server:
```bash
systemctl reload adn-server # SIGHUP — rescans plugins/
```
---
## Server configuration (`PLUGINS`)
Optional block in **`adn-server.yaml`** (not in `adn-server.example.yaml`):
```yaml
PLUGINS:
directory: plugins # default: <project_root>/plugins
master_kill: false # true → unload all plugins
overrides:
my-plugin:
some_key: value # merged into plugin config at load/reload
```
| Key | Role |
|-----|------|
| `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` |
On **SIGHUP**, `PluginManager.rescan()` loads new plugins, unloads removed ones, and calls `on_reload()` when `config.yaml` changed.
### `config.yaml` reserved keys
The loader strips these before passing config to `on_load` / `on_reload`:
| Key | Role |
|-----|------|
| `enabled` | Must be `true` to load the plugin |
| `depends_on` | List of plugin names — topological load order |
| `hot_reload_seconds` | Reserved for future use |
---
## `ServerPlugin` contract
Defined in `src/adn_server/application/plugins/domain/protocol.py`:
| Method | Thread | Role |
|--------|--------|------|
| `name: str` | — | Plugin identifier (directory name) |
| `on_load(bus, config, server_ctx)` | reactor | Read config; subscribe to bus if needed |
| `on_event(event)` | **reactor** | Handle bus events — **O(1), no blocking I/O** |
| `on_reload(config)` | reactor | Optional hot-reload after `config.yaml` change |
| `on_shutdown()` | reactor | Flush buffers; stop background workers |
Factory entry point:
```python
# plugins/<name>/plugin/__init__.py
def create_plugin() -> ServerPlugin:
return MyPlugin()
```
---
## `ServerContext`
Passed to `on_load` as `server_ctx` (`application/plugins/application/context.py`):
| Field | Use |
|-------|-----|
| `config` | Full server config dict |
| `project_root` | Server install path |
| `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 |
**Pattern:** do only dispatch in `on_event`; call `defer_to_thread` for file writes, HTTP, heavy CPU.
---
## Event bus
`PluginBus` (`application/plugins/application/bus.py`) delivers events to every loaded plugin's `on_event`.
- Events are emitted **after voice/data has been forwarded** (`emit_deferred` — next reactor tick).
- Uncaught exceptions increment a per-plugin trip counter; after repeated failures the plugin is **disabled** and `on_shutdown()` is called (circuit breaker).
- `bus.subscribe(handler)` is available for internal handlers; plugins normally implement `on_event` only.
---
## Event types
Pure domain dataclasses in `application/plugins/domain/events.py`. Import in your plugin:
```python
from adn_server.application.plugins.domain.events import (
VoiceCallStart,
VoiceCallFrame,
VoiceCallEnd,
UnitDataStart,
UnitDataFrame,
UnitDataEnd,
)
```
| Event | When |
|-------|------|
| `VoiceCallStart` | Group or private voice call begins |
| `VoiceCallFrame` | One AMBE frame (`dmrpkt`, `frame_type`, `dtype_vseq`) |
| `VoiceCallEnd` | Call ends (`duration_s`, `frame_count`) |
| `UnitDataStart` | Unit-data session begins |
| `UnitDataFrame` | One data frame (`raw_data`, `data_label`, `seq`, `bits`, …) |
| `UnitDataEnd` | Unit-data session ends (`duration_s`, `packet_count`) |
Each event carries a **`CallLegContext`** (`event.context`) with metadata:
| Field | Meaning |
|-------|---------|
| `call_family` | `"GROUP"` or `"PRIVATE"` |
| `direction` | `"RX"` or `"TX"` |
| `origin_system` | Logical system name (bridge leg) |
| `system_mode` | `MASTER`, `PEER`, or `OPENBRIDGE` |
| `peer_id`, `src_id`, `dst_id`, `slot`, `stream_id` | DMR identifiers |
| `server_id` | Reporting server id |
| `is_synthetic`, `is_proxy_ingress` | Synthetic / proxy flags |
| `pkt_time` | Unix timestamp |
| `forwarded_systems` | Tuple of systems this leg was forwarded to |
| `obp_*`, `ber`, `rssi` | OpenBridge / RF metadata when present |
| `extra` | Additional dict (e.g. talker alias on end) |
Use `event.context.to_metadata_dict()` for JSON-serializable output.
Events originate from `VoicePluginBridge` and `DataPluginBridge`, hooked into routing after forward resolution.
---
## Clean architecture in plugins
Use the same inward dependency rule as the core server ([Architecture](../development/architecture.md)):
```mermaid
flowchart TD
subgraph plugin_pkg ["plugins/my-plugin/plugin/"]
impl["application/plugin_impl.py"]
uc["application/*_use_case.py"]
dom["domain/"]
inf["infrastructure/"]
impl --> uc --> dom
inf --> uc
end
coreEvents["adn_server.application.plugins.domain.events"]
impl --> coreEvents
```
| Layer | Responsibility |
|-------|----------------|
| `plugin/__init__.py` | Factory only — `create_plugin()` |
| `application/plugin_impl.py` | `ServerPlugin` adapter: `isinstance` checks, delegate to use cases |
| `application/` | Orchestration (use cases, session state) |
| `domain/` | Pure types and rules — no Twisted, no files, no sockets |
| `infrastructure/` | Writers, HTTP clients, pools — uses `defer_to_thread` |
---
## Reference example
Study `plugins/example/` in the repository root:
- `enabled: false` in committed `config.yaml` (no secrets).
- `DEBUG` log on every event.
- Writes one JSON file per session under `example-events/` on `VoiceCallEnd` / `UnitDataEnd`.
- Tests in `plugins/example/tests/`.
To try it: set `enabled: true`, reload, make a call, inspect `example-events/<stream_id>.json`.
---
## Tests
Plugin tests live next to the plugin, not in `adn-server/tests/`:
```bash
cd plugins/example
python3 -m pytest tests/ -q
```
Framework tests (`test_plugin_bus.py`, `test_voice_bridge.py`, …) remain in the core test suite.
---
## Best practices
1. **Never block the reactor** in `on_event` — delegate I/O and CPU-heavy work.
2. **Version `config.example.yaml`** alongside your plugin; keep `config.yaml` local and gitignored when it holds tokens.
3. **Catch errors** in background workers; uncaught exceptions in `on_event` trip the circuit breaker.
4. **Use `depends_on`** when one plugin must load after another.
5. Copy the **`example`** skeleton when starting a new plugin — rename the directory and implement your use case.

@ -20,6 +20,7 @@ El **ADN DMR Peer Server** es un puente de conferencia [GPL-3.0](https://www.gnu
| Quiero… | Empieza aquí | | Quiero… | Empieza aquí |
|---------|----------------| |---------|----------------|
| Ejecutar y configurar | [Introducción](server/user-guide/introduction.md), [Configuración](server/user-guide/configuration.md) | | Ejecutar y configurar | [Introducción](server/user-guide/introduction.md), [Configuración](server/user-guide/configuration.md) |
| Escribir un plugin drop-in | [Plugins](server/user-guide/plugins.md), skeleton `plugins/example/` |
| TG 4000, 999x, eco | [Números especiales](server/user-guide/special-numbers.md) | | TG 4000, 999x, eco | [Números especiales](server/user-guide/special-numbers.md) |
| Llamadas privadas | [Llamadas privadas](server/user-guide/private-calls.md) | | Llamadas privadas | [Llamadas privadas](server/user-guide/private-calls.md) |
| Voz / TTS | [Voz, anuncios y TTS](server/user-guide/voice-and-tts.md) | | Voz / TTS | [Voz, anuncios y TTS](server/user-guide/voice-and-tts.md) |

@ -28,6 +28,8 @@ Comparación conceptual con **`BRIDGES`** legacy: [BRIDGES vs Subscriptions](bri
Cambios de rendimiento en 2.x (índices, informes, proxy integrado): [Rendimiento (2.x)](performance.md). Cambios de rendimiento en 2.x (índices, informes, proxy integrado): [Rendimiento (2.x)](performance.md).
**Plugins drop-in** (bus de voz/unit-data, contrato `ServerPlugin`): [Guía de plugins](../user-guide/plugins.md). Skeleton de referencia: `plugins/example/` en la raíz del repositorio.
- **`InMemoryAclRouter`** (port `AclRouter`) — solo comprobaciones ACL (`acl_check`). - **`InMemoryAclRouter`** (port `AclRouter`) — solo comprobaciones ACL (`acl_check`).
- **`routing_table_for_report()`** — shim de exportación para monitor/informes (forma legacy BRIDGE_SND); no se usa para reenvíos en runtime. - **`routing_table_for_report()`** — shim de exportación para monitor/informes (forma legacy BRIDGE_SND); no se usa para reenvíos en runtime.

@ -0,0 +1,217 @@
# Plugins
Los **plugins** drop-in amplían el peer server sin modificar el código del núcleo. El framework está en `src/adn_server/application/plugins/`; cada plugin es un directorio bajo `plugins/<name>/` cargado en tiempo de ejecución.
El repositorio incluye un skeleton de referencia: `plugins/example/` (deshabilitado por defecto).
---
## Estructura de directorios
```
plugins/<plugin-name>/
config.yaml # enabled: true|false (+ opciones)
plugin/
__init__.py # debe exportar create_plugin()
domain/ # lógica pura, sin I/O
application/ # casos de uso, adaptador ServerPlugin
infrastructure/ # archivos, HTTP, pools de hilos
tests/
```
Activa un plugin en su `config.yaml` y recarga el servidor:
```bash
systemctl reload adn-server # SIGHUP — reescanea plugins/
```
---
## Configuración del servidor (`PLUGINS`)
Bloque opcional en **`adn-server.yaml`** (no está en `adn-server.example.yaml`):
```yaml
PLUGINS:
directory: plugins # por defecto: <project_root>/plugins
master_kill: false # true → descarga todos los plugins
overrides:
my-plugin:
some_key: value # se fusiona con la config del plugin
```
| Clave | Función |
|-------|---------|
| `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` |
Con **SIGHUP**, `PluginManager.rescan()` carga plugins nuevos, descarga los eliminados y llama `on_reload()` si cambió `config.yaml`.
### Claves reservadas en `config.yaml`
El loader las elimina antes de pasar la config a `on_load` / `on_reload`:
| Clave | Función |
|-------|---------|
| `enabled` | Debe ser `true` para cargar el plugin |
| `depends_on` | Lista de nombres — orden topológico de carga |
| `hot_reload_seconds` | Reservado para uso futuro |
---
## Contrato `ServerPlugin`
Definido en `src/adn_server/application/plugins/domain/protocol.py`:
| Método | Hilo | Función |
|--------|------|---------|
| `name: str` | — | Identificador (nombre del directorio) |
| `on_load(bus, config, server_ctx)` | reactor | Leer config; suscribirse al bus si hace falta |
| `on_event(event)` | **reactor** | Manejar eventos — **O(1), sin I/O bloqueante** |
| `on_reload(config)` | reactor | Hot-reload opcional tras cambio de `config.yaml` |
| `on_shutdown()` | reactor | Vaciar buffers; parar workers |
Punto de entrada:
```python
# plugins/<name>/plugin/__init__.py
def create_plugin() -> ServerPlugin:
return MyPlugin()
```
---
## `ServerContext`
Se pasa a `on_load` como `server_ctx` (`application/plugins/application/context.py`):
| Campo | Uso |
|-------|-----|
| `config` | Dict de configuración completa del servidor |
| `project_root` | Ruta de instalación |
| `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 |
**Patrón:** solo despachar en `on_event`; usar `defer_to_thread` para archivos, HTTP o CPU intensiva.
---
## Bus de eventos
`PluginBus` (`application/plugins/application/bus.py`) entrega eventos al `on_event` de cada plugin cargado.
- Los eventos se emiten **después del forward** de voz/datos (`emit_deferred` — siguiente tick del reactor).
- Excepciones no capturadas incrementan un contador; tras repetir fallos el plugin se **deshabilita** y se llama `on_shutdown()` (circuit breaker).
- `bus.subscribe(handler)` existe para handlers internos; los plugins suelen usar solo `on_event`.
---
## Tipos de evento
Dataclasses puras en `application/plugins/domain/events.py`. Importar en el plugin:
```python
from adn_server.application.plugins.domain.events import (
VoiceCallStart,
VoiceCallFrame,
VoiceCallEnd,
UnitDataStart,
UnitDataFrame,
UnitDataEnd,
)
```
| Evento | Cuándo |
|--------|--------|
| `VoiceCallStart` | Inicio de llamada de voz (grupo o privada) |
| `VoiceCallFrame` | Un frame AMBE (`dmrpkt`, `frame_type`, `dtype_vseq`) |
| `VoiceCallEnd` | Fin de llamada (`duration_s`, `frame_count`) |
| `UnitDataStart` | Inicio de sesión unit-data |
| `UnitDataFrame` | Un frame de datos (`raw_data`, `data_label`, `seq`, `bits`, …) |
| `UnitDataEnd` | Fin de sesión unit-data (`duration_s`, `packet_count`) |
Cada evento incluye **`CallLegContext`** (`event.context`):
| Campo | Significado |
|-------|-------------|
| `call_family` | `"GROUP"` o `"PRIVATE"` |
| `direction` | `"RX"` o `"TX"` |
| `origin_system` | Nombre del sistema lógico (pata del bridge) |
| `system_mode` | `MASTER`, `PEER` u `OPENBRIDGE` |
| `peer_id`, `src_id`, `dst_id`, `slot`, `stream_id` | Identificadores DMR |
| `server_id` | Id de servidor para informes |
| `is_synthetic`, `is_proxy_ingress` | Flags sintético / proxy |
| `pkt_time` | Timestamp Unix |
| `forwarded_systems` | Sistemas a los que se reenvió esta pata |
| `obp_*`, `ber`, `rssi` | Metadatos OpenBridge / RF si aplica |
| `extra` | Dict adicional (p. ej. talker alias al finalizar) |
Usa `event.context.to_metadata_dict()` para salida JSON.
Los eventos provienen de `VoicePluginBridge` y `DataPluginBridge`, enganchados al routing tras el forward.
---
## Clean architecture en plugins
Misma regla de dependencias hacia dentro que el núcleo ([Arquitectura](../development/architecture.md)):
```mermaid
flowchart TD
subgraph plugin_pkg ["plugins/my-plugin/plugin/"]
impl["application/plugin_impl.py"]
uc["application/*_use_case.py"]
dom["domain/"]
inf["infrastructure/"]
impl --> uc --> dom
inf --> uc
end
coreEvents["adn_server.application.plugins.domain.events"]
impl --> coreEvents
```
| Capa | Responsabilidad |
|------|-----------------|
| `plugin/__init__.py` | Solo factory — `create_plugin()` |
| `application/plugin_impl.py` | Adaptador `ServerPlugin`: `isinstance`, delegar a use cases |
| `application/` | Orquestación (casos de uso, estado de sesión) |
| `domain/` | Tipos y reglas puras — sin Twisted, archivos ni sockets |
| `infrastructure/` | Writers, clientes HTTP, pools — usa `defer_to_thread` |
---
## Ejemplo de referencia
Estudia `plugins/example/` en la raíz del repositorio:
- `enabled: false` en el `config.yaml` versionado (sin secretos).
- Log `DEBUG` en cada evento.
- Escribe un JSON por sesión en `example-events/` al `VoiceCallEnd` / `UnitDataEnd`.
- Tests en `plugins/example/tests/`.
Para probarlo: `enabled: true`, reload, haz una llamada, revisa `example-events/<stream_id>.json`.
---
## Tests
Los tests del plugin van junto al plugin, no en `adn-server/tests/`:
```bash
cd plugins/example
python3 -m pytest tests/ -q
```
Los tests del framework (`test_plugin_bus.py`, `test_voice_bridge.py`, …) permanecen en el núcleo.
---
## Buenas prácticas
1. **No bloquear el reactor** en `on_event` — delegar I/O y trabajo pesado.
2. **Versionar `config.example.yaml`** junto al plugin; mantener `config.yaml` local y gitignored si lleva tokens.
3. **Capturar errores** en workers en segundo plano; excepciones en `on_event` activan el circuit breaker.
4. Usar **`depends_on`** cuando un plugin deba cargarse después de otro.
5. Copiar el skeleton **`example`** al crear un plugin nuevo — renombrar el directorio e implementar tu caso de uso.

@ -55,6 +55,7 @@ nav:
- Guía de usuario: - Guía de usuario:
- Introducción: server/user-guide/introduction.md - Introducción: server/user-guide/introduction.md
- Configuración: server/user-guide/configuration.md - Configuración: server/user-guide/configuration.md
- Plugins: server/user-guide/plugins.md
- Talker Alias: server/user-guide/talker-alias.md - Talker Alias: server/user-guide/talker-alias.md
- Bridges y talkgroups: server/user-guide/bridges-and-talkgroups.md - Bridges y talkgroups: server/user-guide/bridges-and-talkgroups.md
- Números especiales (4000, 999x, eco): server/user-guide/special-numbers.md - Números especiales (4000, 999x, eco): server/user-guide/special-numbers.md

@ -55,6 +55,7 @@ nav:
- User guide: - User guide:
- Introduction: server/user-guide/introduction.md - Introduction: server/user-guide/introduction.md
- Configuration: server/user-guide/configuration.md - Configuration: server/user-guide/configuration.md
- Plugins: server/user-guide/plugins.md
- Talker Alias: server/user-guide/talker-alias.md - Talker Alias: server/user-guide/talker-alias.md
- Bridges and talkgroups: server/user-guide/bridges-and-talkgroups.md - Bridges and talkgroups: server/user-guide/bridges-and-talkgroups.md
- Special numbers (4000, 999x, echo): server/user-guide/special-numbers.md - Special numbers (4000, 999x, echo): server/user-guide/special-numbers.md

@ -0,0 +1,36 @@
# Plugins
Drop-in extensions live under `plugins/<name>/`. The server loads them at runtime from this directory (default: `<project_root>/plugins`).
## Reference skeleton
[`example/`](example/) is the only plugin versioned in this repository. It is **disabled by default** (`enabled: false` in `config.yaml`). Use it as a minimal, working template (clean architecture layout, event handling, JSON output on call end).
To try it locally, set `enabled: true` in `plugins/example/config.yaml` and reload the server (`systemctl reload adn-server`).
## Layout
```
plugins/<plugin-name>/
config.yaml # enabled: true|false (+ plugin options)
plugin/
__init__.py # create_plugin()
domain/
application/
infrastructure/
tests/
```
## Documentation
Full framework guide (contract, events, threading, `PLUGINS` server config):
- English: [Plugins (user guide)](../docs/en/server/user-guide/plugins.md)
- Spanish: [Plugins (guía de usuario)](../docs/es/server/user-guide/plugins.md)
Tests for each plugin run from that plugin's directory:
```bash
cd plugins/example
python3 -m pytest tests/ -q
```

@ -0,0 +1,28 @@
# Example plugin
Minimal **reference skeleton** for adn-server drop-in plugins. Disabled by default.
## What it does
- On every bus event: `DEBUG` log line `(EXAMPLE) …`.
- On `VoiceCallEnd` / `UnitDataEnd`: writes one JSON file under `example-events/` with call metadata, duration, and frame/packet counts.
No HTTP, no filters, no secrets — study or copy this tree when authoring your own plugin.
## Enable locally
```bash
# Edit plugins/example/config.yaml → enabled: true
systemctl reload adn-server # or SIGHUP
```
After a voice or unit-data session, check `example-events/<stream_id>.json` under the server project root.
## Tests
```bash
cd plugins/example
python3 -m pytest tests/ -q
```
See also the [Plugins user guide](../../docs/en/server/user-guide/plugins.md).

@ -0,0 +1,4 @@
enabled: false
# Relative to adn-server project root (created on write if missing).
output_dir: example-events

@ -0,0 +1,29 @@
# ADN DMR Peer Server plugin - example factory
#
# 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
###############################################################################
"""Example plugin — drop-in factory."""
from __future__ import annotations
from .application.plugin_impl import ExamplePlugin
def create_plugin() -> ExamplePlugin:
return ExamplePlugin()

@ -0,0 +1,129 @@
# ADN DMR Peer Server plugin - example adapter
#
# 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
###############################################################################
"""ServerPlugin adapter — debug log on every event, JSON file on session end."""
from __future__ import annotations
import logging
from typing import Any
from adn_server.application.plugins.domain.events import (
UnitDataEnd,
UnitDataFrame,
UnitDataStart,
VoiceCallEnd,
VoiceCallFrame,
VoiceCallStart,
)
from ..infrastructure.json_writer import JsonWriter
from .session_use_case import SessionUseCase
logger = logging.getLogger(__name__)
class ExamplePlugin:
name = "example"
def __init__(self) -> None:
self._use_case: SessionUseCase | None = None
self._writer: JsonWriter | None = None
def on_load(self, bus: Any, config: dict[str, Any], server_ctx: Any) -> None:
self._writer = JsonWriter(
project_root=server_ctx.project_root,
output_dir=str(config.get("output_dir", "example-events")),
defer_to_thread=server_ctx.defer_to_thread,
)
self._use_case = SessionUseCase(self._writer)
logger.debug("(EXAMPLE) loaded output_dir=%s", config.get("output_dir"))
def on_reload(self, config: dict[str, Any]) -> None:
if self._writer is not None:
self._writer.set_output_dir(str(config.get("output_dir", "example-events")))
if self._use_case is not None:
self._use_case.shutdown()
logger.debug("(EXAMPLE) reloaded output_dir=%s", config.get("output_dir"))
def on_shutdown(self) -> None:
if self._use_case is not None:
self._use_case.shutdown()
logger.debug("(EXAMPLE) shutdown")
def on_event(self, event: object) -> None:
if self._use_case is None:
return
if isinstance(event, VoiceCallStart):
meta = event.context.to_metadata_dict()
logger.debug("(EXAMPLE) VoiceCallStart stream_id=%s", meta.get("stream_id"))
self._use_case.on_start("voice", meta)
elif isinstance(event, VoiceCallFrame):
meta = event.context.to_metadata_dict()
logger.debug(
"(EXAMPLE) VoiceCallFrame stream_id=%s frame_type=%s",
meta.get("stream_id"),
event.frame_type,
)
self._use_case.on_frame(meta)
elif isinstance(event, VoiceCallEnd):
meta = event.context.to_metadata_dict()
logger.debug(
"(EXAMPLE) VoiceCallEnd stream_id=%s duration_s=%.3f frames=%s",
meta.get("stream_id"),
event.duration_s,
event.frame_count,
)
self._use_case.on_end(
"voice",
meta,
duration_s=event.duration_s,
count=event.frame_count,
)
elif isinstance(event, UnitDataStart):
meta = event.context.to_metadata_dict()
logger.debug(
"(EXAMPLE) UnitDataStart stream_id=%s label=%s",
meta.get("stream_id"),
event.data_label,
)
self._use_case.on_start("unit_data", meta)
elif isinstance(event, UnitDataFrame):
meta = event.context.to_metadata_dict()
logger.debug(
"(EXAMPLE) UnitDataFrame stream_id=%s seq=%s",
meta.get("stream_id"),
event.seq,
)
self._use_case.on_frame(meta)
elif isinstance(event, UnitDataEnd):
meta = event.context.to_metadata_dict()
logger.debug(
"(EXAMPLE) UnitDataEnd stream_id=%s duration_s=%.3f packets=%s",
meta.get("stream_id"),
event.duration_s,
event.packet_count,
)
self._use_case.on_end(
"unit_data",
meta,
duration_s=event.duration_s,
count=event.packet_count,
)

@ -0,0 +1,93 @@
# ADN DMR Peer Server plugin - example session use case
#
# 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
###############################################################################
"""Track voice/data sessions and emit JSON at end."""
from __future__ import annotations
import time
from dataclasses import dataclass, field
from typing import Any
from ..domain.record import build_end_record, session_key
from ..infrastructure.json_writer import JsonWriter
@dataclass
class _Session:
kind: str
meta: dict[str, Any]
started_at: float
count: int = 0
class SessionUseCase:
def __init__(self, writer: JsonWriter) -> None:
self._writer = writer
self._active: dict[tuple[str, int], _Session] = {}
def on_start(self, kind: str, meta: dict[str, Any]) -> None:
key = session_key(meta)
self._active[key] = _Session(
kind=kind,
meta=dict(meta),
started_at=float(meta.get("pkt_time") or time.time()),
)
def on_frame(self, meta: dict[str, Any]) -> None:
key = session_key(meta)
session = self._active.get(key)
if session is not None:
session.count += 1
def on_end(
self,
kind: str,
end_meta: dict[str, Any],
*,
duration_s: float,
count: int | None = None,
) -> None:
key = session_key(end_meta)
session = self._active.pop(key, None)
if session is None:
started_at = float(end_meta.get("pkt_time") or time.time()) - duration_s
start_meta = dict(end_meta)
frame_count = count if count is not None else 0
else:
started_at = session.started_at
start_meta = session.meta
frame_count = count if count is not None else session.count
ended_at = float(end_meta.get("pkt_time") or time.time())
record = build_end_record(
kind=kind,
start_meta=start_meta,
end_meta=end_meta,
duration_s=duration_s,
count=frame_count,
started_at=started_at,
ended_at=ended_at,
)
stream_id = int(end_meta.get("stream_id") or 0)
self._writer.write_record(stream_id, record)
def shutdown(self) -> None:
self._active.clear()

@ -0,0 +1,58 @@
# ADN DMR Peer Server plugin - example record builder
#
# 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
###############################################################################
"""Pure helpers to build JSON records at session end."""
from __future__ import annotations
from datetime import datetime, timezone
from typing import Any
def session_key(meta: dict[str, Any]) -> tuple[str, int]:
return (str(meta.get("origin_system") or ""), int(meta.get("stream_id") or 0))
def _iso_utc(ts: float) -> str:
return datetime.fromtimestamp(ts, tz=timezone.utc).isoformat().replace("+00:00", "Z")
def build_end_record(
*,
kind: str,
start_meta: dict[str, Any],
end_meta: dict[str, Any],
duration_s: float,
count: int,
started_at: float,
ended_at: float,
) -> dict[str, Any]:
"""Flat JSON document written when a voice or unit-data session ends."""
out = {**start_meta, **end_meta}
out["event_kind"] = kind
out["duration_s"] = duration_s
out["started_at"] = _iso_utc(started_at)
out["ended_at"] = _iso_utc(ended_at)
out["pkt_time"] = float(end_meta.get("pkt_time") or ended_at)
if kind == "voice":
out["frame_count"] = count
else:
out["packet_count"] = count
return out

@ -0,0 +1,56 @@
# ADN DMR Peer Server plugin - example JSON writer
#
# 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
###############################################################################
"""Write session JSON files off the reactor thread."""
from __future__ import annotations
import json
import logging
from pathlib import Path
from typing import Any, Callable
logger = logging.getLogger(__name__)
class JsonWriter:
def __init__(
self,
*,
project_root: str,
output_dir: str,
defer_to_thread: Callable[..., Any],
) -> None:
self._root = Path(project_root)
self._output_dir = output_dir
self._defer = defer_to_thread
def set_output_dir(self, output_dir: str) -> None:
self._output_dir = output_dir
def write_record(self, stream_id: int, record: dict[str, Any]) -> None:
self._defer(self._write_sync, stream_id, record)
def _write_sync(self, stream_id: int, record: dict[str, Any]) -> None:
out_dir = self._root / self._output_dir
out_dir.mkdir(parents=True, exist_ok=True)
path = out_dir / f"{stream_id}.json"
path.write_text(json.dumps(record, indent=2, sort_keys=True) + "\n", encoding="utf-8")
logger.debug("(EXAMPLE) wrote %s", path)

@ -0,0 +1,28 @@
# ADN DMR Peer Server plugin - example test path
#
# 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
###############################################################################
from __future__ import annotations
import sys
from pathlib import Path
_ROOT = Path(__file__).resolve().parents[1]
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))

@ -0,0 +1,48 @@
# ADN DMR Peer Server plugin - example tests
#
# 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
###############################################################################
from __future__ import annotations
from plugin.domain.record import build_end_record, session_key
def test_session_key() -> None:
assert session_key({"origin_system": "SYS1", "stream_id": 42}) == ("SYS1", 42)
def test_build_voice_end_record() -> None:
start = {"origin_system": "SYS1", "stream_id": 1, "src_id": 100, "pkt_time": 1000.0}
end = {"origin_system": "SYS1", "stream_id": 1, "dst_id": 200, "pkt_time": 1005.0}
rec = build_end_record(
kind="voice",
start_meta=start,
end_meta=end,
duration_s=5.0,
count=12,
started_at=1000.0,
ended_at=1005.0,
)
assert rec["event_kind"] == "voice"
assert rec["frame_count"] == 12
assert rec["duration_s"] == 5.0
assert rec["src_id"] == 100
assert rec["dst_id"] == 200
assert rec["started_at"].endswith("Z")
assert rec["ended_at"].endswith("Z")

@ -0,0 +1,72 @@
# ADN DMR Peer Server plugin - example session tests
#
# 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
###############################################################################
from __future__ import annotations
import sys
from pathlib import Path
_ROOT = Path(__file__).resolve().parents[1]
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))
from plugin.application.session_use_case import SessionUseCase
from plugin.infrastructure.json_writer import JsonWriter
class _SyncWriter:
def __init__(self) -> None:
self.records: list[tuple[int, dict]] = []
def write_record(self, stream_id: int, record: dict) -> None:
self.records.append((stream_id, record))
def test_voice_session_end_writes_record() -> None:
sink = _SyncWriter()
uc = SessionUseCase(sink) # type: ignore[arg-type]
meta = {"origin_system": "T1", "stream_id": 99, "src_id": 1, "pkt_time": 10.0}
uc.on_start("voice", meta)
uc.on_frame(meta)
uc.on_frame(meta)
uc.on_end("voice", {**meta, "dst_id": 4000}, duration_s=2.5, count=2)
assert len(sink.records) == 1
sid, rec = sink.records[0]
assert sid == 99
assert rec["event_kind"] == "voice"
assert rec["frame_count"] == 2
assert rec["duration_s"] == 2.5
def test_json_writer_delegates_to_thread(tmp_path: Path) -> None:
written: list[Path] = []
def defer(fn, *args): # noqa: ANN001
fn(*args)
written.append(tmp_path / "example-events" / f"{args[0]}.json")
writer = JsonWriter(
project_root=str(tmp_path),
output_dir="example-events",
defer_to_thread=defer,
)
writer.write_record(7, {"stream_id": 7, "event_kind": "voice"})
assert len(written) == 1
assert written[0].is_file()

@ -0,0 +1,41 @@
# ADN DMR Peer Server - plugin system package
#
# 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
###############################################################################
"""Plugin system for adn-server (bus, loader, voice and data events)."""
from .application.manager import PluginManager
from .domain.events import (
UnitDataEnd,
UnitDataFrame,
UnitDataStart,
VoiceCallEnd,
VoiceCallFrame,
VoiceCallStart,
)
__all__ = [
"PluginManager",
"UnitDataEnd",
"UnitDataFrame",
"UnitDataStart",
"VoiceCallEnd",
"VoiceCallFrame",
"VoiceCallStart",
]

@ -0,0 +1,25 @@
# ADN DMR Peer Server - plugin application package
#
# 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
###############################################################################
from .bus import PluginBus
from .context import ServerContext
from .manager import PluginManager
__all__ = ["PluginBus", "PluginManager", "ServerContext"]

@ -0,0 +1,100 @@
# ADN DMR Peer Server - plugin bridge helpers
#
# 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
###############################################################################
"""Shared helpers for voice/data plugin bridges."""
from __future__ import annotations
from collections.abc import Callable
from typing import Any
from adn_server.domain import bytes_4, int_id
UNIT_DATA_DTYPE_LABELS: dict[int, str] = {
3: "UNIT CSBK",
6: "UNIT DATA HEADER",
7: "UNIT VCSBK 1/2 DATA BLOCK",
8: "UNIT VCSBK 3/4 DATA BLOCK",
}
def unit_data_label(dtype_vseq: int) -> str:
return UNIT_DATA_DTYPE_LABELS.get(dtype_vseq, "UNIT DATA")
def int_byte(val: bytes | int | None) -> int | None:
if val is None:
return None
if isinstance(val, int):
return val
if isinstance(val, bytes) and val:
return int.from_bytes(val, "big")
return None
def alias_extra(
config: dict[str, Any],
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
) -> dict[str, Any]:
peer_ids = config.get("_PEER_IDS") or {}
sub_ids = config.get("_SUB_IDS") or {}
tg_ids = config.get("_TG_IDS") or {}
server_ids = config.get("_SERVER_IDS") or {}
extra: dict[str, Any] = {}
pid = int_id(peer_id)
sid = int_id(rf_src)
did = int_id(dst_id)
if pid in peer_ids:
extra["peer_callsign"] = peer_ids[pid]
if sid in sub_ids:
extra["src_callsign"] = sub_ids[sid]
if did in tg_ids:
extra["dst_name"] = tg_ids[did]
srv = int_byte(config.get("GLOBAL", {}).get("SERVER_ID"))
if srv is not None:
server_name = server_ids.get(str(srv))
if server_name:
extra["server_name"] = server_name
return extra
def stream_talker_alias(
config: dict[str, Any],
*,
origin_system: str,
stream_id: int,
rf_src: int,
get_dmra_blocks: Callable[[str, bytes], dict[int, bytes] | None] | None,
) -> str | None:
"""Decoded Talker Alias for a voice stream, when buffered on the origin system."""
if get_dmra_blocks is None:
return None
blocks = get_dmra_blocks(origin_system, bytes_4(stream_id))
if not blocks:
return None
from adn_server.application.talker_alias_use_cases import passthrough_complete
from adn_server.domain.talker_alias import decode_ta_from_blocks
if not passthrough_complete(blocks):
return None
text = decode_ta_from_blocks(blocks)
return text or None

@ -0,0 +1,110 @@
# ADN DMR Peer Server - plugin event bus
#
# 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
###############################################################################
"""In-process plugin event bus with time budget and circuit breaker."""
from __future__ import annotations
import logging
from collections.abc import Callable
from typing import Any
from ..domain.protocol import ServerPlugin
logger = logging.getLogger(__name__)
DEFAULT_BUDGET_S = 0.0001 # 100 µs per on_event
class PluginBus:
def __init__(
self,
*,
budget_s: float = DEFAULT_BUDGET_S,
max_trips: int = 5,
) -> None:
self._handlers: list[Callable[[object], None]] = []
self._plugins: list[ServerPlugin] = []
self._budget_s = budget_s
self._max_trips = max_trips
self._trip_counts: dict[str, int] = {}
self._disabled: set[str] = set()
self._call_later: Callable[..., Any] | None = None
def set_call_later(self, call_later: Callable[..., Any]) -> None:
self._call_later = call_later
def register_plugin(self, plugin: ServerPlugin) -> None:
self._disabled.discard(plugin.name)
self._trip_counts.pop(plugin.name, None)
if plugin in self._plugins:
return
self._plugins.append(plugin)
def unregister_plugin(self, plugin: ServerPlugin) -> None:
self._plugins = [p for p in self._plugins if p is not plugin]
def subscribe(self, handler: Callable[[object], None]) -> None:
self._handlers.append(handler)
def has_subscribers(self) -> bool:
return bool(self._plugins or self._handlers)
def emit(self, event: object) -> None:
if not self.has_subscribers():
return
for handler in self._handlers:
try:
handler(event)
except Exception:
logger.exception("(PLUGIN-BUS) handler failed")
for plugin in list(self._plugins):
if plugin.name in self._disabled:
continue
try:
plugin.on_event(event)
except Exception:
logger.exception("(PLUGIN-BUS) plugin %s failed", plugin.name)
self._trip(plugin)
def emit_deferred(self, event: object) -> None:
"""Schedule emit on next reactor tick (after forward)."""
if not self.has_subscribers():
return
if self._call_later is not None:
self._call_later(0, self.emit, event)
else:
self.emit(event)
def _trip(self, plugin: ServerPlugin) -> None:
count = self._trip_counts.get(plugin.name, 0) + 1
self._trip_counts[plugin.name] = count
if count >= self._max_trips:
logger.error(
"(PLUGIN-BUS) disabling plugin %s after %d exception trip(s)",
plugin.name,
count,
)
self._disabled.add(plugin.name)
try:
plugin.on_shutdown()
except Exception:
logger.exception("(PLUGIN-BUS) shutdown failed for %s", plugin.name)
self.unregister_plugin(plugin)

@ -0,0 +1,35 @@
# ADN DMR Peer Server - plugin server context
#
# 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
###############################################################################
"""Capabilities exposed to plugins (application layer)."""
from __future__ import annotations
from dataclasses import dataclass
from typing import Any, Callable
@dataclass
class ServerContext:
config: dict[str, Any]
project_root: str
defer_to_thread: Callable[..., Any]
call_from_reactor: Callable[..., Any]
call_later: Callable[..., Any]

@ -0,0 +1,162 @@
# ADN DMR Peer Server - plugin unit data event bridge
#
# 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
###############################################################################
"""Bridge unit-data routing → plugin bus (after forward)."""
from __future__ import annotations
from typing import Any
from adn_server.domain import int_id
from adn_server.application.proxy.deployment import is_proxy_inject_only
from ..application.bus import PluginBus
from ..domain.events import CallLegContext, UnitDataEnd, UnitDataFrame, UnitDataStart
from .bridge_common import alias_extra, int_byte, unit_data_label
class DataPluginBridge:
"""Emits unit-data events to PluginBus after routing forward."""
def __init__(self, bus: PluginBus, config: dict[str, Any]) -> None:
self._bus = bus
self._config = config
# (origin_system, slot) -> (stream_id, start_pkt_time)
self._active_streams: dict[tuple[str, int], tuple[int, float]] = {}
def update_config(self, config: dict[str, Any]) -> None:
self._config = config
def notify_after_forward(
self,
*,
route: str,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
seq: int,
slot: int,
frame_type: int,
dtype_vseq: int,
stream_id: bytes,
data: bytes,
pkt_time: float,
source_is_obp: bool,
obp_hops: bytes,
obp_source_server: bytes | None,
obp_ber: bytes,
obp_rssi: bytes,
obp_source_rptr: bytes,
forwarded: list[str],
) -> None:
if not self._bus.has_subscribers():
return
systems_cfg = self._config.get("SYSTEMS", {})
mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER")
server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0
bits = data[15] if len(data) > 15 else 0
dmrpkt = data[20:53] if len(data) >= 53 else b""
sid = int_id(stream_id)
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])
is_new_stream = prev is None or prev[0] != sid
if is_new_stream:
self._active_streams[stream_key] = (sid, pkt_time)
ctx = CallLegContext(
call_family="DATA",
direction="RX",
origin_system=system_name,
system_mode=str(mode),
peer_id=int_id(peer_id),
src_id=int_id(rf_src),
dst_id=int_id(dst_id),
slot=slot,
stream_id=sid,
server_id=server_id,
is_synthetic=False,
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,
obp_hops=int_byte(obp_hops) if source_is_obp else None,
obp_source_rptr_id=int_byte(obp_source_rptr) if source_is_obp else None,
ber=int_byte(obp_ber),
rssi=int_byte(obp_rssi),
forwarded_systems=tuple(forwarded),
extra={
**alias_extra(self._config, peer_id, rf_src, dst_id),
"route": route,
"data_label": unit_data_label(dtype_vseq),
"dtype_vseq": dtype_vseq,
},
)
label = unit_data_label(dtype_vseq)
if is_new_stream or dtype_vseq == 6:
self._bus.emit(
UnitDataStart(
context=ctx,
dtype_vseq=dtype_vseq,
data_label=label,
seq=seq,
frame_type=frame_type,
)
)
frame = UnitDataFrame(
context=ctx,
dmrpkt=dmrpkt,
raw_data=data,
dtype_vseq=dtype_vseq,
data_label=label,
seq=seq,
frame_type=frame_type,
bits=bits,
)
self._bus.emit_deferred(frame)
def _emit_end(
self,
system_name: str,
slot: int,
pkt_time: float,
stream_id: int,
start_time: float,
) -> None:
duration = max(0.0, pkt_time - start_time)
systems_cfg = self._config.get("SYSTEMS", {})
mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER")
server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0
ctx = CallLegContext(
call_family="DATA",
direction="RX",
origin_system=system_name,
system_mode=str(mode),
peer_id=0,
src_id=0,
dst_id=0,
slot=slot,
stream_id=stream_id,
server_id=server_id,
is_synthetic=False,
is_proxy_ingress=is_proxy_inject_only(self._config, system_name),
pkt_time=pkt_time,
)
self._bus.emit(UnitDataEnd(context=ctx, duration_s=duration))

@ -0,0 +1,143 @@
# ADN DMR Peer Server - plugin manager
#
# 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
###############################################################################
"""Plugin lifecycle: discover, load, rescan, shutdown."""
from __future__ import annotations
import logging
from pathlib import Path
from typing import Any
from ..application.bus import PluginBus
from ..application.context import ServerContext
from ..domain.protocol import ServerPlugin
from ..infrastructure.loader import (
load_plugin_from_entry,
plugin_config_for_load,
scan_plugin_dirs,
topo_sort,
)
logger = logging.getLogger(__name__)
class PluginManager:
def __init__(
self,
bus: PluginBus,
ctx: ServerContext,
project_root: str | Path,
) -> None:
self._bus = bus
self._ctx = ctx
self._project_root = Path(project_root)
self._loaded: dict[str, ServerPlugin] = {}
self._configs: dict[str, dict[str, Any]] = {}
@property
def bus(self) -> PluginBus:
return self._bus
def plugins_directory(self, server_config: dict[str, Any]) -> Path:
plugins_cfg = server_config.get("PLUGINS") or {}
if plugins_cfg.get("master_kill"):
return self._project_root / "__disabled__"
custom = plugins_cfg.get("directory")
if custom:
p = Path(custom)
return p if p.is_absolute() else self._project_root / p
return self._project_root / "plugins"
def discover_and_load(self, server_config: dict[str, Any]) -> list[ServerPlugin]:
directory = self.plugins_directory(server_config)
overrides = (server_config.get("PLUGINS") or {}).get("overrides") or {}
entries = topo_sort(scan_plugin_dirs(directory))
loaded: list[ServerPlugin] = []
for entry in entries:
if entry.name in self._loaded:
continue
try:
plugin = load_plugin_from_entry(entry, directory, overrides=overrides)
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)
self._loaded[entry.name] = plugin
self._configs[entry.name] = dict(entry.config)
loaded.append(plugin)
self._bus.register_plugin(plugin)
logger.debug("(PLUGIN-MANAGER) loaded %s", entry.name)
except Exception:
logger.exception("(PLUGIN-MANAGER) failed to load %s", entry.name)
return loaded
def rescan(self, server_config: dict[str, Any]) -> None:
directory = self.plugins_directory(server_config)
if (server_config.get("PLUGINS") or {}).get("master_kill"):
self.shutdown_all()
return
active = {e.name: e for e in topo_sort(scan_plugin_dirs(directory))}
for name in list(self._loaded):
if name not in active:
self._unload_one(name)
overrides = (server_config.get("PLUGINS") or {}).get("overrides") or {}
for entry in topo_sort(scan_plugin_dirs(directory)):
if entry.name not in self._loaded:
try:
plugin = load_plugin_from_entry(entry, directory, overrides=overrides)
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)
self._loaded[entry.name] = plugin
self._configs[entry.name] = dict(entry.config)
self._bus.register_plugin(plugin)
logger.debug("(PLUGIN-MANAGER) loaded %s", entry.name)
except Exception:
logger.exception("(PLUGIN-MANAGER) failed to load %s", entry.name)
continue
new_cfg = entry.config
if new_cfg != self._configs.get(entry.name):
plugin = self._loaded[entry.name]
cfg = plugin_config_for_load(new_cfg)
if overrides.get(entry.name):
cfg = {**cfg, **overrides[entry.name]}
try:
plugin.on_reload(cfg)
except Exception:
logger.exception("(PLUGIN-MANAGER) reload failed for %s", entry.name)
self._configs[entry.name] = dict(new_cfg)
def shutdown_all(self) -> None:
for name in list(self._loaded):
self._unload_one(name)
def _unload_one(self, name: str) -> None:
plugin = self._loaded.pop(name, None)
self._configs.pop(name, None)
if plugin is None:
return
self._bus.unregister_plugin(plugin)
try:
plugin.on_shutdown()
except Exception:
logger.exception("(PLUGIN-MANAGER) shutdown failed for %s", name)
logger.info("(PLUGIN-MANAGER) unloaded %s", name)

@ -0,0 +1,345 @@
# ADN DMR Peer Server - plugin voice event bridge
#
# 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
###############################################################################
"""Bridge routing → plugin bus (after forward)."""
from __future__ import annotations
from collections.abc import Callable
from dataclasses import replace
from typing import Any
from adn_server.domain import HBPF_VOICE, HBPF_VOICE_SYNC, int_id
from adn_server.application.proxy.deployment import is_proxy_inject_only
from ..application.bus import PluginBus
from ..domain.events import CallLegContext, VoiceCallEnd, VoiceCallFrame, VoiceCallStart
from .bridge_common import alias_extra, int_byte, stream_talker_alias
class VoicePluginBridge:
def __init__(
self,
bus: PluginBus,
config: dict[str, Any],
*,
get_dmra_blocks: Callable[[str, bytes], dict[int, bytes] | None] | None = None,
) -> None:
self._bus = bus
self._config = config
self._get_dmra_blocks = get_dmra_blocks
self._plugin_started: set[tuple[str, int]] = set()
self._plugin_ended: set[tuple[str, int]] = set()
def update_config(self, config: dict[str, Any]) -> None:
self._config = config
def has_subscribers(self) -> bool:
return self._bus.has_subscribers()
def _end_extra(
self,
*,
origin_system: str,
stream_id: int,
src_id: int,
base_extra: dict[str, Any],
) -> dict[str, Any]:
extra = dict(base_extra)
ta = stream_talker_alias(
self._config,
origin_system=origin_system,
stream_id=stream_id,
rf_src=src_id,
get_dmra_blocks=self._get_dmra_blocks,
)
if ta:
extra["talker_alias"] = ta
return extra
def _group_ctx_base(
self,
*,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
stream_id: bytes,
pkt_time: float,
source_is_obp: bool,
obp_hops: bytes,
obp_source_server: bytes | None,
obp_ber: bytes,
obp_rssi: bytes,
obp_source_rptr: bytes,
synthetic_announcement: bool,
forwarded: tuple[str, ...] | list[str],
extra: dict[str, Any] | None = None,
) -> dict[str, Any]:
systems_cfg = self._config.get("SYSTEMS", {})
sys_cfg = systems_cfg.get(system_name, {})
mode = sys_cfg.get("MODE", "MASTER")
server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0
return dict(
call_family="GROUP",
direction="RX",
origin_system=system_name,
system_mode=str(mode),
peer_id=int_id(peer_id),
src_id=int_id(rf_src),
dst_id=int_id(dst_id),
slot=slot,
stream_id=int_id(stream_id),
server_id=server_id,
is_synthetic=bool(synthetic_announcement),
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,
obp_hops=int_byte(obp_hops) if source_is_obp else None,
obp_source_rptr_id=int_byte(obp_source_rptr) if source_is_obp else None,
ber=int_byte(obp_ber),
rssi=int_byte(obp_rssi),
forwarded_systems=tuple(forwarded),
extra=extra if extra is not None else alias_extra(self._config, peer_id, rf_src, dst_id),
)
def emit_group_voice_start(
self,
*,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
stream_id: bytes,
pkt_time: float,
source_is_obp: bool,
obp_hops: bytes = b"",
obp_source_server: bytes | None = None,
obp_ber: bytes = b"\x00",
obp_rssi: bytes = b"\x00",
obp_source_rptr: bytes = b"\x00\x00\x00\x00",
synthetic_announcement: bool = False,
forwarded: tuple[str, ...] | list[str] = (),
voice_phase: str | None = None,
) -> bool:
if not self._bus.has_subscribers():
return False
ctx_base = self._group_ctx_base(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
source_is_obp=source_is_obp,
obp_hops=obp_hops,
obp_source_server=obp_source_server,
obp_ber=obp_ber,
obp_rssi=obp_rssi,
obp_source_rptr=obp_source_rptr,
synthetic_announcement=synthetic_announcement,
forwarded=forwarded,
)
key = (ctx_base["origin_system"], ctx_base["stream_id"])
if key in self._plugin_started:
return False
self._plugin_started.add(key)
if voice_phase:
ctx_base = {**ctx_base, "extra": {**ctx_base["extra"], "voice_phase": voice_phase}}
self._bus.emit_deferred(VoiceCallStart(context=CallLegContext(**ctx_base)))
return True
def emit_group_voice_end(
self,
*,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
stream_id: bytes,
pkt_time: float,
duration_s: float,
source_is_obp: bool,
obp_hops: bytes = b"",
obp_source_server: bytes | None = None,
obp_ber: bytes = b"\x00",
obp_rssi: bytes = b"\x00",
obp_source_rptr: bytes = b"\x00\x00\x00\x00",
synthetic_announcement: bool = False,
forwarded: tuple[str, ...] | list[str] = (),
) -> bool:
if not self._bus.has_subscribers():
return False
ctx_base = self._group_ctx_base(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
source_is_obp=source_is_obp,
obp_hops=obp_hops,
obp_source_server=obp_source_server,
obp_ber=obp_ber,
obp_rssi=obp_rssi,
obp_source_rptr=obp_source_rptr,
synthetic_announcement=synthetic_announcement,
forwarded=forwarded,
)
key = (ctx_base["origin_system"], ctx_base["stream_id"])
if key in self._plugin_ended:
return False
self._plugin_ended.add(key)
end_extra = self._end_extra(
origin_system=ctx_base["origin_system"],
stream_id=ctx_base["stream_id"],
src_id=ctx_base["src_id"],
base_extra=ctx_base["extra"],
)
self._bus.emit_deferred(
VoiceCallEnd(
context=CallLegContext(**{**ctx_base, "extra": end_extra}),
duration_s=duration_s,
)
)
return True
def notify_group_after_forward(
self,
*,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
frame_type: int,
dtype_vseq: int,
stream_id: bytes,
data: bytes,
pkt_time: float,
source_is_obp: bool,
obp_hops: bytes,
obp_source_server: bytes | None,
obp_ber: bytes,
obp_rssi: bytes,
obp_source_rptr: bytes,
synthetic_announcement: bool,
forwarded: list[str],
) -> None:
if not self._bus.has_subscribers():
return
if frame_type not in (HBPF_VOICE, HBPF_VOICE_SYNC):
return
ctx_base = self._group_ctx_base(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
source_is_obp=source_is_obp,
obp_hops=obp_hops,
obp_source_server=obp_source_server,
obp_ber=obp_ber,
obp_rssi=obp_rssi,
obp_source_rptr=obp_source_rptr,
synthetic_announcement=synthetic_announcement,
forwarded=forwarded,
)
dmrpkt = data[20:53] if len(data) >= 53 else b""
event = VoiceCallFrame(
context=CallLegContext(**ctx_base),
dmrpkt=dmrpkt,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
)
self._bus.emit_deferred(event)
def notify_private_after_forward(
self,
*,
system_name: str,
peer_id: bytes,
rf_src: bytes,
dst_id: bytes,
slot: int,
frame_type: int,
dtype_vseq: int,
stream_id: bytes,
data: bytes,
pkt_time: float,
synthetic_announcement: bool,
forwarded_targets: list[str],
phase: str,
duration_s: float = 0.0,
) -> None:
if not self._bus.has_subscribers():
return
systems_cfg = self._config.get("SYSTEMS", {})
mode = systems_cfg.get(system_name, {}).get("MODE", "MASTER")
server_id = int_byte(self._config.get("GLOBAL", {}).get("SERVER_ID")) or 0
ctx = CallLegContext(
call_family="PRIVATE",
direction="RX",
origin_system=system_name,
system_mode=str(mode),
peer_id=int_id(peer_id),
src_id=int_id(rf_src),
dst_id=int_id(dst_id),
slot=slot,
stream_id=int_id(stream_id),
server_id=server_id,
is_synthetic=bool(synthetic_announcement),
is_proxy_ingress=is_proxy_inject_only(self._config, system_name),
pkt_time=pkt_time,
forwarded_systems=tuple(forwarded_targets),
extra={**alias_extra(self._config, peer_id, rf_src, dst_id), "private_phase": phase},
)
if phase == "START":
self._bus.emit(VoiceCallStart(context=ctx))
elif phase == "FRAME":
dmrpkt = data[20:53] if len(data) >= 53 else b""
self._bus.emit_deferred(
VoiceCallFrame(
context=ctx,
dmrpkt=dmrpkt,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
)
)
elif phase == "END":
end_extra = self._end_extra(
origin_system=system_name,
stream_id=int_id(stream_id),
src_id=int_id(rf_src),
base_extra=ctx.extra,
)
self._bus.emit(
VoiceCallEnd(
context=replace(ctx, extra=end_extra),
duration_s=duration_s,
)
)

@ -0,0 +1,39 @@
# ADN DMR Peer Server - plugin domain package
#
# 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
###############################################################################
from .events import (
CallLegContext,
UnitDataEnd,
UnitDataFrame,
UnitDataStart,
VoiceCallEnd,
VoiceCallFrame,
VoiceCallStart,
)
__all__ = [
"CallLegContext",
"UnitDataEnd",
"UnitDataFrame",
"UnitDataStart",
"VoiceCallEnd",
"VoiceCallFrame",
"VoiceCallStart",
]

@ -0,0 +1,132 @@
# ADN DMR Peer Server - plugin voice events
#
# 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
###############################################################################
"""Voice plugin events — pure domain types (no Twisted, no I/O)."""
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any
@dataclass(frozen=True)
class CallLegContext:
"""Metadata for one voice leg (ingress or bridge observation)."""
call_family: str # "GROUP" | "PRIVATE"
direction: str # "RX" | "TX"
origin_system: str
system_mode: str # MASTER | PEER | OPENBRIDGE
peer_id: int
src_id: int
dst_id: int
slot: int
stream_id: int
server_id: int
is_synthetic: bool
is_proxy_ingress: bool
pkt_time: float
obp_source_server_id: int | None = None
obp_hops: int | None = None
obp_source_rptr_id: int | None = None
ber: int | None = None
rssi: int | None = None
forwarded_systems: tuple[str, ...] = ()
extra: dict[str, Any] = field(default_factory=dict)
def to_metadata_dict(self) -> dict[str, Any]:
"""JSON-serializable metadata for plugin consumers."""
out: dict[str, Any] = {
"call_family": self.call_family,
"direction": self.direction,
"origin_system": self.origin_system,
"system_mode": self.system_mode,
"peer_id": self.peer_id,
"src_id": self.src_id,
"dst_id": self.dst_id,
"slot": self.slot,
"stream_id": self.stream_id,
"server_id": self.server_id,
"is_synthetic": self.is_synthetic,
"is_proxy_ingress": self.is_proxy_ingress,
"pkt_time": self.pkt_time,
"forwarded_systems": list(self.forwarded_systems),
}
if self.obp_source_server_id is not None:
out["obp_source_server_id"] = self.obp_source_server_id
if self.obp_hops is not None:
out["obp_hops"] = self.obp_hops
if self.obp_source_rptr_id is not None:
out["obp_source_rptr_id"] = self.obp_source_rptr_id
if self.ber is not None:
out["ber"] = self.ber
if self.rssi is not None:
out["rssi"] = self.rssi
if self.extra:
out.update(self.extra)
return out
@dataclass(frozen=True)
class VoiceCallStart:
context: CallLegContext
@dataclass(frozen=True)
class VoiceCallFrame:
context: CallLegContext
dmrpkt: bytes
frame_type: int
dtype_vseq: int
@dataclass(frozen=True)
class VoiceCallEnd:
context: CallLegContext
duration_s: float
frame_count: int = 0
@dataclass(frozen=True)
class UnitDataStart:
context: CallLegContext
dtype_vseq: int
data_label: str
seq: int
frame_type: int
@dataclass(frozen=True)
class UnitDataFrame:
context: CallLegContext
dmrpkt: bytes
raw_data: bytes
dtype_vseq: int
data_label: str
seq: int
frame_type: int
bits: int
@dataclass(frozen=True)
class UnitDataEnd:
context: CallLegContext
duration_s: float
packet_count: int = 0

@ -0,0 +1,44 @@
# ADN DMR Peer Server - plugin contract
#
# 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
###############################################################################
"""Plugin contract (domain port)."""
from __future__ import annotations
from typing import Any, Protocol, runtime_checkable
@runtime_checkable
class ServerPlugin(Protocol):
"""Drop-in plugin loaded from plugins/<name>/plugin/."""
name: str
def on_load(self, bus: Any, config: dict[str, Any], server_ctx: Any) -> None:
"""Subscribe to bus; read plugin config."""
def on_event(self, event: object) -> None:
"""Reactor thread — O(1), no I/O."""
def on_reload(self, config: dict[str, Any]) -> None:
"""Optional hot-reload of plugin config."""
def on_shutdown(self) -> None:
"""Flush buffers; stop workers."""

@ -0,0 +1,35 @@
# ADN DMR Peer Server - plugin infrastructure package
#
# 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
###############################################################################
from .loader import (
import_plugin_module,
load_plugin_from_entry,
plugin_config_for_load,
scan_plugin_dirs,
topo_sort,
)
__all__ = [
"import_plugin_module",
"load_plugin_from_entry",
"plugin_config_for_load",
"scan_plugin_dirs",
"topo_sort",
]

@ -0,0 +1,143 @@
# ADN DMR Peer Server - plugin loader
#
# 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
###############################################################################
"""Discover and load drop-in plugins from plugins/<name>/."""
from __future__ import annotations
import importlib.util
import logging
import sys
from dataclasses import dataclass
from pathlib import Path
from typing import Any
import yaml
from ..domain.protocol import ServerPlugin
logger = logging.getLogger(__name__)
_RESERVED_KEYS = frozenset({"enabled", "depends_on", "hot_reload_seconds"})
@dataclass
class PluginEntry:
name: str
config_path: Path
config: dict[str, Any]
def scan_plugin_dirs(directory: Path) -> list[PluginEntry]:
if not directory.is_dir():
return []
entries: list[PluginEntry] = []
for child in sorted(directory.iterdir()):
if not child.is_dir():
continue
cfg_path = child / "config.yaml"
if not cfg_path.is_file():
continue
try:
raw = yaml.safe_load(cfg_path.read_text(encoding="utf-8")) or {}
except (OSError, yaml.YAMLError) as exc:
logger.warning("(PLUGIN-LOADER) skip %s: %s", child.name, exc)
continue
if not isinstance(raw, dict):
logger.warning("(PLUGIN-LOADER) skip %s: config not a mapping", child.name)
continue
if raw.get("enabled") is not True:
continue
entries.append(PluginEntry(name=child.name, config_path=cfg_path, config=raw))
return entries
def topo_sort(entries: list[PluginEntry]) -> list[PluginEntry]:
by_name = {e.name: e for e in entries}
ordered: list[PluginEntry] = []
seen: set[str] = set()
visiting: set[str] = set()
def visit(name: str) -> None:
if name in seen:
return
if name in visiting:
logger.error("(PLUGIN-LOADER) depends_on cycle at %s", name)
return
entry = by_name.get(name)
if entry is None:
return
visiting.add(name)
deps = entry.config.get("depends_on") or []
if isinstance(deps, list):
for dep in deps:
if isinstance(dep, str):
visit(dep)
visiting.discard(name)
seen.add(name)
ordered.append(entry)
for entry in entries:
visit(entry.name)
return ordered
def plugin_config_for_load(raw: dict[str, Any]) -> dict[str, Any]:
"""Strip loader-reserved keys before passing to plugin."""
return {k: v for k, v in raw.items() if k not in _RESERVED_KEYS}
def import_plugin_module(plugin_dir: Path) -> Any:
"""Load plugins/<name>/plugin/ as a Python package (supports relative imports)."""
init_py = plugin_dir / "__init__.py"
if not init_py.is_file():
raise FileNotFoundError(f"missing {init_py}")
pkg_name = f"adn_plugin_{plugin_dir.parent.name.replace('-', '_')}"
cached = sys.modules.get(pkg_name)
if cached is not None:
return cached
spec = importlib.util.spec_from_file_location(
pkg_name,
init_py,
submodule_search_locations=[str(plugin_dir)],
)
if spec is None or spec.loader is None:
raise ImportError(f"cannot load spec for {init_py}")
mod = importlib.util.module_from_spec(spec)
sys.modules[pkg_name] = mod
spec.loader.exec_module(mod)
return mod
def load_plugin_from_entry(
entry: PluginEntry,
plugins_root: Path,
*,
overrides: dict[str, Any] | None = None,
) -> ServerPlugin:
plugin_dir = plugins_root / entry.name / "plugin"
mod = import_plugin_module(plugin_dir)
create = getattr(mod, "create_plugin", None)
if not callable(create):
raise AttributeError(f"{entry.name}: plugin/__init__.py must export create_plugin()")
plugin = create()
if not hasattr(plugin, "on_load") or not hasattr(plugin, "on_event"):
raise TypeError(f"{entry.name}: create_plugin() must return ServerPlugin")
return plugin

@ -52,13 +52,8 @@ _CSV_FAMILIES = {
"GROUP VOICE": "GROUP", "GROUP VOICE": "GROUP",
"PRIVATE VOICE": "PRIVATE", "PRIVATE VOICE": "PRIVATE",
"UNIT DATA": "UNIT", "UNIT DATA": "UNIT",
# routing_use_cases._dtype_labels emits one of these instead of the generic # routing_use_cases emits these labels for dtype_vseq in (6, 7, 8). CSBK (3) is
# "UNIT DATA" label for dtype_vseq in (3, 6, 7, 8) — private-call CSBK setup # still logged at INFO but no longer emitted to monitor (prelude storm filter).
# signaling and SMS/GPS (ARS/LRRP) header + VCSBK blocks. Without these,
# parse_bridge_event_csv() returns None and the event is silently dropped
# ("voice_event not emitted (unmapped CSV)"), so a private call whose CSBK
# handshake never escalates to voice (e.g. destination unreachable) never
# shows up in the monitor at all.
"UNIT CSBK": "UNIT", "UNIT CSBK": "UNIT",
"UNIT DATA HEADER": "UNIT", "UNIT DATA HEADER": "UNIT",
"UNIT VCSBK 1/2 DATA BLOCK": "UNIT", "UNIT VCSBK 1/2 DATA BLOCK": "UNIT",

@ -53,6 +53,7 @@ from .helpers import (
hbp_ingress_new_stream_collision, hbp_ingress_new_stream_collision,
master_per_peer_slot_contention, master_per_peer_slot_contention,
tg_has_active_conversation, tg_has_active_conversation,
unit_data_reportable,
) )
from .peer_downlink_index import count_connected_peers from .peer_downlink_index import count_connected_peers
@ -384,8 +385,10 @@ class HbpForwardMixin:
"(%s) UNIT trace -> %s stream=%s payload33=%s", "(%s) UNIT trace -> %s stream=%s payload33=%s",
source_system, d_system, int_id(stream_id), dmrpkt.hex(), source_system, d_system, int_id(stream_id), dmrpkt.hex(),
) )
self._send_routing_event( _dtype_vseq = data[15] & 0xF if len(data) > 15 else 0
"UNIT DATA,DATA,TX,{},{},{},{},{},{}".format( if unit_data_reportable(_dtype_vseq):
d_system, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id), self._send_routing_event(
"UNIT DATA,DATA,TX,{},{},{},{},{},{}".format(
d_system, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id),
)
) )
)

@ -815,6 +815,90 @@ def _obp_stream_active(st: dict[str, Any], pkt_time: float) -> bool:
return True return True
def obp_ingress_stream_on_system(
obp_status: dict[Any, Any],
stream_id: bytes,
dst_id: bytes,
rf_src: bytes,
) -> int:
best_sid: bytes | None = None
best_first: float | None = None
for sid, st in obp_status.items():
if not isinstance(sid, (bytes, bytearray)) or not isinstance(st, dict):
continue
if "H_LC" in st:
continue
if st.get("TGID") != dst_id:
continue
if not _same_rf_source(st.get("RFS", b"\x00\x00\x00"), rf_src):
continue
if st.get("_fin"):
continue
first = st.get("1ST")
if first is None:
continue
if best_first is None or float(first) < best_first:
best_first = float(first)
best_sid = sid
return int_id(best_sid if best_sid is not None else stream_id)
def obp_cross_system_winner(
protocols: dict[str, Any],
systems_cfg: dict[str, Any],
system_name: str,
stream_id: bytes,
dst_id: bytes,
) -> str:
hr_times: dict[str, float] = {}
for other_name, proto in protocols.items():
if systems_cfg.get(other_name, {}).get("MODE") != "OPENBRIDGE":
continue
obp_status = getattr(proto, "STATUS", None)
if not isinstance(obp_status, dict):
continue
ent = obp_status.get(stream_id)
if not isinstance(ent, dict) or "1ST" not in ent or ent.get("TGID") != dst_id:
continue
hr_times[other_name] = float(ent["1ST"])
if not hr_times:
return system_name
return min(hr_times, key=hr_times.get)
def obp_is_canonical_ingress(
protocols: dict[str, Any],
systems_cfg: dict[str, Any],
system_name: str,
stream_id: bytes,
dst_id: bytes,
rf_src: bytes,
) -> bool:
src_proto = protocols.get(system_name)
obp_status = getattr(src_proto, "STATUS", None) if src_proto else None
if not isinstance(obp_status, dict):
return False
sid = int_id(stream_id)
if sid != obp_ingress_stream_on_system(obp_status, stream_id, dst_id, rf_src):
return False
return system_name == obp_cross_system_winner(
protocols, systems_cfg, system_name, stream_id, dst_id
)
def obp_status_plugin_voice(
protocols: dict[str, Any],
system_name: str,
stream_id: bytes,
) -> bool:
proto = protocols.get(system_name)
status = getattr(proto, "STATUS", None) if proto else None
if not isinstance(status, dict):
return False
ent = status.get(stream_id)
return isinstance(ent, dict) and bool(ent.get("_plugin_voice"))
def tg_has_active_conversation( def tg_has_active_conversation(
protocols: dict[str, Any], protocols: dict[str, Any],
systems_cfg: dict[str, Any], systems_cfg: dict[str, Any],
@ -1006,6 +1090,15 @@ def unit_data_hbp_target_idle(
) )
def unit_data_reportable(dtype_vseq: int) -> bool:
"""True when a unit-data frame should emit monitor/report events.
CSBK prelude (dtype 3) is routed but omitted from BRDG_EVENT — same policy as
data-log ``log_dtypes`` — to avoid SMS/GPS setup storms in logs and monitors.
"""
return dtype_vseq in (6, 7, 8)
def is_unit_data_ingress( def is_unit_data_ingress(
call_type: str, call_type: str,
dtype_vseq: int, dtype_vseq: int,

@ -52,7 +52,7 @@ from typing import Any
from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, int_id from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, int_id
from ...domain.dmr import decode from ...domain.dmr import decode
from ...domain.dmr.const import LC_OPT from ...domain.dmr.const import LC_OPT
from .helpers import group_voice_tg_ingress_collision from .helpers import group_voice_tg_ingress_collision, obp_is_canonical_ingress, unit_data_reportable
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -152,6 +152,7 @@ class ObpForwardMixin:
stream_id: bytes, stream_id: bytes,
data: bytes, data: bytes,
obp_hops: bytes, obp_hops: bytes,
synthetic_announcement: bool = False,
) -> bool: ) -> bool:
"""Port of bridge_master.routerOBP.dmrd_received group/vcsbk (~2269-2411). False = drop packet.""" """Port of bridge_master.routerOBP.dmrd_received group/vcsbk (~2269-2411). False = drop packet."""
pkt_time = time.time() pkt_time = time.time()
@ -221,6 +222,23 @@ class ObpForwardMixin:
system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id) system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id)
) )
) )
bridge = self._voice_plugin_bridge
if bridge is not None and bridge.has_subscribers() and obp_is_canonical_ingress(
protocols, systems_cfg, system_name, stream_id, dst_id, rf_src,
) and bridge.emit_group_voice_start(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
source_is_obp=True,
obp_hops=obp_hops if obp_hops else b"",
synthetic_announcement=synthetic_announcement,
voice_phase="INGRESS",
):
status[stream_id]["_plugin_voice"] = True
else: else:
st = status[stream_id] st = status[stream_id]
if "packets" in st: if "packets" in st:
@ -473,8 +491,10 @@ class ObpForwardMixin:
logger.warning("(%s) send_data_to_obp %s failed: %s", source_system, target, exc) logger.warning("(%s) send_data_to_obp %s failed: %s", source_system, target, exc)
return return
logger.debug("(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s", source_system, target, int_id(dst_id)) logger.debug("(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s", source_system, target, int_id(dst_id))
self._send_routing_event( _dtype_vseq = data[15] & 0xF if len(data) > 15 else 0
"UNIT DATA,DATA,TX,{},{},{},{},{},{}".format( if unit_data_reportable(_dtype_vseq):
target, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id), self._send_routing_event(
"UNIT DATA,DATA,TX,{},{},{},{},{},{}".format(
target, int_id(stream_id), int_id(peer_id), int_id(rf_src), 1, int_id(dst_id),
)
) )
)

@ -296,6 +296,20 @@ class RoutingTimerMixin:
system_name, int_id(stream_id), int_id(peer), int_id(rfs), 1, int_id(tgid), duration system_name, int_id(stream_id), int_id(peer), int_id(rfs), 1, int_id(tgid), duration
) )
) )
bridge = self._voice_plugin_bridge
if st.get("_plugin_voice") and bridge is not None:
bridge.emit_group_voice_end(
system_name=system_name,
peer_id=peer,
rf_src=rfs,
dst_id=tgid,
slot=1,
stream_id=stream_id,
pkt_time=last,
duration_s=duration,
source_is_obp=True,
)
st.pop("_plugin_voice", None)
st["_to"] = True st["_to"] = True
self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now) self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now)
continue continue
@ -331,6 +345,19 @@ class RoutingTimerMixin:
_slot.get("RX_TIME", 0) - _slot.get("RX_START", 0), _slot.get("RX_TIME", 0) - _slot.get("RX_START", 0),
) )
) )
bridge = self._voice_plugin_bridge
if bridge is not None:
bridge.emit_group_voice_end(
system_name=system_name,
peer_id=_slot.get("RX_PEER", b"\x00\x00\x00\x00"),
rf_src=_slot.get("RX_RFS", b"\x00\x00\x00"),
dst_id=_slot.get("RX_TGID", b"\x00\x00\x00"),
slot=slot,
stream_id=_slot.get("RX_STREAM_ID", b"\x00"),
pkt_time=_slot.get("RX_TIME", now),
duration_s=_slot.get("RX_TIME", 0) - _slot.get("RX_START", 0),
source_is_obp=False,
)
if _slot.get("RX_TIME", 0) < now - 60: if _slot.get("RX_TIME", 0) < now - 60:
_slot["RX_STREAM_ID"] = b"\x00" _slot["RX_STREAM_ID"] = b"\x00"
tx_streams = _slot.get("TX_STREAMS") tx_streams = _slot.get("TX_STREAMS")

@ -36,7 +36,17 @@ from typing import Any
from bitarray import bitarray from bitarray import bitarray
from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4, int_id from ..domain import (
HBPF_DATA_SYNC,
HBPF_SLT_VHEAD,
HBPF_SLT_VTERM,
HBPF_VOICE,
HBPF_VOICE_SYNC,
STREAM_TO,
bytes_3,
bytes_4,
int_id,
)
from ..domain.dmr import bptc from ..domain.dmr import bptc
from .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder from .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder
from .reporting_use_cases import ReportingUseCases from .reporting_use_cases import ReportingUseCases
@ -51,11 +61,13 @@ from .routing.helpers import (
obp_deferred_bridge_tx_leg, obp_deferred_bridge_tx_leg,
obp_flat_bridge_tx_idle, obp_flat_bridge_tx_idle,
obp_publish_flat_bridge_tx, obp_publish_flat_bridge_tx,
obp_status_plugin_voice,
obp_sync_flat_bridge_tx_times, obp_sync_flat_bridge_tx_times,
obp_target_bcsq_quenches_stream, obp_target_bcsq_quenches_stream,
resolve_voice_peer_id, resolve_voice_peer_id,
slot_has_active_voice, slot_has_active_voice,
unit_data_hbp_target_idle, unit_data_hbp_target_idle,
unit_data_reportable,
) )
from .routing.lc_ta import LcTaMixin from .routing.lc_ta import LcTaMixin
from .routing.obp_forward import ObpForwardMixin from .routing.obp_forward import ObpForwardMixin
@ -96,6 +108,8 @@ class RoutingUseCases(
call_later: Any = None, call_later: Any = None,
encode_emblc: DmrEmbeddedLcEncoder | None = None, encode_emblc: DmrEmbeddedLcEncoder | None = None,
ta_emblc_encoder: TalkerAliasEmblcEncoder | None = None, ta_emblc_encoder: TalkerAliasEmblcEncoder | None = None,
voice_plugin_bridge: Any = None,
data_plugin_bridge: Any = None,
) -> None: ) -> None:
self._acl_router = acl_router self._acl_router = acl_router
self._config = config self._config = config
@ -114,6 +128,8 @@ class RoutingUseCases(
raise TypeError("encode_emblc and ta_emblc_encoder are required (wire from main.py)") raise TypeError("encode_emblc and ta_emblc_encoder are required (wire from main.py)")
self._encode_emblc = encode_emblc self._encode_emblc = encode_emblc
self._talker_alias = TalkerAliasUseCases(config, ta_emblc_encoder=ta_emblc_encoder) self._talker_alias = TalkerAliasUseCases(config, ta_emblc_encoder=ta_emblc_encoder)
self._voice_plugin_bridge = voice_plugin_bridge
self._data_plugin_bridge = data_plugin_bridge
# (source_system, stream_id) -> {rf_src, peer, targets, timer} # (source_system, stream_id) -> {rf_src, peer, targets, timer}
self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {} self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {}
# Passthrough DMRA/embed relay already applied for this source stream. # Passthrough DMRA/embed relay already applied for this source stream.
@ -281,6 +297,7 @@ class RoutingUseCases(
stream_id, stream_id,
data, data,
obp_hops if obp_use_parsed else b"", obp_hops if obp_use_parsed else b"",
synthetic_announcement=synthetic_announcement,
): ):
return return
# SINGLE_MODE in-band VTERM on another TG (e.g. 9990 echo) deactivates static OFF # SINGLE_MODE in-band VTERM on another TG (e.g. 9990 echo) deactivates static OFF
@ -321,6 +338,18 @@ class RoutingUseCases(
int(synthetic_announcement), int(synthetic_announcement),
) )
) )
if self._voice_plugin_bridge is not None:
self._voice_plugin_bridge.emit_group_voice_start(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
source_is_obp=False,
synthetic_announcement=synthetic_announcement,
)
elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
if not _obp_grp: if not _obp_grp:
duration = 0.0 duration = 0.0
@ -361,6 +390,19 @@ class RoutingUseCases(
int(synthetic_announcement), int(synthetic_announcement),
) )
) )
if self._voice_plugin_bridge is not None:
self._voice_plugin_bridge.emit_group_voice_end(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=pkt_time,
duration_s=duration,
source_is_obp=False,
synthetic_announcement=synthetic_announcement,
)
has_source = bool( has_source = bool(
self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int) self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int)
) )
@ -951,6 +993,26 @@ class RoutingUseCases(
int(synthetic_announcement), int(synthetic_announcement),
) )
) )
if ost.get("_plugin_voice") and self._voice_plugin_bridge is not None:
self._voice_plugin_bridge.emit_group_voice_end(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
stream_id=stream_id,
pkt_time=_end_t,
duration_s=call_duration,
source_is_obp=True,
obp_hops=obp_hops if obp_use_parsed else b"",
obp_source_server=obp_source_server,
obp_ber=obp_ber,
obp_rssi=obp_rssi,
obp_source_rptr=obp_source_rptr,
synthetic_announcement=synthetic_announcement,
forwarded=tuple(forwarded),
)
ost.pop("_plugin_voice", None)
ost["_fin"] = True ost["_fin"] = True
self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, _end_t) self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, _end_t)
ost["lastSeq"] = False ost["lastSeq"] = False
@ -960,6 +1022,35 @@ class RoutingUseCases(
"(ROUTER) Bridged TG %s from %s -> %s", "(ROUTER) Bridged TG %s from %s -> %s",
relay_table_key, system_name, ", ".join(forwarded), relay_table_key, system_name, ", ".join(forwarded),
) )
if (
call_type in ("group", "vcsbk")
and self._voice_plugin_bridge is not None
and frame_type in (HBPF_VOICE, HBPF_VOICE_SYNC)
and (
not source_is_obp
or obp_status_plugin_voice(protocols, system_name, stream_id)
)
):
self._voice_plugin_bridge.notify_group_after_forward(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
source_is_obp=source_is_obp,
obp_hops=obp_hops if obp_use_parsed else b"",
obp_source_server=obp_source_server,
obp_ber=obp_ber,
obp_rssi=obp_rssi,
obp_source_rptr=obp_source_rptr,
synthetic_announcement=synthetic_announcement,
forwarded=forwarded,
)
return True return True
# ── Unit DATA path (SMS, GPS, CSBK) — legacy routerOBP/routerHBP unit data branch ── # ── Unit DATA path (SMS, GPS, CSBK) — legacy routerOBP/routerHBP unit data branch ──
@ -1015,6 +1106,7 @@ class RoutingUseCases(
dmrpkt = data[20:53] if len(data) >= 53 else b"" dmrpkt = data[20:53] if len(data) >= 53 else b""
_bits = data[15] if len(data) > 15 else 0 _bits = data[15] if len(data) > 15 else 0
_int_dst_id = int_id(dst_id) _int_dst_id = int_id(dst_id)
_forwarded: list[str] = []
systems_cfg = self._config.get("SYSTEMS", {}) systems_cfg = self._config.get("SYSTEMS", {})
source_is_obp = systems_cfg.get(system_name, {}).get("MODE") == "OPENBRIDGE" source_is_obp = systems_cfg.get(system_name, {}).get("MODE") == "OPENBRIDGE"
global_cfg = self._config.get("GLOBAL", {}) global_cfg = self._config.get("GLOBAL", {})
@ -1123,19 +1215,21 @@ class RoutingUseCases(
system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot, system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), _int_dst_id, slot,
) )
_dtype_labels = {3: "UNIT CSBK", 6: "UNIT DATA HEADER", 7: "UNIT VCSBK 1/2 DATA BLOCK", 8: "UNIT VCSBK 3/4 DATA BLOCK"} if unit_data_reportable(dtype_vseq):
_label = _dtype_labels.get(dtype_vseq, "UNIT DATA") _dtype_labels = {6: "UNIT DATA HEADER", 7: "UNIT VCSBK 1/2 DATA BLOCK", 8: "UNIT VCSBK 3/4 DATA BLOCK"}
self._send_routing_event( _label = _dtype_labels.get(dtype_vseq, "UNIT DATA")
"{},DATA,RX,{},{},{},{},{},{}".format( self._send_routing_event(
_label, system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, _int_dst_id, "{},DATA,RX,{},{},{},{},{},{}".format(
_label, system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, _int_dst_id,
)
) )
)
# DATA-GATEWAY forwarding (legacy ~2281-2284 / ~3083-3087) # DATA-GATEWAY forwarding (legacy ~2281-2284 / ~3083-3087)
if global_cfg.get("DATA_GATEWAY"): if global_cfg.get("DATA_GATEWAY"):
dg_cfg = systems_cfg.get("DATA-GATEWAY", {}) dg_cfg = systems_cfg.get("DATA-GATEWAY", {})
if dg_cfg.get("MODE") == "OPENBRIDGE" and dg_cfg.get("ENABLED"): if dg_cfg.get("MODE") == "OPENBRIDGE" and dg_cfg.get("ENABLED"):
logger.debug("(%s) DATA packet sent to DATA-GATEWAY", system_name) logger.debug("(%s) DATA packet sent to DATA-GATEWAY", system_name)
_forwarded.append("DATA-GATEWAY")
self._send_data_to_obp( self._send_data_to_obp(
system_name, "DATA-GATEWAY", data, dmrpkt, pkt_time, stream_id, system_name, "DATA-GATEWAY", data, dmrpkt, pkt_time, stream_id,
dst_id, peer_id, rf_src, _bits, slot, dst_id, peer_id, rf_src, _bits, slot,
@ -1151,6 +1245,7 @@ class RoutingUseCases(
if sys_name == "DATA-GATEWAY": if sys_name == "DATA-GATEWAY":
continue continue
if sys_cfg.get("MODE") == "OPENBRIDGE" and sys_cfg.get("VER", 1) > 1 and _int_dst_id >= 1000000: if sys_cfg.get("MODE") == "OPENBRIDGE" and sys_cfg.get("VER", 1) > 1 and _int_dst_id >= 1000000:
_forwarded.append(sys_name)
self._send_data_to_obp( self._send_data_to_obp(
system_name, sys_name, data, dmrpkt, pkt_time, stream_id, system_name, sys_name, data, dmrpkt, pkt_time, stream_id,
dst_id, peer_id, rf_src, _bits, slot, dst_id, peer_id, rf_src, _bits, slot,
@ -1172,6 +1267,7 @@ class RoutingUseCases(
_hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5)
if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime):
_tmp_bits = _bits ^ (1 << 7) if slot != _d_slot else _bits _tmp_bits = _bits ^ (1 << 7) if slot != _d_slot else _bits
_forwarded.append(_d_system)
self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_d_peer_id) self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_d_peer_id)
elif not is_private_subscriber_dst(dst_id): elif not is_private_subscriber_dst(dst_id):
self._log_unit_data_hbp_busy( self._log_unit_data_hbp_busy(
@ -1199,6 +1295,7 @@ class RoutingUseCases(
_hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5)
if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime):
_tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits
_forwarded.append(_d_system)
self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer) self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer)
elif not is_private_subscriber_dst(dst_id): elif not is_private_subscriber_dst(dst_id):
self._log_unit_data_hbp_busy( self._log_unit_data_hbp_busy(
@ -1214,6 +1311,7 @@ class RoutingUseCases(
_hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5)
if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime):
_tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits _tmp_bits = _bits ^ (1 << 7) if slot != 2 else _bits
_forwarded.append(_d_system)
self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer) self._send_data_to_hbp(system_name, _d_system, _d_slot, dst_id, _tmp_bits, data, dmrpkt, rf_src, stream_id, peer_id, d_peer_id=_to_peer)
elif not is_private_subscriber_dst(dst_id): elif not is_private_subscriber_dst(dst_id):
self._log_unit_data_hbp_busy( self._log_unit_data_hbp_busy(
@ -1224,6 +1322,29 @@ class RoutingUseCases(
if _matched: if _matched:
break break
if self._data_plugin_bridge is not None:
self._data_plugin_bridge.notify_after_forward(
route="unit",
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
seq=seq,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
source_is_obp=source_is_obp,
obp_hops=_hops,
obp_source_server=_source_server if source_is_obp else None,
obp_ber=_ber,
obp_rssi=_rssi,
obp_source_rptr=_source_rptr,
forwarded=_forwarded,
)
def _pvt_call_received( def _pvt_call_received(
self, self,
system_name: str, system_name: str,
@ -1251,6 +1372,7 @@ class RoutingUseCases(
if source_proto: if source_proto:
source_status.setdefault(slot, {}) source_status.setdefault(slot, {})
slot_st = source_status[slot] if source_proto and slot in source_status else {} slot_st = source_status[slot] if source_proto and slot in source_status else {}
_notify_pvt_start = False
_unit_data = is_unit_data_ingress( _unit_data = is_unit_data_ingress(
"unit", dtype_vseq, stream_id, slot_st.get("RX_STREAM_ID"), "unit", dtype_vseq, stream_id, slot_st.get("RX_STREAM_ID"),
) )
@ -1261,6 +1383,7 @@ class RoutingUseCases(
system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot, system_name, int_id(stream_id), int_id(rf_src), int_id(peer_id), int_id(dst_id), slot,
) )
return return
_notify_pvt_start = True
slot_st["RX_START"] = pkt_time slot_st["RX_START"] = pkt_time
self._pvt_same_system_dst_peer_id = None self._pvt_same_system_dst_peer_id = None
if dst_id in sub_map: if dst_id in sub_map:
@ -1397,6 +1520,28 @@ class RoutingUseCases(
self._send_to_system(_target, send_data) self._send_to_system(_target, send_data)
except Exception as e: except Exception as e:
logger.warning("(ROUTER) send_to_system %s failed: %s", _target, e) logger.warning("(ROUTER) send_to_system %s failed: %s", _target, e)
if self._data_plugin_bridge is not None and _unit_data:
self._data_plugin_bridge.notify_after_forward(
route="private",
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
seq=seq,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
source_is_obp=False,
obp_hops=b"",
obp_source_server=None,
obp_ber=data[53:54] if len(data) > 53 else b"\x00",
obp_rssi=data[54:55] if len(data) > 54 else b"\x00",
obp_source_rptr=b"\x00\x00\x00\x00",
forwarded=list(getattr(self, "_pvt_targets", [])),
)
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM: if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM:
self._pvt_targets = [] self._pvt_targets = []
call_duration = pkt_time - slot_st.get("RX_START", pkt_time) call_duration = pkt_time - slot_st.get("RX_START", pkt_time)
@ -1430,6 +1575,57 @@ class RoutingUseCases(
) )
) )
self._pvt_same_system_dst_peer_id = None self._pvt_same_system_dst_peer_id = None
if self._voice_plugin_bridge is not None and not _unit_data:
_pvt_fwd = list(getattr(self, "_pvt_targets", []))
if _notify_pvt_start:
self._voice_plugin_bridge.notify_private_after_forward(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
synthetic_announcement=False,
forwarded_targets=_pvt_fwd,
phase="START",
)
elif frame_type in (HBPF_VOICE, HBPF_VOICE_SYNC):
self._voice_plugin_bridge.notify_private_after_forward(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
synthetic_announcement=False,
forwarded_targets=_pvt_fwd,
phase="FRAME",
)
elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM and slot_st.get("RX_TYPE") != HBPF_SLT_VTERM:
self._voice_plugin_bridge.notify_private_after_forward(
system_name=system_name,
peer_id=peer_id,
rf_src=rf_src,
dst_id=dst_id,
slot=slot,
frame_type=frame_type,
dtype_vseq=dtype_vseq,
stream_id=stream_id,
data=data,
pkt_time=pkt_time,
synthetic_announcement=False,
forwarded_targets=_pvt_fwd,
phase="END",
duration_s=call_duration,
)
if slot_st: if slot_st:
if _unit_data: if _unit_data:
# Keep stream continuity for multi-frame unit data without marking the # Keep stream continuity for multi-frame unit data without marking the

@ -38,6 +38,11 @@ from adn_server.application import (
VoiceUseCases, VoiceUseCases,
) )
from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases from adn_server.application.dynamic_tg_use_cases import DynamicTgUseCases
from adn_server.application.plugins.application.bus import PluginBus
from adn_server.application.plugins.application.context import ServerContext
from adn_server.application.plugins.application.data_bridge import DataPluginBridge
from adn_server.application.plugins.application.manager import PluginManager
from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge
from adn_server.application.proxy.deployment import ( from adn_server.application.proxy.deployment import (
is_obp_proxy_managed, is_obp_proxy_managed,
is_proxy_inject_only, is_proxy_inject_only,
@ -431,6 +436,19 @@ def run_peer_server(
call_from_reactor=reactor.callFromThread, call_from_reactor=reactor.callFromThread,
) )
plugin_bus = PluginBus()
plugin_bus.set_call_later(reactor.callLater)
server_ctx = ServerContext(
config=config,
project_root=project_root,
defer_to_thread=threads.deferToThread,
call_from_reactor=reactor.callFromThread,
call_later=reactor.callLater,
)
plugin_manager = PluginManager(plugin_bus, server_ctx, project_root)
voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks)
data_plugin_bridge = DataPluginBridge(plugin_bus, config)
routing_use_cases = RoutingUseCases( routing_use_cases = RoutingUseCases(
acl_router, acl_router,
config, config,
@ -445,7 +463,11 @@ def run_peer_server(
call_later=reactor.callLater, call_later=reactor.callLater,
encode_emblc=encode_emblc, encode_emblc=encode_emblc,
ta_emblc_encoder=default_ta_emblc_encoder, ta_emblc_encoder=default_ta_emblc_encoder,
voice_plugin_bridge=voice_plugin_bridge,
data_plugin_bridge=data_plugin_bridge,
) )
plugin_manager.discover_and_load(config)
reactor.addSystemEventTrigger("before", "shutdown", plugin_manager.shutdown_all)
routing_use_cases.apply_startup_subscriptions() routing_use_cases.apply_startup_subscriptions()
dynamic_tg_uc = DynamicTgUseCases( dynamic_tg_uc = DynamicTgUseCases(
dynamic_tg_store, dynamic_tg_store,
@ -756,6 +778,9 @@ def run_peer_server(
obp_proxy_state = start_obp_proxy_service(config, protocols, logger=logger) obp_proxy_state = start_obp_proxy_service(config, protocols, logger=logger)
if result.added or result.removed or result.updated or result.rebound: if result.added or result.removed or result.updated or result.rebound:
_on_config_systems_changed() _on_config_systems_changed()
voice_plugin_bridge.update_config(config)
data_plugin_bridge.update_config(config)
plugin_manager.rescan(config)
def _do_config_reload() -> None: def _do_config_reload() -> None:
nonlocal report_mqtt, proxy_state, obp_proxy_state nonlocal report_mqtt, proxy_state, obp_proxy_state

@ -72,6 +72,17 @@ def reopen_file_handlers(logger: logging.Logger | None = None) -> int:
return count return count
def _silence_noisy_library_loggers() -> None:
"""Keep third-party DEBUG chatter out of the server log when LOG_LEVEL=DEBUG."""
for name in (
"watchdog",
"watchdog.observers",
"watchdog.observers.inotify",
"watchdog.observers.inotify_buffer",
):
logging.getLogger(name).setLevel(logging.WARNING)
def reapply_log_level(log_config: dict[str, Any]) -> str: def reapply_log_level(log_config: dict[str, Any]) -> str:
"""Apply LOGGER.LOG_LEVEL after SIGHUP reload (handlers unchanged). """Apply LOGGER.LOG_LEVEL after SIGHUP reload (handlers unchanged).
@ -90,6 +101,7 @@ def reapply_log_level(log_config: dict[str, Any]) -> str:
handler.setLevel(level) handler.setLevel(level)
for handler in app_logger.handlers: for handler in app_logger.handlers:
handler.setLevel(level) handler.setLevel(level)
_silence_noisy_library_loggers()
return logging.getLevelName(level) return logging.getLevelName(level)
@ -134,4 +146,5 @@ def setup_logging(log_config: dict[str, Any]) -> logging.Logger:
logging.basicConfig(level=level, handlers=handlers or [logging.NullHandler()], force=True) logging.basicConfig(level=level, handlers=handlers or [logging.NullHandler()], force=True)
logger = logging.getLogger(log_name) logger = logging.getLogger(log_name)
logger.setLevel(level) logger.setLevel(level)
_silence_noisy_library_loggers()
return logger return logger

@ -0,0 +1,35 @@
# ADN DMR Peer Server - tests plugin bridge helpers
#
# 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
###############################################################################
from __future__ import annotations
from adn_server.application.plugins.application.bridge_common import alias_extra
def test_alias_extra_resolves_server_name_with_string_server_ids_key():
config = {
"GLOBAL": {"SERVER_ID": 7302},
"_SERVER_IDS": {"7302": "ADN_7302_Chile"},
"_PEER_IDS": {},
"_SUB_IDS": {},
"_TG_IDS": {},
}
extra = alias_extra(config, b"\x00\x01\x1d\x12", b"\x00!\x0a", b"\x00\x03H")
assert extra["server_name"] == "ADN_7302_Chile"

@ -0,0 +1,98 @@
# ADN DMR Peer Server - tests OBP canonical ingress helpers
#
# 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
###############################################################################
"""Core OBP canonical ingress selection for plugin lifecycle."""
from __future__ import annotations
from types import SimpleNamespace
from adn_server.application.routing.helpers import obp_is_canonical_ingress
def _stream_id(n: int) -> bytes:
return n.to_bytes(4, "big")
def _rf(n: int) -> bytes:
return n.to_bytes(3, "big")
def _tg(n: int) -> bytes:
return n.to_bytes(3, "big")
def test_canonical_ingress_single_stream_on_obp() -> None:
sid = _stream_id(42)
status = {
sid: {"1ST": 1.0, "TGID": _tg(0x33420), "RFS": _rf(0x334202)},
}
proto = SimpleNamespace(STATUS=status)
protocols = {"OBP-CL": proto}
systems_cfg = {"OBP-CL": {"MODE": "OPENBRIDGE"}}
assert obp_is_canonical_ingress(
protocols, systems_cfg, "OBP-CL", sid, _tg(0x33420), _rf(0x334202),
)
def test_canonical_ingress_picks_earliest_same_tg_rf() -> None:
winner = _stream_id(42)
loser = _stream_id(99)
status = {
winner: {"1ST": 1.0, "TGID": _tg(0x33420), "RFS": _rf(0x334202)},
loser: {"1ST": 2.0, "TGID": _tg(0x33420), "RFS": _rf(0x334202)},
}
proto = SimpleNamespace(STATUS=status)
protocols = {"OBP-CL": proto}
systems_cfg = {"OBP-CL": {"MODE": "OPENBRIDGE"}}
assert obp_is_canonical_ingress(
protocols, systems_cfg, "OBP-CL", winner, _tg(0x33420), _rf(0x334202),
)
assert not obp_is_canonical_ingress(
protocols, systems_cfg, "OBP-CL", loser, _tg(0x33420), _rf(0x334202),
)
def test_obp_status_plugin_voice_flag() -> None:
from adn_server.application.routing.helpers import obp_status_plugin_voice
sid = _stream_id(42)
proto = SimpleNamespace(STATUS={sid: {"_plugin_voice": True}})
protocols = {"OBP-CL": proto}
assert obp_status_plugin_voice(protocols, "OBP-CL", sid)
assert not obp_status_plugin_voice(protocols, "OBP-CL", _stream_id(99))
def test_canonical_ingress_cross_obp_winner() -> None:
sid = _stream_id(42)
proto_a = SimpleNamespace(
STATUS={sid: {"1ST": 2.0, "TGID": _tg(0x33420), "RFS": _rf(0x334202)}},
)
proto_b = SimpleNamespace(
STATUS={sid: {"1ST": 1.0, "TGID": _tg(0x33420), "RFS": _rf(0x334202)}},
)
protocols = {"OBP-A": proto_a, "OBP-B": proto_b}
systems_cfg = {"OBP-A": {"MODE": "OPENBRIDGE"}, "OBP-B": {"MODE": "OPENBRIDGE"}}
assert not obp_is_canonical_ingress(
protocols, systems_cfg, "OBP-A", sid, _tg(0x33420), _rf(0x334202),
)
assert obp_is_canonical_ingress(
protocols, systems_cfg, "OBP-B", sid, _tg(0x33420), _rf(0x334202),
)

@ -0,0 +1,82 @@
# ADN DMR Peer Server - tests plugin bus
#
# 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
###############################################################################
"""Plugin bus circuit breaker behaviour."""
from __future__ import annotations
import time
from typing import Any
from adn_server.application.plugins.application.bus import PluginBus
class _SlowPlugin:
name = "slow"
def on_event(self, event: object) -> None:
time.sleep(0.002)
def on_shutdown(self) -> None:
pass
class _BrokenPlugin:
name = "broken"
def on_event(self, event: object) -> None:
raise RuntimeError("boom")
def on_shutdown(self) -> None:
pass
def test_plugin_bus_does_not_disable_on_budget_overrun() -> None:
bus = PluginBus(budget_s=0.0001, max_trips=2)
plugin = _SlowPlugin()
bus.register_plugin(plugin)
for _ in range(5):
bus.emit("evt")
assert plugin in bus._plugins # noqa: SLF001
assert "slow" not in bus._disabled # noqa: SLF001
def test_plugin_bus_disables_on_exception() -> None:
bus = PluginBus(max_trips=2)
plugin = _BrokenPlugin()
bus.register_plugin(plugin)
for _ in range(2):
bus.emit("evt")
assert plugin not in bus._plugins # noqa: SLF001
assert "broken" in bus._disabled # noqa: SLF001
def test_plugin_bus_re_register_after_disable() -> None:
bus = PluginBus(max_trips=1)
plugin = _BrokenPlugin()
bus.register_plugin(plugin)
bus.emit("evt")
assert "broken" in bus._disabled # noqa: SLF001
assert plugin not in bus._plugins # noqa: SLF001
plugin2 = _BrokenPlugin()
bus.register_plugin(plugin2)
assert "broken" not in bus._disabled # noqa: SLF001
assert plugin2 in bus._plugins # noqa: SLF001

@ -0,0 +1,65 @@
# ADN DMR Peer Server - tests stream talker alias helper
#
# 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
###############################################################################
"""Plugin bridge helper: resolve Talker Alias from buffered DMRA blocks."""
from __future__ import annotations
from adn_server.application.plugins.application.bridge_common import stream_talker_alias
from adn_server.domain import bytes_4
from adn_server.domain.talker_alias import build_dmra_packets, parse_dmra_packet
def test_stream_talker_alias_returns_none_without_callback() -> None:
assert (
stream_talker_alias(
{},
origin_system="MASTER-A",
stream_id=1,
rf_src=2,
get_dmra_blocks=None,
)
is None
)
def test_stream_talker_alias_decodes_buffered_blocks() -> None:
rf_src = bytes_4(3120001)
stream_id = 0x90909090
blocks: dict[int, bytes] = {}
for pkt in build_dmra_packets(rf_src, "CE5RPY Test", "utf8"):
parsed = parse_dmra_packet(pkt)
assert parsed is not None
_, block_id, payload = parsed
blocks[block_id] = payload
def get_dmra_blocks(system: str, sid: bytes) -> dict[int, bytes] | None:
if system == "MASTER-A" and sid == bytes_4(stream_id):
return blocks
return None
text = stream_talker_alias(
{},
origin_system="MASTER-A",
stream_id=stream_id,
rf_src=3120001,
get_dmra_blocks=get_dmra_blocks,
)
assert text == "CE5RPY Test"

@ -0,0 +1,137 @@
# ADN DMR Peer Server - tests voice plugin bridge
#
# 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
###############################################################################
"""VoicePluginBridge emit helpers and frame-only notify."""
from __future__ import annotations
from adn_server.application.plugins.application.bus import PluginBus
from adn_server.application.plugins.application.voice_bridge import VoicePluginBridge
from adn_server.application.plugins.domain.events import VoiceCallEnd, VoiceCallFrame, VoiceCallStart
from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE
def _stream_id(n: int) -> bytes:
return n.to_bytes(4, "big")
def _rf(n: int) -> bytes:
return n.to_bytes(3, "big")
def _tg(n: int) -> bytes:
return n.to_bytes(3, "big")
def _bridge() -> tuple[VoicePluginBridge, list[object]]:
events: list[object] = []
bus = PluginBus()
bus.subscribe(events.append)
config = {
"GLOBAL": {"SERVER_ID": 1},
"SYSTEMS": {"OBP-CL": {"MODE": "OPENBRIDGE"}},
}
return VoicePluginBridge(bus, config), events
def test_emit_group_voice_start_is_idempotent() -> None:
bridge, events = _bridge()
kw = dict(
system_name="OBP-CL",
peer_id=b"\x00\x00\x01\x00",
rf_src=_rf(0x334202),
dst_id=_tg(0x33420),
slot=1,
stream_id=_stream_id(42),
pkt_time=100.0,
source_is_obp=True,
voice_phase="INGRESS",
)
assert bridge.emit_group_voice_start(**kw) is True
assert bridge.emit_group_voice_start(**kw) is False
assert len(events) == 1
assert isinstance(events[0], VoiceCallStart)
assert events[0].context.extra.get("voice_phase") == "INGRESS"
def test_emit_group_voice_end_is_idempotent() -> None:
bridge, events = _bridge()
kw = dict(
system_name="OBP-CL",
peer_id=b"\x00\x00\x01\x00",
rf_src=_rf(0x334202),
dst_id=_tg(0x33420),
slot=1,
stream_id=_stream_id(42),
pkt_time=110.0,
duration_s=10.0,
source_is_obp=True,
)
assert bridge.emit_group_voice_end(**kw) is True
assert bridge.emit_group_voice_end(**kw) is False
assert len(events) == 1
assert isinstance(events[0], VoiceCallEnd)
def test_notify_group_after_forward_emits_frames_only() -> None:
bridge, events = _bridge()
bridge.notify_group_after_forward(
system_name="OBP-CL",
peer_id=b"\x00\x00\x01\x00",
rf_src=_rf(0x334202),
dst_id=_tg(0x33420),
slot=1,
stream_id=_stream_id(42),
frame_type=HBPF_VOICE,
dtype_vseq=1,
data=b"\x00" * 60,
pkt_time=101.0,
source_is_obp=True,
obp_hops=b"\x00",
obp_source_server=None,
obp_ber=b"\x00",
obp_rssi=b"\x00",
obp_source_rptr=b"\x00",
synthetic_announcement=False,
forwarded=["HBP-1"],
)
assert len(events) == 1
assert isinstance(events[0], VoiceCallFrame)
bridge.notify_group_after_forward(
system_name="OBP-CL",
peer_id=b"\x00\x00\x01\x00",
rf_src=_rf(0x334202),
dst_id=_tg(0x33420),
slot=1,
stream_id=_stream_id(42),
frame_type=HBPF_DATA_SYNC,
dtype_vseq=HBPF_SLT_VHEAD,
data=b"\x00" * 60,
pkt_time=101.0,
source_is_obp=True,
obp_hops=b"\x00",
obp_source_server=None,
obp_ber=b"\x00",
obp_rssi=b"\x00",
obp_source_rptr=b"\x00",
synthetic_announcement=False,
forwarded=[],
)
assert len(events) == 1

@ -56,7 +56,7 @@ def test_unit_data_reports_rx_event_when_reporting_enabled() -> None:
@pytest.mark.behavior @pytest.mark.behavior
def test_unit_data_csbk_new_stream_is_handled() -> None: def test_unit_data_csbk_new_stream_is_handled() -> None:
"""Regression: CSBK dtype 3 on new stream is accepted and reported.""" """Regression: CSBK dtype 3 on new stream is accepted but not monitor-reported."""
scenario = DeterministicScenario(enable_reporting=True) scenario = DeterministicScenario(enable_reporting=True)
base = PacketSpec(call_type="unit", dst_id=5003, stream_id=0x63636363, slot=2) base = PacketSpec(call_type="unit", dst_id=5003, stream_id=0x63636363, slot=2)
@ -66,7 +66,8 @@ def test_unit_data_csbk_new_stream_is_handled() -> None:
) )
assert_inject_ok(ok) assert_inject_ok(ok)
assert_report_event(scenario, "UNIT CSBK") assert scenario.report_factory is not None
assert not any("UNIT CSBK" in ev for ev in scenario.report_factory.events)
slot = scenario.protocols["MASTER-A"].STATUS.get(2, {}) slot = scenario.protocols["MASTER-A"].STATUS.get(2, {})
assert slot.get("RX_TYPE", HBPF_SLT_VTERM) == HBPF_SLT_VTERM assert slot.get("RX_TYPE", HBPF_SLT_VTERM) == HBPF_SLT_VTERM

Loading…
Cancel
Save

Powered by TurnKey Linux.