diff --git a/.gitignore b/.gitignore index c4056c3..1620cb1 100644 --- a/.gitignore +++ b/.gitignore @@ -36,6 +36,18 @@ json/* # Internal 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/ tests/fixtures/sessions/_capture/ diff --git a/docs/en/README.md b/docs/en/README.md index 0a0bf90..9e8c1f5 100644 --- a/docs/en/README.md +++ b/docs/en/README.md @@ -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 | |------------|------------| | 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) | | Private calls | [Private calls](server/user-guide/private-calls.md) | | Voice / TTS | [Voice, announcements, and TTS](server/user-guide/voice-and-tts.md) | diff --git a/docs/en/server/development/architecture.md b/docs/en/server/development/architecture.md index 087af46..3ed747d 100644 --- a/docs/en/server/development/architecture.md +++ b/docs/en/server/development/architecture.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). +**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`). - **`routing_table_for_report()`** — export shim for monitor/report (legacy BRIDGE_SND shape); not used for runtime forwards. diff --git a/docs/en/server/user-guide/plugins.md b/docs/en/server/user-guide/plugins.md new file mode 100644 index 0000000..3acbab2 --- /dev/null +++ b/docs/en/server/user-guide/plugins.md @@ -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//` loaded at runtime. + +The repository ships one reference skeleton: `plugins/example/` (disabled by default). + +--- + +## Directory layout + +``` +plugins// + 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: /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//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//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/.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. diff --git a/docs/es/README.md b/docs/es/README.md index 728bd2b..b75504f 100644 --- a/docs/es/README.md +++ b/docs/es/README.md @@ -20,6 +20,7 @@ El **ADN DMR Peer Server** es un puente de conferencia [GPL-3.0](https://www.gnu | Quiero… | Empieza aquí | |---------|----------------| | 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) | | Llamadas privadas | [Llamadas privadas](server/user-guide/private-calls.md) | | Voz / TTS | [Voz, anuncios y TTS](server/user-guide/voice-and-tts.md) | diff --git a/docs/es/server/development/architecture.md b/docs/es/server/development/architecture.md index 1db3780..59686a4 100644 --- a/docs/es/server/development/architecture.md +++ b/docs/es/server/development/architecture.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). +**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`). - **`routing_table_for_report()`** — shim de exportación para monitor/informes (forma legacy BRIDGE_SND); no se usa para reenvíos en runtime. diff --git a/docs/es/server/user-guide/plugins.md b/docs/es/server/user-guide/plugins.md new file mode 100644 index 0000000..8b92866 --- /dev/null +++ b/docs/es/server/user-guide/plugins.md @@ -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//` cargado en tiempo de ejecución. + +El repositorio incluye un skeleton de referencia: `plugins/example/` (deshabilitado por defecto). + +--- + +## Estructura de directorios + +``` +plugins// + 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: /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//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//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/.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. diff --git a/mkdocs.es.yml b/mkdocs.es.yml index 7fb3f8e..286d983 100644 --- a/mkdocs.es.yml +++ b/mkdocs.es.yml @@ -55,6 +55,7 @@ nav: - Guía de usuario: - Introducción: server/user-guide/introduction.md - Configuración: server/user-guide/configuration.md + - Plugins: server/user-guide/plugins.md - Talker Alias: server/user-guide/talker-alias.md - Bridges y talkgroups: server/user-guide/bridges-and-talkgroups.md - Números especiales (4000, 999x, eco): server/user-guide/special-numbers.md diff --git a/mkdocs.yml b/mkdocs.yml index 1d82929..1e1966e 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -55,6 +55,7 @@ nav: - User guide: - Introduction: server/user-guide/introduction.md - Configuration: server/user-guide/configuration.md + - Plugins: server/user-guide/plugins.md - Talker Alias: server/user-guide/talker-alias.md - Bridges and talkgroups: server/user-guide/bridges-and-talkgroups.md - Special numbers (4000, 999x, echo): server/user-guide/special-numbers.md diff --git a/plugins/.gitkeep b/plugins/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/plugins/README.md b/plugins/README.md new file mode 100644 index 0000000..d69ccd1 --- /dev/null +++ b/plugins/README.md @@ -0,0 +1,36 @@ +# Plugins + +Drop-in extensions live under `plugins//`. The server loads them at runtime from this directory (default: `/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// + 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 +``` diff --git a/plugins/example/README.md b/plugins/example/README.md new file mode 100644 index 0000000..75e9942 --- /dev/null +++ b/plugins/example/README.md @@ -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/.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). diff --git a/plugins/example/config.yaml b/plugins/example/config.yaml new file mode 100644 index 0000000..461c012 --- /dev/null +++ b/plugins/example/config.yaml @@ -0,0 +1,4 @@ +enabled: false + +# Relative to adn-server project root (created on write if missing). +output_dir: example-events diff --git a/plugins/example/plugin/__init__.py b/plugins/example/plugin/__init__.py new file mode 100644 index 0000000..ee5b280 --- /dev/null +++ b/plugins/example/plugin/__init__.py @@ -0,0 +1,29 @@ +# ADN DMR Peer Server plugin - example factory +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Example plugin — drop-in factory.""" + +from __future__ import annotations + +from .application.plugin_impl import ExamplePlugin + + +def create_plugin() -> ExamplePlugin: + return ExamplePlugin() diff --git a/plugins/example/plugin/application/plugin_impl.py b/plugins/example/plugin/application/plugin_impl.py new file mode 100644 index 0000000..7652c41 --- /dev/null +++ b/plugins/example/plugin/application/plugin_impl.py @@ -0,0 +1,129 @@ +# ADN DMR Peer Server plugin - example adapter +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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, + ) diff --git a/plugins/example/plugin/application/session_use_case.py b/plugins/example/plugin/application/session_use_case.py new file mode 100644 index 0000000..617aeb1 --- /dev/null +++ b/plugins/example/plugin/application/session_use_case.py @@ -0,0 +1,93 @@ +# ADN DMR Peer Server plugin - example session use case +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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() diff --git a/plugins/example/plugin/domain/record.py b/plugins/example/plugin/domain/record.py new file mode 100644 index 0000000..3df01fb --- /dev/null +++ b/plugins/example/plugin/domain/record.py @@ -0,0 +1,58 @@ +# ADN DMR Peer Server plugin - example record builder +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/plugins/example/plugin/infrastructure/json_writer.py b/plugins/example/plugin/infrastructure/json_writer.py new file mode 100644 index 0000000..97df1db --- /dev/null +++ b/plugins/example/plugin/infrastructure/json_writer.py @@ -0,0 +1,56 @@ +# ADN DMR Peer Server plugin - example JSON writer +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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) diff --git a/plugins/example/tests/conftest.py b/plugins/example/tests/conftest.py new file mode 100644 index 0000000..a9c1590 --- /dev/null +++ b/plugins/example/tests/conftest.py @@ -0,0 +1,28 @@ +# ADN DMR Peer Server plugin - example test path +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +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)) diff --git a/plugins/example/tests/test_record.py b/plugins/example/tests/test_record.py new file mode 100644 index 0000000..d120a6c --- /dev/null +++ b/plugins/example/tests/test_record.py @@ -0,0 +1,48 @@ +# ADN DMR Peer Server plugin - example tests +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +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") diff --git a/plugins/example/tests/test_session_use_case.py b/plugins/example/tests/test_session_use_case.py new file mode 100644 index 0000000..8379b94 --- /dev/null +++ b/plugins/example/tests/test_session_use_case.py @@ -0,0 +1,72 @@ +# ADN DMR Peer Server plugin - example session tests +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +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() diff --git a/src/adn_server/application/plugins/__init__.py b/src/adn_server/application/plugins/__init__.py new file mode 100644 index 0000000..e84cc8f --- /dev/null +++ b/src/adn_server/application/plugins/__init__.py @@ -0,0 +1,41 @@ +# ADN DMR Peer Server - plugin system package +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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", +] diff --git a/src/adn_server/application/plugins/application/__init__.py b/src/adn_server/application/plugins/application/__init__.py new file mode 100644 index 0000000..5fde489 --- /dev/null +++ b/src/adn_server/application/plugins/application/__init__.py @@ -0,0 +1,25 @@ +# ADN DMR Peer Server - plugin application package +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +from .bus import PluginBus +from .context import ServerContext +from .manager import PluginManager + +__all__ = ["PluginBus", "PluginManager", "ServerContext"] diff --git a/src/adn_server/application/plugins/application/bridge_common.py b/src/adn_server/application/plugins/application/bridge_common.py new file mode 100644 index 0000000..665097d --- /dev/null +++ b/src/adn_server/application/plugins/application/bridge_common.py @@ -0,0 +1,100 @@ +# ADN DMR Peer Server - plugin bridge helpers +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/src/adn_server/application/plugins/application/bus.py b/src/adn_server/application/plugins/application/bus.py new file mode 100644 index 0000000..f535952 --- /dev/null +++ b/src/adn_server/application/plugins/application/bus.py @@ -0,0 +1,110 @@ +# ADN DMR Peer Server - plugin event bus +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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) diff --git a/src/adn_server/application/plugins/application/context.py b/src/adn_server/application/plugins/application/context.py new file mode 100644 index 0000000..d294fe2 --- /dev/null +++ b/src/adn_server/application/plugins/application/context.py @@ -0,0 +1,35 @@ +# ADN DMR Peer Server - plugin server context +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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] diff --git a/src/adn_server/application/plugins/application/data_bridge.py b/src/adn_server/application/plugins/application/data_bridge.py new file mode 100644 index 0000000..8ca505e --- /dev/null +++ b/src/adn_server/application/plugins/application/data_bridge.py @@ -0,0 +1,162 @@ +# ADN DMR Peer Server - plugin unit data event bridge +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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)) diff --git a/src/adn_server/application/plugins/application/manager.py b/src/adn_server/application/plugins/application/manager.py new file mode 100644 index 0000000..5a54418 --- /dev/null +++ b/src/adn_server/application/plugins/application/manager.py @@ -0,0 +1,143 @@ +# ADN DMR Peer Server - plugin manager +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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) diff --git a/src/adn_server/application/plugins/application/voice_bridge.py b/src/adn_server/application/plugins/application/voice_bridge.py new file mode 100644 index 0000000..c90ed43 --- /dev/null +++ b/src/adn_server/application/plugins/application/voice_bridge.py @@ -0,0 +1,345 @@ +# ADN DMR Peer Server - plugin voice event bridge +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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, + ) + ) diff --git a/src/adn_server/application/plugins/domain/__init__.py b/src/adn_server/application/plugins/domain/__init__.py new file mode 100644 index 0000000..a29c671 --- /dev/null +++ b/src/adn_server/application/plugins/domain/__init__.py @@ -0,0 +1,39 @@ +# ADN DMR Peer Server - plugin domain package +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +from .events import ( + CallLegContext, + UnitDataEnd, + UnitDataFrame, + UnitDataStart, + VoiceCallEnd, + VoiceCallFrame, + VoiceCallStart, +) + +__all__ = [ + "CallLegContext", + "UnitDataEnd", + "UnitDataFrame", + "UnitDataStart", + "VoiceCallEnd", + "VoiceCallFrame", + "VoiceCallStart", +] diff --git a/src/adn_server/application/plugins/domain/events.py b/src/adn_server/application/plugins/domain/events.py new file mode 100644 index 0000000..d7d3658 --- /dev/null +++ b/src/adn_server/application/plugins/domain/events.py @@ -0,0 +1,132 @@ +# ADN DMR Peer Server - plugin voice events +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/src/adn_server/application/plugins/domain/protocol.py b/src/adn_server/application/plugins/domain/protocol.py new file mode 100644 index 0000000..6f29012 --- /dev/null +++ b/src/adn_server/application/plugins/domain/protocol.py @@ -0,0 +1,44 @@ +# ADN DMR Peer Server - plugin contract +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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//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.""" diff --git a/src/adn_server/application/plugins/infrastructure/__init__.py b/src/adn_server/application/plugins/infrastructure/__init__.py new file mode 100644 index 0000000..f854260 --- /dev/null +++ b/src/adn_server/application/plugins/infrastructure/__init__.py @@ -0,0 +1,35 @@ +# ADN DMR Peer Server - plugin infrastructure package +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +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", +] diff --git a/src/adn_server/application/plugins/infrastructure/loader.py b/src/adn_server/application/plugins/infrastructure/loader.py new file mode 100644 index 0000000..452587b --- /dev/null +++ b/src/adn_server/application/plugins/infrastructure/loader.py @@ -0,0 +1,143 @@ +# ADN DMR Peer Server - plugin loader +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Discover and load drop-in plugins from plugins//.""" + +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//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 diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index 4de9252..ace28ef 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -52,13 +52,8 @@ _CSV_FAMILIES = { "GROUP VOICE": "GROUP", "PRIVATE VOICE": "PRIVATE", "UNIT DATA": "UNIT", - # routing_use_cases._dtype_labels emits one of these instead of the generic - # "UNIT DATA" label for dtype_vseq in (3, 6, 7, 8) — private-call CSBK setup - # 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. + # routing_use_cases emits these labels for dtype_vseq in (6, 7, 8). CSBK (3) is + # still logged at INFO but no longer emitted to monitor (prelude storm filter). "UNIT CSBK": "UNIT", "UNIT DATA HEADER": "UNIT", "UNIT VCSBK 1/2 DATA BLOCK": "UNIT", diff --git a/src/adn_server/application/routing/hbp_forward.py b/src/adn_server/application/routing/hbp_forward.py index 049afb3..65fe3a7 100644 --- a/src/adn_server/application/routing/hbp_forward.py +++ b/src/adn_server/application/routing/hbp_forward.py @@ -53,6 +53,7 @@ from .helpers import ( hbp_ingress_new_stream_collision, master_per_peer_slot_contention, tg_has_active_conversation, + unit_data_reportable, ) from .peer_downlink_index import count_connected_peers @@ -384,8 +385,10 @@ class HbpForwardMixin: "(%s) UNIT trace -> %s stream=%s payload33=%s", source_system, d_system, int_id(stream_id), dmrpkt.hex(), ) - 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), + _dtype_vseq = data[15] & 0xF if len(data) > 15 else 0 + if unit_data_reportable(_dtype_vseq): + 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), + ) ) - ) diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index 3256058..4a3d40f 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -815,6 +815,90 @@ def _obp_stream_active(st: dict[str, Any], pkt_time: float) -> bool: 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( protocols: 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( call_type: str, dtype_vseq: int, diff --git a/src/adn_server/application/routing/obp_forward.py b/src/adn_server/application/routing/obp_forward.py index 91b6c1b..6a78ca3 100644 --- a/src/adn_server/application/routing/obp_forward.py +++ b/src/adn_server/application/routing/obp_forward.py @@ -52,7 +52,7 @@ from typing import Any from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, int_id from ...domain.dmr import decode 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__) @@ -152,6 +152,7 @@ class ObpForwardMixin: stream_id: bytes, data: bytes, obp_hops: bytes, + synthetic_announcement: bool = False, ) -> bool: """Port of bridge_master.routerOBP.dmrd_received group/vcsbk (~2269-2411). False = drop packet.""" 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) ) ) + 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: st = status[stream_id] if "packets" in st: @@ -473,8 +491,10 @@ class ObpForwardMixin: logger.warning("(%s) send_data_to_obp %s failed: %s", source_system, target, exc) return logger.debug("(%s) UNIT Data Bridged to OBP System: %s DST_ID: %s", source_system, target, 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), + _dtype_vseq = data[15] & 0xF if len(data) > 15 else 0 + if unit_data_reportable(_dtype_vseq): + 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), + ) ) - ) diff --git a/src/adn_server/application/routing/timers.py b/src/adn_server/application/routing/timers.py index 260df00..abc70c5 100644 --- a/src/adn_server/application/routing/timers.py +++ b/src/adn_server/application/routing/timers.py @@ -296,6 +296,20 @@ class RoutingTimerMixin: 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 self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, now) continue @@ -331,6 +345,19 @@ class RoutingTimerMixin: _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: _slot["RX_STREAM_ID"] = b"\x00" tx_streams = _slot.get("TX_STREAMS") diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 5106a76..2a49fe2 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -36,7 +36,17 @@ from typing import Any 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 .ports import AclRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder from .reporting_use_cases import ReportingUseCases @@ -51,11 +61,13 @@ from .routing.helpers import ( obp_deferred_bridge_tx_leg, obp_flat_bridge_tx_idle, obp_publish_flat_bridge_tx, + obp_status_plugin_voice, obp_sync_flat_bridge_tx_times, obp_target_bcsq_quenches_stream, resolve_voice_peer_id, slot_has_active_voice, unit_data_hbp_target_idle, + unit_data_reportable, ) from .routing.lc_ta import LcTaMixin from .routing.obp_forward import ObpForwardMixin @@ -96,6 +108,8 @@ class RoutingUseCases( call_later: Any = None, encode_emblc: DmrEmbeddedLcEncoder | None = None, ta_emblc_encoder: TalkerAliasEmblcEncoder | None = None, + voice_plugin_bridge: Any = None, + data_plugin_bridge: Any = None, ) -> None: self._acl_router = acl_router self._config = config @@ -114,6 +128,8 @@ class RoutingUseCases( raise TypeError("encode_emblc and ta_emblc_encoder are required (wire from main.py)") self._encode_emblc = encode_emblc 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} self._both_ta_wait: dict[tuple[str, bytes], dict[str, Any]] = {} # Passthrough DMRA/embed relay already applied for this source stream. @@ -281,6 +297,7 @@ class RoutingUseCases( stream_id, data, obp_hops if obp_use_parsed else b"", + synthetic_announcement=synthetic_announcement, ): return # SINGLE_MODE in-band VTERM on another TG (e.g. 9990 echo) deactivates static OFF @@ -321,6 +338,18 @@ class RoutingUseCases( 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: if not _obp_grp: duration = 0.0 @@ -361,6 +390,19 @@ class RoutingUseCases( 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( self._voice_relay_tables_with_active_source(system_name, bridge_match_slot, dst_int) ) @@ -951,6 +993,26 @@ class RoutingUseCases( 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 self._obp_emit_end_tx_for_forward_legs(stream_id, system_name, _end_t) ost["lastSeq"] = False @@ -960,6 +1022,35 @@ class RoutingUseCases( "(ROUTER) Bridged TG %s from %s -> %s", 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 # ── 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"" _bits = data[15] if len(data) > 15 else 0 _int_dst_id = int_id(dst_id) + _forwarded: list[str] = [] systems_cfg = self._config.get("SYSTEMS", {}) source_is_obp = systems_cfg.get(system_name, {}).get("MODE") == "OPENBRIDGE" 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, ) - _dtype_labels = {3: "UNIT CSBK", 6: "UNIT DATA HEADER", 7: "UNIT VCSBK 1/2 DATA BLOCK", 8: "UNIT VCSBK 3/4 DATA BLOCK"} - _label = _dtype_labels.get(dtype_vseq, "UNIT DATA") - self._send_routing_event( - "{},DATA,RX,{},{},{},{},{},{}".format( - _label, system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, _int_dst_id, + if unit_data_reportable(dtype_vseq): + _dtype_labels = {6: "UNIT DATA HEADER", 7: "UNIT VCSBK 1/2 DATA BLOCK", 8: "UNIT VCSBK 3/4 DATA BLOCK"} + _label = _dtype_labels.get(dtype_vseq, "UNIT DATA") + self._send_routing_event( + "{},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) if global_cfg.get("DATA_GATEWAY"): dg_cfg = systems_cfg.get("DATA-GATEWAY", {}) if dg_cfg.get("MODE") == "OPENBRIDGE" and dg_cfg.get("ENABLED"): logger.debug("(%s) DATA packet sent to DATA-GATEWAY", system_name) + _forwarded.append("DATA-GATEWAY") self._send_data_to_obp( system_name, "DATA-GATEWAY", data, dmrpkt, pkt_time, stream_id, dst_id, peer_id, rf_src, _bits, slot, @@ -1151,6 +1245,7 @@ class RoutingUseCases( if sys_name == "DATA-GATEWAY": continue 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( system_name, sys_name, data, dmrpkt, pkt_time, stream_id, dst_id, peer_id, rf_src, _bits, slot, @@ -1172,6 +1267,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _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) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( @@ -1199,6 +1295,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _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) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( @@ -1214,6 +1311,7 @@ class RoutingUseCases( _hangtime = _d_sys_cfg.get("GROUP_HANGTIME", 5) if unit_data_hbp_target_idle(_dst_slot, pkt_time, _hangtime): _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) elif not is_private_subscriber_dst(dst_id): self._log_unit_data_hbp_busy( @@ -1224,6 +1322,29 @@ class RoutingUseCases( if _matched: 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( self, system_name: str, @@ -1251,6 +1372,7 @@ class RoutingUseCases( if source_proto: source_status.setdefault(slot, {}) 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", 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, ) return + _notify_pvt_start = True slot_st["RX_START"] = pkt_time self._pvt_same_system_dst_peer_id = None if dst_id in sub_map: @@ -1397,6 +1520,28 @@ class RoutingUseCases( self._send_to_system(_target, send_data) except Exception as 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: self._pvt_targets = [] 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 + 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 _unit_data: # Keep stream continuity for multi-frame unit data without marking the diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index 9afb12d..5483e8b 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -38,6 +38,11 @@ from adn_server.application import ( VoiceUseCases, ) 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 ( is_obp_proxy_managed, is_proxy_inject_only, @@ -431,6 +436,19 @@ def run_peer_server( 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( acl_router, config, @@ -445,7 +463,11 @@ def run_peer_server( call_later=reactor.callLater, encode_emblc=encode_emblc, 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() dynamic_tg_uc = DynamicTgUseCases( dynamic_tg_store, @@ -756,6 +778,9 @@ def run_peer_server( obp_proxy_state = start_obp_proxy_service(config, protocols, logger=logger) if result.added or result.removed or result.updated or result.rebound: _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: nonlocal report_mqtt, proxy_state, obp_proxy_state diff --git a/src/adn_server/infrastructure/logging_config.py b/src/adn_server/infrastructure/logging_config.py index 9a26557..248a632 100644 --- a/src/adn_server/infrastructure/logging_config.py +++ b/src/adn_server/infrastructure/logging_config.py @@ -72,6 +72,17 @@ def reopen_file_handlers(logger: logging.Logger | None = None) -> int: 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: """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) for handler in app_logger.handlers: handler.setLevel(level) + _silence_noisy_library_loggers() 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) logger = logging.getLogger(log_name) logger.setLevel(level) + _silence_noisy_library_loggers() return logger diff --git a/tests/application/test_bridge_common.py b/tests/application/test_bridge_common.py new file mode 100644 index 0000000..201d426 --- /dev/null +++ b/tests/application/test_bridge_common.py @@ -0,0 +1,35 @@ +# ADN DMR Peer Server - tests plugin bridge helpers +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +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" diff --git a/tests/application/test_obp_canonical_ingress.py b/tests/application/test_obp_canonical_ingress.py new file mode 100644 index 0000000..3d82fcd --- /dev/null +++ b/tests/application/test_obp_canonical_ingress.py @@ -0,0 +1,98 @@ +# ADN DMR Peer Server - tests OBP canonical ingress helpers +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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), + ) diff --git a/tests/application/test_plugin_bus.py b/tests/application/test_plugin_bus.py new file mode 100644 index 0000000..7bebc01 --- /dev/null +++ b/tests/application/test_plugin_bus.py @@ -0,0 +1,82 @@ +# ADN DMR Peer Server - tests plugin bus +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/tests/application/test_stream_talker_alias.py b/tests/application/test_stream_talker_alias.py new file mode 100644 index 0000000..3d8f062 --- /dev/null +++ b/tests/application/test_stream_talker_alias.py @@ -0,0 +1,65 @@ +# ADN DMR Peer Server - tests stream talker alias helper +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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" diff --git a/tests/application/test_voice_bridge.py b/tests/application/test_voice_bridge.py new file mode 100644 index 0000000..b404889 --- /dev/null +++ b/tests/application/test_voice_bridge.py @@ -0,0 +1,137 @@ +# ADN DMR Peer Server - tests voice plugin bridge +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/tests/routing/test_unit_data_ingress.py b/tests/routing/test_unit_data_ingress.py index 988b643..47a47fe 100644 --- a/tests/routing/test_unit_data_ingress.py +++ b/tests/routing/test_unit_data_ingress.py @@ -56,7 +56,7 @@ def test_unit_data_reports_rx_event_when_reporting_enabled() -> None: @pytest.mark.behavior 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) 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_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, {}) assert slot.get("RX_TYPE", HBPF_SLT_VTERM) == HBPF_SLT_VTERM