From 62f429aad90272e1684f554ac8a1999cf15eed43 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 13 Jul 2026 15:48:53 -0400 Subject: [PATCH] fix: inject scheduled voice via routing and restore monitor TX events Route announcements/TTS through synthetic PTT on the proxy MASTER (SERVER_ID peer, normal dmrd_received forwarding). Emit START/END TX report events for inject so monitor fans out to SYSTEM-N; configurable server voice DMR_ID. --- adn-voice.example.yaml | 9 + docs/en/server/user-guide/special-numbers.md | 26 +- docs/en/server/user-guide/voice-and-tts.md | 25 +- docs/es/server/user-guide/special-numbers.md | 26 +- docs/es/server/user-guide/voice-and-tts.md | 25 +- src/adn_server/application/ident_use_cases.py | 3 +- .../routing/announcement_ptt_inject.py | 99 +++++ src/adn_server/application/routing/helpers.py | 42 +- .../application/routing_use_cases.py | 7 +- src/adn_server/application/server_voice.py | 114 ++++++ src/adn_server/application/voice_use_cases.py | 386 ++++++++++++++---- .../infrastructure/bootstrap/peer_server.py | 30 ++ .../twisted_adapters/udp_hbp.py | 11 +- .../infrastructure/voice/tts_engine.py | 144 ++++++- ...test_master_server_broadcast_contention.py | 23 +- tests/application/test_monitor_topology.py | 17 + tests/application/test_server_voice.py | 105 +++++ tests/harness/voice_helpers.py | 32 +- .../test_hbp_repeat_options_filter.py | 7 +- tests/routing/test_announcement_ptt_inject.py | 77 ++++ .../test_obp_announcement_contention.py | 21 +- .../voice/test_announcement_anticollision.py | 5 +- tests/voice/test_broadcast_queue.py | 5 +- tests/voice/test_inject_ptt_peer.py | 122 ++++++ tests/voice/test_scheduled_announcement.py | 2 + tests/voice/test_scheduled_tts.py | 30 +- tests/voice/test_tts_engine_serial.py | 123 ++++++ tests/voice/test_voice_config_reload.py | 26 ++ 28 files changed, 1397 insertions(+), 145 deletions(-) create mode 100644 src/adn_server/application/routing/announcement_ptt_inject.py create mode 100644 src/adn_server/application/server_voice.py create mode 100644 tests/application/test_server_voice.py create mode 100644 tests/routing/test_announcement_ptt_inject.py create mode 100644 tests/voice/test_inject_ptt_peer.py create mode 100644 tests/voice/test_tts_engine_serial.py diff --git a/adn-voice.example.yaml b/adn-voice.example.yaml index f9554a2..2c690d2 100644 --- a/adn-voice.example.yaml +++ b/adn-voice.example.yaml @@ -7,6 +7,11 @@ # Each announcement/TTS has its own LANGUAGE. ENABLED: true activates that item. VOICE: + # Optional — omit on legacy installs; default DMR_ID 1000001. + # Per-item DMR_ID on ANNOUNCEMENTS / TTS_ANNOUNCEMENTS overrides VOICE.DMR_ID. + # Callsign resolves from subscriber DB for that DMR ID. + # DMR_ID: 1000001 + # Voice ident only (MASTER systems with VOICE_IDENT in adn-server.yaml). # Leave empty: ident auto-detects from systems. Or set "es_ES,en_GB" to preload. # NOT used by ANNOUNCEMENTS or TTS — those use each item's LANGUAGE. @@ -21,11 +26,13 @@ VOICE: # FILE: filename without .ambe, looked up in Audio//ondemand/ # MODE: "hourly" (once per hour) or "interval" (every INTERVAL seconds) # ENABLED: true = plays; false = configured but inactive (edit to true when ready) + # DMR_ID: optional RF source for this clip (defaults to VOICE.DMR_ID) # --------------------------------------------------------------------------- ANNOUNCEMENTS: - ENABLED: true FILE: announcement1 TG: 2 + # DMR_ID: 1000001 MODE: interval INTERVAL: 60 LANGUAGE: es_ES @@ -77,12 +84,14 @@ VOICE: - ENABLED: true FILE: texto1 TG: 2 + # DMR_ID: 1000001 MODE: interval INTERVAL: 60 LANGUAGE: es_ES - ENABLED: false FILE: texto2 TG: 2 + # DMR_ID: 1000001 MODE: hourly INTERVAL: 3600 LANGUAGE: es_ES diff --git a/docs/en/server/user-guide/special-numbers.md b/docs/en/server/user-guide/special-numbers.md index 6c10220..fe9142c 100644 --- a/docs/en/server/user-guide/special-numbers.md +++ b/docs/en/server/user-guide/special-numbers.md @@ -2,19 +2,21 @@ Several **destination IDs** are reserved for **control or services**. They are handled in protocol layers and/or the bridge router, not as normal group traffic. -## ID 5000 — server voice source (not “announcement TG”) {#id-5000--server-voice-source-not-announcement-tg} +## Server voice source ID (configurable, not “announcement TG”) {#server-voice-source-id-configurable} -**Important:** **5000** is the **RF source ID** the server uses when it **transmits** automated voice. Radios and dashboards show **caller ID 5000** for that traffic. +**Important:** The server uses a configurable **DMR ID** (`VOICE.DMR_ID` in `adn-voice.yaml`). Default: **1000001**. Each announcement/TTS row may set its own **`DMR_ID`**; when omitted, it inherits `VOICE.DMR_ID`. Callsign comes from the subscriber DB for that ID. | Traffic | Destination in the DMR packet | Notes | |---------|----------------------------------|--------| -| **Scheduled AMBE** (`ANNOUNCEMENTS`) | Whatever **`TG`** you set in `adn-voice.yaml` | Source ID **5000**. | -| **TTS** (`TTS_ANNOUNCEMENTS`) | Same — configured **`TG`** | Source ID **5000**. | -| **On-demand** (TG **9991–9999**) | **TG 9** | Short info clips; source **5000** (see [Voice, announcements, and TTS](voice-and-tts.md)). | -| **Disconnected / reflector prompts** | **TG 9** | Source **5000**. | -| **Voice ident** | **All-call** (`16777215`) or **`OVERRIDE_IDENT_TG`** if set | Source **5000**. | +| **Scheduled AMBE** (`ANNOUNCEMENTS`) | Whatever **`TG`** you set in `adn-voice.yaml` | Item `DMR_ID` or `VOICE.DMR_ID` (default **1000001**). | +| **TTS** (`TTS_ANNOUNCEMENTS`) | Same — configured **`TG`** | Same. | +| **On-demand** (TG **9991–9999**) | **TG 9** | Short info clips; configurable source (see [Voice, announcements, and TTS](voice-and-tts.md)). | +| **Disconnected / reflector prompts** | **TG 9** | Same. | +| **Voice ident** | **All-call** (`16777215`) or **`OVERRIDE_IDENT_TG`** if set | Same. | -You do **not** “monitor TG 5000” to hear scheduled announcements: you monitor the **configured announcement TG** (e.g. 2, 9, 26811). **5000** appears as the **transmitter ID** on those calls. +You do **not** “monitor” the server voice ID to hear scheduled announcements: you monitor the **configured announcement TG** (e.g. 2, 9, 26811). The configured ID appears as the **transmitter** on those calls. + +The legacy local ID **5000** is still recognised on some on-demand playback paths for compatibility, but the new default is **1000001** (avoids `SUB_ACL: DENY:0-1000000` blocks on OBP). ### Destination TG 5000 (inbound group) @@ -28,7 +30,7 @@ If a group call arrives with **destination TG 5000** and there is **no** existin For **on-demand** playback (after you key **9991–9999**) and for **disconnected / reflector** voice lines, the server transmits **group** packets with: -- **Source ID 5000** +- **DMR ID** = item `DMR_ID` or `VOICE.DMR_ID` (default **1000001**) - **Destination TG 9** - **Timeslot 2** (the code drives the **TS2** slot for that hotspot) @@ -92,7 +94,7 @@ Operationally: if users report “bridges drop too easily” after OPTIONS updat - Trigger path is **private VTERM** for destination **9991–9999**, then async playback generation. - Works from **MASTER** and **PEER** paths. -The **audio** is sent with **source ID 5000** and **destination TG 9** in the generated stream. File layout: [Voice, announcements, and TTS](voice-and-tts.md). +The **audio** is sent with the **configured source ID** and **destination TG 9** in the generated stream. File layout: [Voice, announcements, and TTS](voice-and-tts.md). ## TG 9990 — echo (in-band) @@ -118,11 +120,11 @@ Many small IDs (0–5, 9, etc.) and the **999x service** range are excluded from | ID / range | Role | |------------|------| -| **5000** (source) | Server-generated voice (announcements, TTS, prompts, ident) — **caller ID** on receivers | +| **VOICE.DMR_ID** (source) | Server-generated voice — **caller ID** on receivers (default **1000001**) | | **5000** (destination) | No auto user-activated bridge if missing from `BRIDGES` | | **4000** (group) | Deactivate dynamic bridges | | **4000** (unit) | Disconnect dynamics; not routed as PC | -| **9991–9999** | On-demand / information audio (trigger TG); playback uses src **5000** → TG **9** (TS2) | +| **9991–9999** | On-demand / information audio (trigger TG); playback uses `VOICE.DMR_ID` → TG **9** (TS2) | | **9** | Service/prompt lane (short server audio); reserved for auto-bridges; internal **TGID** on TS2 legs | | **9990** | Echo bridge TG (with ECHO system) | | **16777215** | All-call (default voice-ident destination unless overridden) | diff --git a/docs/en/server/user-guide/voice-and-tts.md b/docs/en/server/user-guide/voice-and-tts.md index d8f8076..7180dc6 100644 --- a/docs/en/server/user-guide/voice-and-tts.md +++ b/docs/en/server/user-guide/voice-and-tts.md @@ -7,9 +7,26 @@ Template: `adn-voice.example.yaml`. -## Server voice identity (ID 5000) +## Server voice identity (configurable `DMR_ID`) -All **server-originated** voice uses **RF source ID 5000** in the DMR stream so clients can tell infrastructure traffic from user radios: +All **server-originated** voice uses a configurable **DMR ID** in `adn-voice.yaml`: + +```yaml +VOICE: + DMR_ID: 1000001 # optional; default 1000001 + + TTS_ANNOUNCEMENTS: + - ENABLED: true + FILE: texto1 + TG: 730500 + DMR_ID: 3109898 # optional; inherits VOICE.DMR_ID when omitted +``` + +- **`VOICE.DMR_ID`** — default RF source. **Optional** — existing `adn-voice.yaml` files need no change. +- **Per-item `DMR_ID`** — override on each `ANNOUNCEMENTS` / `TTS_ANNOUNCEMENTS` row. Also optional. +- **Callsign** — resolved from the subscriber DB (`users` / `SUB_MAP`) for that DMR ID; not in voice config. + +Applies to: - Scheduled **ANNOUNCEMENTS** and **TTS_ANNOUNCEMENTS** (destination = the **`TG`** field in each item). - **On-demand** clips triggered by dialling **9991–9999** (playback uses destination **TG 9** on **TS2**; you still key **999x** to request the file). @@ -18,7 +35,7 @@ All **server-originated** voice uses **RF source ID 5000** in the DMR stream so See [TG 9 — local service lane](special-numbers.md#tg-9-local-service-lane-prompts-and-bridge-plumbing) for why TG 9 is used and what must be enabled on the hotspot. - **Voice ident** (destination **all-call** or **`OVERRIDE_IDENT_TG`**). -See [Special numbers — ID 5000](special-numbers.md#id-5000--server-voice-source-not-announcement-tg) for the full table. +See [Special numbers — server voice source ID](special-numbers.md#server-voice-source-id-configurable) for the full table. ## Features @@ -45,6 +62,8 @@ Configure **`TTS_VOCODER_CMD`** or **`TTS_AMBESERVER_HOST`** / **`TTS_AMBESERVER When scheduling announcements, the server may **wait** if target slots are busy, **drop** targets if a live QSO appears mid-stream, and only mark **hourly** announcement state after a successful target list — this avoids clobbering live traffic. +Scheduled and TTS broadcasts inject each frame as a **synthetic hotspot PTT** on `PROXY.TARGET_SYSTEM` (the same MASTER used by the integrated proxy). Routing fans out to bridged hotspots and OPENBRIDGE legs; local peers on that MASTER still receive frames via `send_system`. + ## Broadcast queue Parallel **broadcasts** on **different** TGs may run concurrently; **same-TG** broadcasts are **serialised** so one announcement completes before another on that TG. diff --git a/docs/es/server/user-guide/special-numbers.md b/docs/es/server/user-guide/special-numbers.md index 3934847..84f3df9 100644 --- a/docs/es/server/user-guide/special-numbers.md +++ b/docs/es/server/user-guide/special-numbers.md @@ -2,19 +2,21 @@ Varios **IDs de destino** están reservados para **control o servicios**. Se gestionan en capas de protocolo y/o en el router de bridges, no como tráfico de grupo normal. -## ID 5000 — fuente de voz del servidor (no es «TG de anuncio») {#id-5000--server-voice-source-not-announcement-tg} +## ID de voz del servidor — fuente configurable (no es «TG de anuncio») {#server-voice-source-id-configurable} -**Importante:** **5000** es el **ID de fuente RF** que el servidor usa cuando **transmite** voz automatizada. Las radios y paneles muestran **ID de llamada 5000** en ese tráfico. +**Importante:** El servidor usa un **DMR ID** configurable (`VOICE.DMR_ID` en `adn-voice.yaml`). Predeterminado: **1000001**. Cada anuncio/TTS puede definir **`DMR_ID`** propio; si se omite, hereda `VOICE.DMR_ID`. El callsign lo resuelve la base de suscriptores para ese ID. | Tráfico | Destino en el paquete DMR | Notas | |---------|---------------------------|--------| -| **AMBE programado** (`ANNOUNCEMENTS`) | El **`TG`** que definas en `adn-voice.yaml` | ID de fuente **5000**. | -| **TTS** (`TTS_ANNOUNCEMENTS`) | Igual — **`TG`** configurado | ID de fuente **5000**. | -| **Bajo demanda** (TG **9991–9999**) | **TG 9** | Clips informativos cortos; fuente **5000** (ver [Voz, anuncios y TTS](voice-and-tts.md)). | -| **Desconectado / reflector** | **TG 9** | Fuente **5000**. | -| **Ident por voz** | **All-call** (`16777215`) o **`OVERRIDE_IDENT_TG`** si está definido | Fuente **5000**. | +| **AMBE programado** (`ANNOUNCEMENTS`) | El **`TG`** que definas en `adn-voice.yaml` | `DMR_ID` del ítem o `VOICE.DMR_ID` (predeterminado **1000001**). | +| **TTS** (`TTS_ANNOUNCEMENTS`) | Igual — **`TG`** configurado | Igual. | +| **Bajo demanda** (TG **9991–9999**) | **TG 9** | Clips informativos cortos; fuente configurable (ver [Voz, anuncios y TTS](voice-and-tts.md)). | +| **Desconectado / reflector** | **TG 9** | Igual. | +| **Ident por voz** | **All-call** (`16777215`) o **`OVERRIDE_IDENT_TG`** si está definido | Igual. | -**No** «monitorizas TG 5000» para oír anuncios programados: monitorizas el **TG de anuncio configurado** (p. ej. 2, 9, 26811). **5000** aparece como **ID del transmisor** en esas llamadas. +**No** «monitorizas» el ID de voz del servidor para oír anuncios programados: monitorizas el **TG de anuncio configurado** (p. ej. 2, 9, 26811). El ID configurado aparece como **transmisor** en esas llamadas. + +El ID local heredado **5000** sigue reconocido en algunos caminos de reproducción bajo demanda por compatibilidad, pero el predeterminado nuevo es **1000001** (evita bloqueos con `SUB_ACL: DENY:0-1000000` en OBP). ### TG de destino 5000 (grupo entrante) @@ -28,7 +30,7 @@ Si llega una llamada de grupo con **TG de destino 5000** y **no** hay fila `BRID Para reproducción **bajo demanda** (tras marcar **9991–9999**) y para líneas de voz de **desconectado / reflector**, el servidor transmite paquetes de **grupo** con: -- **ID de fuente 5000** +- **DMR ID** = `DMR_ID` del ítem o `VOICE.DMR_ID` (predeterminado **1000001**) - **TG de destino 9** - **Timeslot 2** (el código usa el slot **TS2** para ese hotspot) @@ -92,7 +94,7 @@ A nivel operativo: si usuarios reportan «bridges que se caen demasiado fácil» - La ruta de disparo es **private VTERM** para destino **9991–9999**, seguida de generación/reproducción asíncrona. - Funciona desde rutas **MASTER** y **PEER**. -El **audio** se envía con **ID de fuente 5000** y **TG de destino 9** en el flujo generado. Estructura de ficheros: [Voz, anuncios y TTS](voice-and-tts.md). +El **audio** se envía con el **ID de fuente configurado** y **TG de destino 9** en el flujo generado. Estructura de ficheros: [Voz, anuncios y TTS](voice-and-tts.md). ## TG 9990 — eco (en banda) @@ -118,11 +120,11 @@ Muchos IDs pequeños (0–5, 9, etc.) y el rango de **servicio 999x** quedan exc | ID / rango | Rol | |------------|-----| -| **5000** (fuente) | Voz generada por el servidor (anuncios, TTS, mensajes, ident) — **ID de llamada** en receptores | +| **VOICE.DMR_ID** (fuente) | Voz generada por el servidor — **ID de llamada** en receptores (predeterminado **1000001**) | | **5000** (destino) | Sin bridge UA automático si falta en `BRIDGES` | | **4000** (grupo) | Desactivar bridges dinámicos | | **4000** (unitaria) | Desconectar dinámicos; no enrutada como PC | -| **9991–9999** | Audio informativo / bajo demanda (TG de disparo); reproducción usa fuente **5000** → TG **9** (TS2) | +| **9991–9999** | Audio informativo / bajo demanda (TG de disparo); reproducción usa `VOICE.DMR_ID` → TG **9** (TS2) | | **9** | Carril de servicio/mensajes (audio corto del servidor); reservada para auto-bridges; **TGID** interno en patas TS2 | | **9990** | TG de bridge de eco (con sistema ECHO) | | **16777215** | All-call (destino por defecto de ident por voz salvo sobrescritura) | diff --git a/docs/es/server/user-guide/voice-and-tts.md b/docs/es/server/user-guide/voice-and-tts.md index f700589..bfc5705 100644 --- a/docs/es/server/user-guide/voice-and-tts.md +++ b/docs/es/server/user-guide/voice-and-tts.md @@ -7,9 +7,26 @@ Plantilla: `adn-voice.example.yaml`. -## Identidad de voz del servidor (ID 5000) +## Identidad de voz del servidor (`DMR_ID` configurable) -Toda la voz **originada por el servidor** usa **ID de fuente RF 5000** en el flujo DMR para que los clientes distingan el tráfico de infraestructura del de usuarios: +Toda la voz **originada por el servidor** usa un **DMR ID** configurable en `adn-voice.yaml`: + +```yaml +VOICE: + DMR_ID: 1000001 # opcional; predeterminado 1000001 + + TTS_ANNOUNCEMENTS: + - ENABLED: true + FILE: texto1 + TG: 730500 + DMR_ID: 3109898 # opcional; hereda VOICE.DMR_ID si se omite +``` + +- **`VOICE.DMR_ID`** — ID RF predeterminado. **Opcional** — sin cambios en `adn-voice.yaml` existente. +- **`DMR_ID` por ítem** — override en cada fila de `ANNOUNCEMENTS` / `TTS_ANNOUNCEMENTS`. También opcional. +- **Callsign** — lo resuelve la base de suscriptores (`users` / `SUB_MAP`) para ese DMR ID; no va en voice config. + +Aplica a: - **ANNOUNCEMENTS** y **TTS_ANNOUNCEMENTS** programados (destino = el campo **`TG`** de cada ítem). - Clips **bajo demanda** disparados marcando **9991–9999** (la reproducción usa TG de destino **9** en **TS2**; sigues marcando **999x** para solicitar el fichero). @@ -19,7 +36,7 @@ Ver [TG 9 — carril de servicio local](special-numbers.md#tg-9-local-service-la - **Ident por voz** (destino **all-call** o **`OVERRIDE_IDENT_TG`**). -Ver [Números especiales — ID 5000](special-numbers.md#id-5000--server-voice-source-not-announcement-tg) para la tabla completa. +Ver [Números especiales — ID de voz del servidor](special-numbers.md#server-voice-source-id-configurable) para la tabla completa. ## Funciones @@ -46,6 +63,8 @@ Configura **`TTS_VOCODER_CMD`** o **`TTS_AMBESERVER_HOST`** / **`TTS_AMBESERVER_ Al programar anuncios, el servidor puede **esperar** si los slots objetivo están ocupados, **descartar** objetivos si aparece un QSO en vivo a mitad, y solo marcar estado de anuncio **horario** tras una lista de objetivos con éxito — evita pisar tráfico en vivo. +Los anuncios programados y TTS inyectan cada trama como **PTT sintético de hotspot** en `PROXY.TARGET_SYSTEM` (el mismo MASTER que usa el proxy integrado). El enrutamiento reparte a hotspots puenteados y piernas OPENBRIDGE; los peers locales en ese MASTER siguen recibiendo tramas vía `send_system`. + ## Cola de emisión Las **emisiones** en paralelo en **TG distintas** pueden ejecutarse a la vez; las emisiones en el **mismo TG** se **serializan** para que un anuncio termine antes que otro en ese TG. diff --git a/src/adn_server/application/ident_use_cases.py b/src/adn_server/application/ident_use_cases.py index 41d2680..043e404 100644 --- a/src/adn_server/application/ident_use_cases.py +++ b/src/adn_server/application/ident_use_cases.py @@ -31,6 +31,7 @@ import time from typing import Any, Callable from ..domain import HBPF_SLT_VTERM, bytes_3, int_id +from .server_voice import server_voice_rf_src_bytes logger = logging.getLogger(__name__) @@ -119,7 +120,7 @@ class IdentUseCases: continue _all_call = bytes_3(16777215) - _source_id = bytes_3(5000) + _source_id = server_voice_rf_src_bytes(self._config) _dst_id = b"" override_tg = sys_cfg.get("OVERRIDE_IDENT_TG") if override_tg is not None and int(override_tg) > 0 and int(override_tg) < 16777215: diff --git a/src/adn_server/application/routing/announcement_ptt_inject.py b/src/adn_server/application/routing/announcement_ptt_inject.py new file mode 100644 index 0000000..adfa7ed --- /dev/null +++ b/src/adn_server/application/routing/announcement_ptt_inject.py @@ -0,0 +1,99 @@ +# ADN DMR Peer Server - announcement synthetic PTT ingress +# +# 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 +############################################################################### +# +# Derived from ADN DMR Server / FreeDMR / HBlink. Original license: +############################################################################### +# Copyright (C) 2026 Joaquin Madrid Belando, EA5GVK +# Copyright (C) 2020 Simon Adlem, G7RZU +# Copyright (C) 2016-2019 Cortney T. Buffington, N0MJS +# +# 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 +############################################################################### + +"""Inject scheduled announcements as a synthetic hotspot PTT on the proxy MASTER.""" + +from __future__ import annotations + +from typing import Any + +from ..proxy.deployment import proxy_target_system +from .helpers import parse_dmrd_burst_fields + + +def announcement_ptt_system(config: dict[str, Any]) -> str | None: + """MASTER used for synthetic PTT ingress (same as integrated hotspot proxy target).""" + target = proxy_target_system(config) + systems_cfg = config.get("SYSTEMS", {}) + if target: + sys_cfg = systems_cfg.get(target, {}) + if sys_cfg.get("MODE") == "MASTER" and sys_cfg.get("ENABLED", True): + return target + for name, sys_cfg in systems_cfg.items(): + if sys_cfg.get("MODE") != "MASTER": + continue + if not sys_cfg.get("ENABLED", True): + continue + if sys_cfg.get("PEERS"): + return name + return None + + +def inject_announcement_ptt( + routing: Any, + master_system: str, + pkt: bytes, + *, + pkt_time: float, + server_id: bytes, +) -> bool | None: + """Feed one DMRD frame through ``dmrd_received`` (HBP path, same as proxy inject).""" + burst = parse_dmrd_burst_fields(pkt) + if burst is None: + return False + slot, frame_type, dtype_vseq, stream_id, dst_id, call_type = burst + seq = pkt[4] if len(pkt) > 4 else 0 + rf_src = pkt[5:8] if len(pkt) > 7 else b"\x00\x00\x00" + peer_id = pkt[11:15] if len(pkt) >= 15 else server_id + return routing.dmrd_received( + master_system, + peer_id, + rf_src, + dst_id, + seq, + slot, + call_type, + frame_type, + dtype_vseq, + stream_id, + pkt, + ingress_pkt_time=pkt_time, + ) diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index ab143c9..8fb6c8c 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -48,6 +48,7 @@ from typing import Any from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, bytes_3, bytes_4, int_id from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO +from ..server_voice import DEFAULT_SERVER_VOICE_ID PeerVoiceSlotRow = dict[str, Any] PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow] @@ -461,17 +462,30 @@ def inject_only_defer_obp_hbp_slot_contention( return source_is_hbp and connected_count > 1 -_SERVER_VOICE_RF_SRC = 5000 +_SERVER_VOICE_RF_SRC_LEGACY = 5000 -def master_slot_holds_server_broadcast(slot_st: dict[str, Any], pkt_time: float) -> bool: +def master_slot_holds_server_broadcast( + slot_st: dict[str, Any], + pkt_time: float, + *, + server_voice_rf_src: int | None = None, + server_voice_rf_srcs: frozenset[int] | None = None, +) -> bool: """True when server scheduled/TTS voice holds the flat MASTER slot TX row. - Announcements stamp TX_TYPE=VHEAD, TX_RFS=5000, and refresh TX_TIME each frame. + Announcements stamp TX_TYPE=VHEAD, TX_RFS=server voice ID, and refresh TX_TIME each frame. Re-apply global slot contention for OBP→MASTER when held so inject-only defer does not interleave mesh voice (legacy ``bridge_master`` TX_TGID/TX_TIME rules). """ - if int_id(slot_st.get("TX_RFS", b"\x00\x00\x00")) != _SERVER_VOICE_RF_SRC: + tx_rfs = int_id(slot_st.get("TX_RFS", b"\x00\x00\x00")) + if server_voice_rf_srcs is not None: + if tx_rfs not in server_voice_rf_srcs: + return False + elif server_voice_rf_src is not None: + if tx_rfs != server_voice_rf_src: + return False + elif tx_rfs != DEFAULT_SERVER_VOICE_ID and tx_rfs != _SERVER_VOICE_RF_SRC_LEGACY: return False tx_type = slot_st.get("TX_TYPE") if tx_type is None or tx_type == HBPF_SLT_VTERM: @@ -1052,8 +1066,13 @@ def is_on_demand_service_dst(dst_id: int) -> bool: return 9991 <= int(dst_id) <= 9999 -def is_server_originated_voice(packet: bytes) -> bool: - """True for server group playback (src 5000 -> TG 9 TS2), e.g. on-demand / disconnected.""" +def is_server_originated_voice( + packet: bytes, + *, + server_voice_rf_src: int | None = None, + server_voice_rf_srcs: frozenset[int] | None = None, +) -> bool: + """True for server group playback (server voice ID -> TG 9 TS2), e.g. on-demand / disconnected.""" if len(packet) < 11: return False burst = parse_dmrd_burst_fields(packet) @@ -1062,7 +1081,16 @@ def is_server_originated_voice(packet: bytes) -> bool: slot, _, _, _, dst_id, call_type = burst if call_type not in ("group", "vcsbk"): return False - return int_id(packet[5:8]) == 5000 and int_id(dst_id) == 9 and slot == 2 + src = int_id(packet[5:8]) + if server_voice_rf_srcs is not None: + if src not in server_voice_rf_srcs: + return False + elif server_voice_rf_src is not None: + if src != server_voice_rf_src and src != _SERVER_VOICE_RF_SRC_LEGACY: + return False + elif src != DEFAULT_SERVER_VOICE_ID and src != _SERVER_VOICE_RF_SRC_LEGACY: + return False + return int_id(dst_id) == 9 and slot == 2 def parse_dmrd_route_fields(packet: bytes) -> tuple[int, int, str] | None: diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 2d9c5ea..f7257b9 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -64,6 +64,7 @@ from .routing.store_authority_mixin import StoreAuthorityMixin from .routing.subscription_table import SubscriptionTableMixin from .routing.timers import RoutingTimerMixin from .routing.voice_subscription import VoiceSubscriptionMixin +from .server_voice import all_server_voice_ids from .talker_alias_use_cases import TalkerAliasUseCases logger = logging.getLogger(__name__) @@ -692,7 +693,11 @@ class RoutingUseCases( _ts_st["TX_TYPE"] = HBPF_SLT_VTERM _apply_master_slot_contention = ( not _defer_slot_contention - or master_slot_holds_server_broadcast(_ts_st, pkt_time) + or master_slot_holds_server_broadcast( + _ts_st, + pkt_time, + server_voice_rf_srcs=all_server_voice_ids(self._config), + ) ) if ( _apply_master_slot_contention diff --git a/src/adn_server/application/server_voice.py b/src/adn_server/application/server_voice.py new file mode 100644 index 0000000..8d9420d --- /dev/null +++ b/src/adn_server/application/server_voice.py @@ -0,0 +1,114 @@ +# ADN DMR Peer Server - server-originated voice identity +# +# 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 +############################################################################### + +"""RF source ID for scheduled announcements, TTS, and server playback. + +``VOICE.DMR_ID`` is optional. When absent (legacy ``adn-voice.yaml``), the +default is **1000001**. Per-item ``DMR_ID`` on announcement rows is also +optional and inherits the global default. Callsign/display name comes from the +subscriber alias DB (``_SUB_IDS`` / users file) for that DMR ID — not from voice +config. Invalid DMR_ID values are ignored and fall back to the default. +""" + +from __future__ import annotations + +from typing import Any + +from ..domain import bytes_3 + +DEFAULT_SERVER_VOICE_ID = 1000001 +LEGACY_SERVER_VOICE_ID = 5000 + + +def _parse_dmr_id(raw: Any) -> int | None: + if raw is None or raw == "": + return None + try: + return int(raw) + except (TypeError, ValueError): + return None + + +def server_voice_dmr_id(config: dict[str, Any] | None) -> int: + """Global default RF source (``VOICE.DMR_ID``), default 1000001.""" + if not config: + return DEFAULT_SERVER_VOICE_ID + voice = config.get("VOICE") + if not isinstance(voice, dict): + return DEFAULT_SERVER_VOICE_ID + for key in ("DMR_ID", "SRC_ID", "ID"): + parsed = _parse_dmr_id(voice.get(key)) + if parsed is not None: + return parsed + return DEFAULT_SERVER_VOICE_ID + + +def server_voice_id(config: dict[str, Any] | None) -> int: + """Alias for :func:`server_voice_dmr_id` (global default RF source).""" + return server_voice_dmr_id(config) + + +def server_voice_src_id(config: dict[str, Any] | None) -> int: + """Deprecated alias for :func:`server_voice_dmr_id`.""" + return server_voice_dmr_id(config) + + +def announcement_item_dmr_id( + item: dict[str, Any] | None, + config: dict[str, Any] | None, +) -> int: + """Per-item ``DMR_ID`` override, else global ``VOICE.DMR_ID``.""" + if isinstance(item, dict): + parsed = _parse_dmr_id(item.get("DMR_ID")) + if parsed is not None: + return parsed + return server_voice_dmr_id(config) + + +def announcement_item_source_bytes( + item: dict[str, Any] | None, + config: dict[str, Any] | None, +) -> bytes: + return bytes_3(announcement_item_dmr_id(item, config)) + + +def server_voice_rf_src_bytes(config: dict[str, Any] | None) -> bytes: + return bytes_3(server_voice_dmr_id(config)) + + +def all_server_voice_ids(config: dict[str, Any] | None) -> frozenset[int]: + """All RF source IDs used by server voice (global, per-item, legacy).""" + ids = {server_voice_dmr_id(config), LEGACY_SERVER_VOICE_ID} + if not config: + return frozenset(ids) + voice = config.get("VOICE") + if not isinstance(voice, dict): + return frozenset(ids) + for key in ("ANNOUNCEMENTS", "TTS_ANNOUNCEMENTS"): + for entry in voice.get(key) or []: + if isinstance(entry, dict): + parsed = _parse_dmr_id(entry.get("DMR_ID")) + if parsed is not None: + ids.add(parsed) + return frozenset(ids) + + +def is_server_voice_rf_src(rf_src: int, config: dict[str, Any] | None) -> bool: + return int(rf_src) in all_server_voice_ids(config) diff --git a/src/adn_server/application/voice_use_cases.py b/src/adn_server/application/voice_use_cases.py index f6495d7..2408eb5 100644 --- a/src/adn_server/application/voice_use_cases.py +++ b/src/adn_server/application/voice_use_cases.py @@ -31,8 +31,13 @@ import time from datetime import datetime from typing import Any, Callable -from ..domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4 +from ..domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id from .ports import VoiceProvider +from .server_voice import ( + announcement_item_source_bytes, + server_voice_id, + server_voice_rf_src_bytes, +) logger = logging.getLogger(__name__) @@ -55,6 +60,9 @@ class VoiceUseCases: call_later: Callable[..., Any] | None = None, start_looping_call: Callable[[Callable[[], None], float, bool], Any] | None = None, defer_to_thread: Callable[..., Any] | None = None, + inject_announcement_ptt: Callable[[bytes, float], bool | None] | None = None, + send_routing_event: Callable[[str], None] | None = None, + announcement_ptt_system: str | None = None, ) -> None: self._voice = voice_provider self._config = config @@ -65,6 +73,10 @@ class VoiceUseCases: self._call_later = call_later self._start_looping_call = start_looping_call self._defer_to_thread = defer_to_thread + self._inject_announcement_ptt = inject_announcement_ptt + self._send_routing_event = send_routing_event + self._announcement_ptt_system = announcement_ptt_system + self._voice_report_state: dict[str, Any] | None = None self._ann_tasks: dict[int, Any] = {} self._tts_tasks: dict[int, Any] = {} self._announcement_running: dict[int, bool] = {} @@ -74,6 +86,9 @@ class VoiceUseCases: self._broadcast_queue: list[dict[str, Any]] = [] self._broadcast_active_tgs: set[str] = set() + def _server_source_id(self) -> bytes: + return server_voice_rf_src_bytes(self._config) + def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]: """Load AMBE words for given languages (legacy readAMBE.readfiles).""" return self._voice.get_ambe_words(languages, audio_path) @@ -82,14 +97,114 @@ class VoiceUseCases: """Generate HBP voice packets for phrase (legacy mk_voice.pkt_gen).""" return self._voice.pkt_gen(rf_src, dst_id, peer, slot, phrase) + def _global_server_id_bytes(self) -> bytes: + server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") + return bytes_4(int_id(server_id)) + + def _active_bridge_slots_for_tg(self, tg: int, system: str) -> set[int]: + bridges = self._routing_table_for_report() if self._routing_table_for_report else {} + entries = bridges.get(str(tg), []) + if not isinstance(entries, list): + return set() + out: set[int] = set() + for be in entries: + if not isinstance(be, dict): + continue + if be.get("SYSTEM") != system or not be.get("ACTIVE"): + continue + ts = be.get("TS") + if ts is not None: + out.add(int(ts)) + return out + + def _inject_ptt_slot_busy( + self, + slot: dict[str, Any], + tg: int, + sys_cfg: dict[str, Any], + wire_ts: int, + ) -> bool: + """True when MASTER slot cannot accept synthetic PTT for ``tg``.""" + from .routing.helpers import master_dynamic_tg_slots, master_slot_holds_server_broadcast + from .server_voice import all_server_voice_ids + + if wire_ts in master_dynamic_tg_slots(sys_cfg, int(tg)): + return False + ptt_system = self._announcement_ptt_system or "" + if wire_ts in self._active_bridge_slots_for_tg(tg, ptt_system): + return False + if slot.get("RX_TYPE") == HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM: + return False + if master_slot_holds_server_broadcast( + slot, + time.time(), + server_voice_rf_srcs=all_server_voice_ids(self._config), + ): + return False + return True + + def _inject_ptt_slot_order(self, tg: int, sys_cfg: dict[str, Any]) -> list[int]: + """Prefer the RF slot where ``tg`` is already a dynamic UA session.""" + from .routing.helpers import master_dynamic_tg_slots + + dynamic = sorted(master_dynamic_tg_slots(sys_cfg, int(tg)), reverse=True) + preferred = list(dynamic) + if self._announcement_ptt_system: + for ts in sorted( + self._active_bridge_slots_for_tg(tg, self._announcement_ptt_system), + reverse=True, + ): + if ts not in preferred: + preferred.append(ts) + ordered: list[int] = [] + for ts in [*preferred, 2, 1]: + if ts not in ordered: + ordered.append(ts) + return ordered + + def _build_inject_ptt_targets( + self, label: str, tg: int, + ) -> tuple[list[dict[str, Any]], int]: + """Synthetic PTT on proxy MASTER: routing creates the UA bridge on inject.""" + targets: list[dict[str, Any]] = [] + busy_count = 0 + ptt_system = self._announcement_ptt_system + if not ptt_system or not self._inject_announcement_ptt: + return targets, busy_count + protocols = self._get_protocols() if self._get_protocols else {} + systems_cfg = self._config.get("SYSTEMS", {}) + if ptt_system not in protocols or ptt_system not in systems_cfg: + return targets, busy_count + if systems_cfg[ptt_system].get("MODE") != "MASTER": + return targets, busy_count + sys_obj = protocols[ptt_system] + status = getattr(sys_obj, "STATUS", None) + if not status: + return targets, busy_count + sys_cfg = systems_cfg[ptt_system] + for ts in self._inject_ptt_slot_order(tg, sys_cfg): + slot_index = 2 if ts == 2 else 1 + slot = status.get(slot_index) + if not slot: + continue + if self._inject_ptt_slot_busy(slot, tg, sys_cfg, ts): + logger.debug("(%s) System %s TS%s busy (QSO active), skipping", label, ptt_system, ts) + busy_count += 1 + continue + targets.append({"sys_obj": sys_obj, "name": ptt_system, "slot": slot, "ts": ts}) + break + return targets, busy_count + def _build_announcement_targets( self, tg_int: int, tg_str: str, label: str ) -> tuple[list[dict[str, Any]], int]: - """MASTER systems with active bridge for tg_int and idle slot (legacy target list). + """MASTER targets for announcement/TTS broadcast. - Returns (targets, busy_count). busy_count increments when a candidate slot is skipped - because a QSO is active (RX/TX not VTERM), for anti-collision retry scheduling. + Inject path: synthetic PTT on the TG's dynamic UA slot when applicable. + Legacy path: MASTER systems with an ACTIVE bridge row for the TG. """ + if self._inject_announcement_ptt: + return self._build_inject_ptt_targets(label, tg_int) targets: list[dict[str, Any]] = [] busy_count = 0 protocols = self._get_protocols() if self._get_protocols else {} @@ -135,6 +250,18 @@ class VoiceUseCases: targets.append({"sys_obj": sys_obj, "name": sys_name, "slot": slot, "ts": ts}) return targets, busy_count + def _announcement_packet_peer( + self, + tg: int, + targets: list[dict[str, Any]], + server_id: bytes, + ) -> bytes: + """DMRD peer field: always GLOBAL SERVER_ID on inject (never a hotspot radio id).""" + del tg, targets + if self._inject_announcement_ptt: + return self._global_server_id_bytes() + return bytes_4(int_id(server_id)) + def _send_filtered_by_tg( self, sys_obj: Any, pkt: bytes, tg: int, ts: int, bridges: dict[str, list[dict[str, Any]]] ) -> int: @@ -146,9 +273,150 @@ class VoiceUseCases: return -1 return 0 - def _mark_slots_busy(self, targets: list[dict[str, Any]]) -> None: + def _emit_announcement_voice_event( + self, + action: str, + trx: str, + system: str, + stream_id: bytes, + slot: int, + tg: int, + rf_src: int, + duration: float | None = None, + ) -> None: + if not self._send_routing_event: + return + parts = [ + "GROUP VOICE", + action, + trx, + system, + str(int_id(stream_id)), + str(rf_src), + str(rf_src), + str(slot), + str(tg), + ] + if duration is not None: + parts.append(f"{duration:.2f}") + self._send_routing_event(",".join(parts)) + + def _maybe_begin_legacy_voice_report( + self, + targets: list[dict[str, Any]], + pkts_by_ts: dict[int, list[bytes]], + tg: int, + ) -> None: + # Inject: routing emits START,RX (SERVER_ID) + OBP TX; same-MASTER downlink has no + # bridge leg, so emit START,TX on the proxy MASTER for monitor fan-out (SYSTEM-N). + self._begin_legacy_voice_report(targets, pkts_by_ts, tg) + + def _maybe_end_legacy_voice_report(self) -> None: + self._end_legacy_voice_report() + + def _begin_legacy_voice_report( + self, + targets: list[dict[str, Any]], + pkts_by_ts: dict[int, list[bytes]], + tg: int, + ) -> None: + if not self._send_routing_event or not targets: + return + wire_ts = targets[0]["ts"] + pkts = pkts_by_ts.get(wire_ts) or [] + if not pkts: + return + stream_id = pkts[0][16:20] + rf_src = int_id(pkts[0][5:8]) + systems: list[tuple[str, int]] = [] + seen: set[tuple[str, int]] = set() + for t in targets: + key = (t["name"], t["ts"]) + if key in seen: + continue + seen.add(key) + systems.append(key) + self._emit_announcement_voice_event( + "START", "TX", t["name"], stream_id, t["ts"], tg, rf_src + ) + self._voice_report_state = { + "stream_id": stream_id, + "rf_src": pkts[0][5:8], + "start": time.time(), + "tg": tg, + "systems": systems, + } + + def _end_legacy_voice_report(self) -> None: + state = self._voice_report_state + self._voice_report_state = None + if not state or not self._send_routing_event: + return + duration = time.time() - float(state["start"]) + stream_id = state["stream_id"] + tg = int(state["tg"]) + for name, ts in state["systems"]: + self._emit_announcement_voice_event( + "END", "TX", name, stream_id, ts, tg, int_id(state["rf_src"]), duration + ) + + def _send_announcement_packets( + self, + targets: list[dict[str, Any]], + pkts_by_ts: dict[int, list[bytes]], + pkt_idx: int, + source_id: bytes, + dst_id: bytes, + tg: int, + label: str, + ) -> None: + """One frame: synthetic PTT via routing, or legacy per-MASTER send_system.""" + if self._inject_announcement_ptt: + wire_ts = targets[0]["ts"] if targets else 2 + pkt = pkts_by_ts[wire_ts][pkt_idx] + if self._inject_announcement_ptt(pkt, time.time()) is False: + logger.warning("(%s) Routing rejected announcement frame %s", label, pkt_idx) + return + bridges = self._routing_table_for_report() if self._routing_table_for_report else {} + now = time.time() + for t in targets: + try: + sys_obj = t["sys_obj"] + slot = t["slot"] + t_ts = t["ts"] + pkt = pkts_by_ts[t_ts][pkt_idx] + stream_id = pkt[16:20] + if stream_id not in sys_obj.STATUS: + sys_obj.STATUS[stream_id] = { + "START": now, + "CONTENTION": False, + "RFS": source_id, + "TGID": dst_id, + "LAST": now, + } + slot["TX_TGID"] = dst_id + slot["TX_RFS"] = source_id + else: + sys_obj.STATUS[stream_id]["LAST"] = now + slot["TX_TIME"] = now + self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges) + except Exception as e: + logger.error( + "(%s) Error sending packet %s to %s/TS%s: %s", + label, + pkt_idx, + t.get("name"), + t.get("ts"), + e, + ) + + def _mark_slots_busy( + self, + targets: list[dict[str, Any]], + source_id: bytes | None = None, + ) -> None: """Mark target slots busy (TX_TYPE=VHEAD) to prevent TS conflict.""" - server_rfs = bytes_3(5000) + server_rfs = source_id if source_id is not None else self._server_source_id() now = time.time() for t in targets: try: @@ -181,10 +449,10 @@ class VoiceUseCases: 'source_id': source_id, 'dst_id': dst_id, 'tg': tg, 'num': num, 'label': label, }) _pos = len(self._broadcast_queue) - logger.info('(%s) Enqueued broadcast for same TG %s (position %s in queue)', label, tg, _pos) + logger.info('(%s) Same TG %s still on air; deferring next playback (pending %s)', label, tg, _pos) else: self._broadcast_active_tgs.add(_tg_key) - self._mark_slots_busy(targets) + self._mark_slots_busy(targets, source_id) logger.info('(%s) Starting broadcast immediately for TG %s (active TGs: %s)', label, tg, len(self._broadcast_active_tgs)) if self._call_later: if _type == 'ann': @@ -207,8 +475,8 @@ class VoiceUseCases: _label = _next['label'] _tg_key = str(_next['tg']) self._broadcast_active_tgs.add(_tg_key) - self._mark_slots_busy(_next['targets']) - logger.info('(%s) Starting broadcast from queue for TG %s (%s remaining, active TGs: %s)', _label, _next['tg'], len(self._broadcast_queue), len(self._broadcast_active_tgs)) + self._mark_slots_busy(_next['targets'], _next['source_id']) + logger.info('(%s) Starting deferred same-TG playback for TG %s (%s pending, %s active TG(s))', _label, _next['tg'], len(self._broadcast_queue), len(self._broadcast_active_tgs)) if self._call_later: if _type == 'ann': self._call_later(0.5, self._announcement_send_broadcast, _next['targets'], _next['pkts_by_ts'], 0, _next['source_id'], _next['dst_id'], _next['tg'], _next['num'], _label, None) @@ -219,14 +487,20 @@ class VoiceUseCases: if tg is not None: self._broadcast_active_tgs.discard(str(tg)) if self._broadcast_queue: - logger.info('(QUEUE) Broadcast finished for TG %s, checking queue (%s queued, active TGs: %s)', tg, len(self._broadcast_queue), len(self._broadcast_active_tgs)) + logger.info( + '(BROADCAST) TG %s announcement/TTS playback finished; %s same-TG deferred, %s active TG(s)', + tg, len(self._broadcast_queue), len(self._broadcast_active_tgs), + ) if self._call_later: self._call_later(_BROADCAST_GAP, self._start_next_broadcast) else: if not self._broadcast_active_tgs: - logger.info('(QUEUE) All broadcasts finished, queue empty') + logger.info('(BROADCAST) All announcement/TTS playbacks finished') else: - logger.info('(QUEUE) Broadcast finished for TG %s, %s TGs still active', tg, len(self._broadcast_active_tgs)) + logger.info( + '(BROADCAST) TG %s announcement/TTS playback finished; %s other TG(s) still on air', + tg, len(self._broadcast_active_tgs), + ) def _announcement_send_broadcast( self, @@ -244,6 +518,7 @@ class VoiceUseCases: total = len(pkts_by_ts.get(1, [])) if pkt_idx >= total or not targets: self._mark_slots_free(targets) + self._maybe_end_legacy_voice_report() for t in targets: try: obj = t.get("sys_obj") @@ -283,6 +558,7 @@ class VoiceUseCases: targets.remove(t) if not targets: self._announcement_running[ann_idx] = False + self._maybe_end_legacy_voice_report() logger.info( "(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s", label, @@ -291,31 +567,10 @@ class VoiceUseCases: ) self._broadcast_finished(tg) return - bridges = self._routing_table_for_report() if self._routing_table_for_report else {} now = time.time() - for t in targets: - try: - sys_obj = t["sys_obj"] - slot = t["slot"] - t_ts = t["ts"] - pkt = pkts_by_ts[t_ts][pkt_idx] - stream_id = pkt[16:20] - if stream_id not in sys_obj.STATUS: - sys_obj.STATUS[stream_id] = { - "START": now, - "CONTENTION": False, - "RFS": source_id, - "TGID": dst_id, - "LAST": now, - } - slot["TX_TGID"] = dst_id - slot["TX_RFS"] = source_id - else: - sys_obj.STATUS[stream_id]["LAST"] = now - slot["TX_TIME"] = now - self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges) - except Exception as e: - logger.error("(%s) Error sending packet %s to %s/TS%s: %s", label, pkt_idx, t.get("name"), t.get("ts"), e) + if pkt_idx == 0: + self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg) + self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label) if next_time is None: next_time = now + _FRAME_INTERVAL else: @@ -369,7 +624,7 @@ class VoiceUseCases: if not _file or not _tg: return _dst_id = bytes_3(_tg) - _source_id = bytes_3(5000) + _source_id = announcement_item_source_bytes(item, self._config) server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") if not isinstance(server_id, bytes): server_id = bytes_3(int(server_id)) @@ -400,9 +655,10 @@ class VoiceUseCases: if mode == "hourly": self._announcement_last_hour[ann_idx] = datetime.now().hour _say_list = [_say] + pkt_peer = self._announcement_packet_peer(_tg, targets, server_id) pkts_by_ts = { - 1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _say_list)), - 2: list(self.pkt_gen(_source_id, _dst_id, server_id, 1, _say_list)), + 1: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 0, _say_list)), + 2: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 1, _say_list)), } ts1_count = sum(1 for t in targets if t["ts"] == 1) ts2_count = sum(1 for t in targets if t["ts"] == 2) @@ -475,7 +731,12 @@ class VoiceUseCases: return logger.info("(%s) Playing TTS file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang) _dst_id = bytes_3(_tg) - _source_id = bytes_3(5000) + tts_list = self._config.get("VOICE", {}).get("TTS_ANNOUNCEMENTS") or [] + tts_item = tts_list[tts_idx] if 0 <= tts_idx < len(tts_list) else None + _source_id = announcement_item_source_bytes( + tts_item if isinstance(tts_item, dict) else None, + self._config, + ) server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") if not isinstance(server_id, bytes): server_id = bytes_3(int(server_id)) @@ -515,9 +776,10 @@ class VoiceUseCases: if mode == "hourly": self._tts_last_hour[tts_idx] = datetime.now().hour _say_list = [_say] + pkt_peer = self._announcement_packet_peer(_tg, targets, server_id) pkts_by_ts = { - 1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _say_list)), - 2: list(self.pkt_gen(_source_id, _dst_id, server_id, 1, _say_list)), + 1: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 0, _say_list)), + 2: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 1, _say_list)), } logger.info("(%s) Broadcasting %s packets to %s targets", label, len(pkts_by_ts[1]), len(targets)) self._enqueue_broadcast('tts', targets, pkts_by_ts, _source_id, _dst_id, _tg, tts_idx, label) @@ -547,6 +809,7 @@ class VoiceUseCases: total = len(pkts_by_ts.get(1, [])) if pkt_idx >= total or not targets: self._mark_slots_free(targets) + self._maybe_end_legacy_voice_report() for t in targets: try: obj = t.get("sys_obj") @@ -586,6 +849,7 @@ class VoiceUseCases: targets.remove(t) if not targets: self._tts_running[tts_idx] = False + self._maybe_end_legacy_voice_report() logger.info( "(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s", label, @@ -594,31 +858,10 @@ class VoiceUseCases: ) self._broadcast_finished(tg) return - bridges = self._routing_table_for_report() if self._routing_table_for_report else {} now = time.time() - for t in targets: - try: - sys_obj = t["sys_obj"] - slot = t["slot"] - t_ts = t["ts"] - pkt = pkts_by_ts[t_ts][pkt_idx] - stream_id = pkt[16:20] - if stream_id not in sys_obj.STATUS: - sys_obj.STATUS[stream_id] = { - "START": now, - "CONTENTION": False, - "RFS": source_id, - "TGID": dst_id, - "LAST": now, - } - slot["TX_TGID"] = dst_id - slot["TX_RFS"] = source_id - else: - sys_obj.STATUS[stream_id]["LAST"] = now - slot["TX_TIME"] = now - self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges) - except Exception as e: - logger.error("(%s) Error sending packet %s to %s: %s", label, pkt_idx, t.get("name"), e) + if pkt_idx == 0: + self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg) + self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label) if next_time is None: next_time = now + _FRAME_INTERVAL else: @@ -659,12 +902,12 @@ class VoiceUseCases: logger.info("(%s) Playing on-demand AMBE file: %s (ID: %s)", system, file_number, file_number) time.sleep(1) _say = [pairs] - speech = self.pkt_gen(bytes_3(5000), bytes_3(9), bytes_4(9), 1, _say) + _source_id = self._server_source_id() + speech = self.pkt_gen(_source_id, bytes_3(9), bytes_4(9), 1, _say) time.sleep(1) _slot = protocol.STATUS.get(2) if not _slot: return - _source_id = bytes_3(5000) _dst_id = bytes_3(9) _next_time = time.time() _pkt_count = 0 @@ -708,7 +951,8 @@ class VoiceUseCases: else: _say.append(words.get("notlinked") or silence) _say.append(silence) - speech = self.pkt_gen(bytes_3(5000), bytes_3(9), bytes_4(9), 1, _say) + _source_id = self._server_source_id() + speech = self.pkt_gen(_source_id, bytes_3(9), bytes_4(9), 1, _say) time.sleep(1) _slot = protocol.STATUS.get(2) if not _slot: @@ -720,7 +964,7 @@ class VoiceUseCases: _delay = _next_time - time.time() if _delay > 0.001: time.sleep(_delay) - self._call_from_reactor(protocol.send_voice_packet, pkt, bytes_3(5000), bytes_3(9), _slot) + self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, bytes_3(9), _slot) logger.debug("(%s) disconnected voice thread end", system) def apply_voice_config(self) -> None: diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index ace687b..8143464 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -53,6 +53,7 @@ from adn_server.application.runtime_context import ( ) from adn_server.application.subscription.echo_seed import seed_echo_routing_table from adn_server.application.subscription.store_sync import replace_store_from_routing_table +from adn_server.domain import bytes_4 from adn_server.domain.dmr.bptc import encode_emblc from adn_server.infrastructure.acl_router import InMemoryAclRouter from adn_server.infrastructure.config_normalizer import ( @@ -374,6 +375,35 @@ def run_peer_server( ) routing_table_for_report = routing_use_cases.routing_table_for_report voice_use_cases._routing_table_for_report = routing_table_for_report + from adn_server.application.routing.announcement_ptt_inject import ( + announcement_ptt_system, + inject_announcement_ptt, + ) + + _ptt_system = announcement_ptt_system(config) + + def _inject_announcement_ptt(pkt: bytes, pkt_time: float) -> bool | None: + if not _ptt_system: + return False + server_id = config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00") + if not isinstance(server_id, bytes): + server_id = bytes_4(int(server_id or 0) & 0xFFFFFFFF) + accepted = inject_announcement_ptt( + routing_use_cases, + _ptt_system, + pkt, + pkt_time=pkt_time, + server_id=server_id, + ) + proto = protocols.get(_ptt_system) + send_system = getattr(proto, "send_system", None) if proto is not None else None + if callable(send_system): + send_system(pkt) + return accepted + + voice_use_cases._inject_announcement_ptt = _inject_announcement_ptt + voice_use_cases._announcement_ptt_system = _ptt_system + voice_use_cases._send_routing_event = reporting_use_cases.send_routing_event report_factory.set_routing_table(routing_table_for_report()) report_factory.set_systems(config.get("SYSTEMS", {})) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 26b1d3b..3350c4a 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -39,6 +39,7 @@ from twisted.internet import reactor, task from twisted.internet.protocol import DatagramProtocol from ...application.proxy.deployment import is_proxy_inject_only +from ...application.server_voice import all_server_voice_ids from ...application.routing.downlink import ( DownlinkContext, iter_downlink_voice_slots, @@ -398,7 +399,9 @@ class HBPProtocol(DatagramProtocol): connected = self._cached_connected_peer_count() if connected <= 0: return () - if is_server_originated_voice(packet): + if is_server_originated_voice( + packet, server_voice_rf_srcs=all_server_voice_ids(self._CONFIG) + ): slot_st = self.STATUS.get(2, {}) rx_peer = slot_st.get("RX_PEER", b"") if rx_peer and rx_peer in self._peers: @@ -575,8 +578,10 @@ class HBPProtocol(DatagramProtocol): slot, tgid, call_type = parsed if call_type not in ("group", "vcsbk"): return True - # Server playback (5000 -> TG 9): not in OPTIONS; deliver to requesting RX peer on TS2. - if is_server_originated_voice(packet): + # Server playback (server voice ID -> TG 9): not in OPTIONS; deliver to requesting RX peer on TS2. + if is_server_originated_voice( + packet, server_voice_rf_srcs=all_server_voice_ids(self._CONFIG) + ): slot_st = self.STATUS.get(2, {}) rx_peer = slot_st.get("RX_PEER", b"") if rx_peer and rx_peer != b"\x00\x00\x00\x00": diff --git a/src/adn_server/infrastructure/voice/tts_engine.py b/src/adn_server/infrastructure/voice/tts_engine.py index a0322fe..c9cd09f 100644 --- a/src/adn_server/infrastructure/voice/tts_engine.py +++ b/src/adn_server/infrastructure/voice/tts_engine.py @@ -34,11 +34,21 @@ import os import socket import struct import subprocess +import threading +import time import wave from typing import Any logger = logging.getLogger(__name__) +# One TTS conversion at a time (shared AMBEServer / vocoder cannot handle parallel sessions). +_tts_conversion_lock = threading.Lock() + +_AMBESERVER_FRAME_TIMEOUT_S = 5.0 +_AMBESERVER_ENCODE_TIMEOUT_S = 120.0 +_AMBESERVER_MAX_CONSECUTIVE_ERRORS = 10 +_AMBESERVER_PROGRESS_EVERY_FRAMES = 100 + _LANG_MAP = { "es_ES": "es", "en_GB": "en", "en_US": "en", "fr_FR": "fr", "de_DE": "de", "it_IT": "it", "pt_PT": "pt", "pt_BR": "pt", @@ -168,7 +178,7 @@ def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int) return False try: sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - sock.settimeout(5.0) + sock.settimeout(_AMBESERVER_FRAME_TIMEOUT_S) except Exception as e: logger.error("(TTS-AMBESERVER) Error creating UDP socket: %s", e) return False @@ -212,12 +222,33 @@ def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int) return False _total_frames = wf.getnframes() _sample_rate = wf.getframerate() - logger.info("(TTS-AMBESERVER) WAV: %d samples, %d Hz", _total_frames, _sample_rate) + _duration_s = _total_frames / _sample_rate if _sample_rate else 0.0 + _pcm_frames = (_total_frames + DV3K_SAMPLES_PER_FRAME - 1) // DV3K_SAMPLES_PER_FRAME + logger.info( + "(TTS-AMBESERVER) WAV: %d samples, %d Hz, duration: %.1fs (~%d AMBE frames)", + _total_frames, + _sample_rate, + _duration_s, + _pcm_frames, + ) _raw_frames = wf.readframes(_total_frames) wf.close() _samples = list(struct.unpack("<" + "h" * _total_frames, _raw_frames)) _ambe_frames: list[bytes] = [] + _frames_sent = 0 + _frames_error = 0 + _consecutive_errors = 0 + _encode_deadline = time.monotonic() + _AMBESERVER_ENCODE_TIMEOUT_S for i in range(0, len(_samples), DV3K_SAMPLES_PER_FRAME): + _frame_idx = i // DV3K_SAMPLES_PER_FRAME + if time.monotonic() >= _encode_deadline: + logger.error( + "(TTS-AMBESERVER) Encoding timeout after %ss (%d/%d frames)", + int(_AMBESERVER_ENCODE_TIMEOUT_S), + _frames_sent, + _pcm_frames, + ) + break _chunk = _samples[i : i + DV3K_SAMPLES_PER_FRAME] if len(_chunk) < DV3K_SAMPLES_PER_FRAME: _chunk = _chunk + [0] * (DV3K_SAMPLES_PER_FRAME - len(_chunk)) @@ -228,8 +259,38 @@ def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int) _ambe_data = _parse_ambe_response(data) if _ambe_data is not None: _ambe_frames.append(_ambe_data) - except (socket.timeout, Exception): - pass + _frames_sent += 1 + _consecutive_errors = 0 + else: + _frames_error += 1 + _consecutive_errors += 1 + logger.debug( + "(TTS-AMBESERVER) Frame %d: non-AMBE response (type: 0x%02X)", + _frame_idx, + data[3] if len(data) > 3 else 0, + ) + except socket.timeout: + _frames_error += 1 + _consecutive_errors += 1 + logger.warning("(TTS-AMBESERVER) Timeout on frame %d", _frame_idx) + except Exception as e: + _frames_error += 1 + _consecutive_errors += 1 + logger.error("(TTS-AMBESERVER) Error on frame %d: %s", _frame_idx, e) + if _consecutive_errors >= _AMBESERVER_MAX_CONSECUTIVE_ERRORS: + logger.error( + "(TTS-AMBESERVER) Aborting after %d consecutive frame errors (%d/%d sent)", + _consecutive_errors, + _frames_sent, + _pcm_frames, + ) + break + if _frame_idx > 0 and _frame_idx % _AMBESERVER_PROGRESS_EVERY_FRAMES == 0: + logger.info( + "(TTS-AMBESERVER) Progress: %d/%d frames encoded", + _frames_sent, + _pcm_frames, + ) sock.close() if not _ambe_frames: logger.error("(TTS-AMBESERVER) No AMBE frames received") @@ -241,7 +302,12 @@ def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int) except Exception as e: logger.error("(TTS-AMBESERVER) Error writing AMBE: %s", e) return False - logger.info("(TTS-AMBESERVER) Encoding completed: %s", ambe_path) + logger.info( + "(TTS-AMBESERVER) Encoding completed: %d frames (%d errors): %s", + _frames_sent, + _frames_error, + ambe_path, + ) return True @@ -254,24 +320,33 @@ def _cleanup(files: list[str]) -> None: logger.warning("(TTS) cleanup remove %s failed: %s", f, e) -def text_to_ambe( +def _cleanup(files: list[str]) -> None: + for f in files: + try: + if os.path.isfile(f): + os.remove(f) + except Exception as e: + logger.warning("(TTS) cleanup remove %s failed: %s", f, e) + + +def _ambe_cache_valid(txt_path: str, ambe_path: str) -> bool: + if not os.path.isfile(ambe_path): + return False + if not os.path.isfile(txt_path): + return True + return os.path.getmtime(ambe_path) > os.path.getmtime(txt_path) + + +def _text_to_ambe_uncached( txt_path: str, ambe_path: str, language: str, vocoder_cmd: str, - ambeserver_host: str = "", - ambeserver_port: int = 2460, - volume_db: int = 0, - speed: float = 1.0, + ambeserver_host: str, + ambeserver_port: int, + volume_db: int, + speed: float, ) -> bool: - """Convert .txt to .ambe (gTTS -> mp3 -> ffmpeg -> wav -> vocoder/AMBEServer).""" - if not os.path.isfile(txt_path): - logger.warning("(TTS) Text file not found: %s", txt_path) - return False - if os.path.isfile(ambe_path): - if os.path.getmtime(ambe_path) > os.path.getmtime(txt_path): - logger.info("(TTS) Using cached AMBE (newer than .txt): %s", ambe_path) - return True with open(txt_path, "r", encoding="utf-8") as f: text = f.read().strip() if not text: @@ -308,6 +383,39 @@ def text_to_ambe( return True +def text_to_ambe( + txt_path: str, + ambe_path: str, + language: str, + vocoder_cmd: str, + ambeserver_host: str = "", + ambeserver_port: int = 2460, + volume_db: int = 0, + speed: float = 1.0, +) -> bool: + """Convert .txt to .ambe (gTTS -> mp3 -> ffmpeg -> wav -> vocoder/AMBEServer).""" + if not os.path.isfile(txt_path): + logger.warning("(TTS) Text file not found: %s", txt_path) + return False + if _ambe_cache_valid(txt_path, ambe_path): + logger.info("(TTS) Using cached AMBE (newer than .txt): %s", ambe_path) + return True + with _tts_conversion_lock: + if _ambe_cache_valid(txt_path, ambe_path): + logger.info("(TTS) Using cached AMBE (newer than .txt): %s", ambe_path) + return True + return _text_to_ambe_uncached( + txt_path, + ambe_path, + language, + vocoder_cmd, + ambeserver_host, + ambeserver_port, + volume_db, + speed, + ) + + def ensure_tts_ambe(config: dict[str, Any], item: dict[str, Any], audio_path: str) -> str | None: """Ensure .ambe exists for TTS item; create from .txt if needed. Returns path or None.""" if not item.get("ENABLED", False): diff --git a/tests/application/test_master_server_broadcast_contention.py b/tests/application/test_master_server_broadcast_contention.py index cf1af84..ca1bec1 100644 --- a/tests/application/test_master_server_broadcast_contention.py +++ b/tests/application/test_master_server_broadcast_contention.py @@ -1,6 +1,24 @@ # ADN DMR Peer Server - server broadcast slot contention # # 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 +############################################################################### + +"""MASTER slot contention while server scheduled voice holds TX row.""" from __future__ import annotations @@ -8,6 +26,7 @@ from adn_server.application.routing.helpers import ( hbp_slot_blocks_group_voice, master_slot_holds_server_broadcast, ) +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4 @@ -15,7 +34,7 @@ def test_master_slot_holds_server_broadcast_detects_announcement_row() -> None: slot = { "TX_TYPE": HBPF_SLT_VHEAD, "TX_TIME": 100.0, - "TX_RFS": bytes_3(5000), + "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID), "TX_TGID": bytes_3(91), } assert master_slot_holds_server_broadcast(slot, 100.1) is True @@ -37,7 +56,7 @@ def test_obp_blocked_when_server_broadcast_holds_slot() -> None: "RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VHEAD, "TX_TIME": 200.0, - "TX_RFS": bytes_3(5000), + "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID), "TX_TGID": bytes_3(91), } blocked = hbp_slot_blocks_group_voice( diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py index ae060cd..f703f26 100644 --- a/tests/application/test_monitor_topology.py +++ b/tests/application/test_monitor_topology.py @@ -367,6 +367,23 @@ def test_remap_voice_event_for_single_peer() -> None: assert mapped.startswith("GROUP VOICE,START,TX,SYSTEM-7,") +def test_remap_announcement_tx_fans_out_to_ua_only_peer() -> None: + """Inject START,TX on SYSTEM reaches peers with dynamic UA but no static OPTIONS TG.""" + peer = bytes_4(730039210) + peers = {peer: _peer(options=b"TS2=730500;SINGLE=0;")} + config = _proxy_config(peers) + sys_cfg = config["SYSTEMS"]["SYSTEM"] + pk = bytes_4(730039210) + sys_cfg.setdefault("_PEER_UA_MULTI_TGS", {}).setdefault(pk, {})[2] = {730600} + peer_slots = {peer: 7} + raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,1000001,1000001,2,730600" + events = remap_inject_proxy_voice_events( + raw, config, config["SYSTEMS"], peer_slots, + ) + systems = {ev.split(",")[3] for ev in events} + assert systems == {"SYSTEM-7"} + + def test_non_proxy_systems_pass_through_unchanged() -> None: systems = { "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, diff --git a/tests/application/test_server_voice.py b/tests/application/test_server_voice.py new file mode 100644 index 0000000..052f933 --- /dev/null +++ b/tests/application/test_server_voice.py @@ -0,0 +1,105 @@ +# ADN DMR Peer Server - tests application server voice identity +# +# 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 +############################################################################### + +"""Server voice DMR_ID config (global default and per-announcement override).""" + +from __future__ import annotations + +from adn_server.application.server_voice import ( + DEFAULT_SERVER_VOICE_ID, + all_server_voice_ids, + announcement_item_dmr_id, + announcement_item_source_bytes, + server_voice_dmr_id, + server_voice_id, + server_voice_rf_src_bytes, +) +from adn_server.domain import bytes_3, int_id + + +def test_server_voice_defaults_when_voice_section_missing() -> None: + assert server_voice_dmr_id({}) == DEFAULT_SERVER_VOICE_ID + assert server_voice_id({}) == DEFAULT_SERVER_VOICE_ID + assert int_id(server_voice_rf_src_bytes({})) == DEFAULT_SERVER_VOICE_ID + + +def test_server_voice_reads_dmr_id_from_config() -> None: + config = {"VOICE": {"DMR_ID": 3109898}} + assert server_voice_dmr_id(config) == 3109898 + assert server_voice_rf_src_bytes(config) == bytes_3(3109898) + + +def test_server_voice_legacy_keys_still_work() -> None: + assert server_voice_dmr_id({"VOICE": {"SRC_ID": 2000002}}) == 2000002 + assert server_voice_dmr_id({"VOICE": {"ID": 3000003}}) == 3000003 + + +def test_server_voice_invalid_dmr_id_falls_back_to_default() -> None: + config = {"VOICE": {"DMR_ID": "not-a-number"}} + assert server_voice_dmr_id(config) == DEFAULT_SERVER_VOICE_ID + + +def test_announcement_item_dmr_id_override() -> None: + config = {"VOICE": {"DMR_ID": 1000001}} + item = {"DMR_ID": 3109898, "TG": 2} + assert announcement_item_dmr_id(item, config) == 3109898 + assert announcement_item_dmr_id({"TG": 2}, config) == 1000001 + + +def test_all_server_voice_ids_collects_global_and_per_item() -> None: + config = { + "VOICE": { + "DMR_ID": 1000001, + "ANNOUNCEMENTS": [{"DMR_ID": 2000002}], + "TTS_ANNOUNCEMENTS": [{"DMR_ID": 3000003}], + } + } + assert all_server_voice_ids(config) == frozenset({1000001, 2000002, 3000003, 5000}) + + +def test_legacy_voice_yaml_without_dmr_id_uses_defaults() -> None: + """Pre-migration adn-voice.yaml (no DMR_ID keys) keeps working.""" + legacy = { + "VOICE": { + "ANNOUNCEMENTS": [ + {"ENABLED": True, "TG": 91, "FILE": "welcome", "LANGUAGE": "en_GB", "MODE": "interval", "INTERVAL": 60}, + ], + "TTS_ANNOUNCEMENTS": [ + {"ENABLED": True, "TG": 92, "FILE": "texto1", "LANGUAGE": "es_ES", "MODE": "interval", "INTERVAL": 60}, + ], + } + } + ann = legacy["VOICE"]["ANNOUNCEMENTS"][0] + tts = legacy["VOICE"]["TTS_ANNOUNCEMENTS"][0] + assert server_voice_dmr_id(legacy) == DEFAULT_SERVER_VOICE_ID + assert announcement_item_dmr_id(ann, legacy) == DEFAULT_SERVER_VOICE_ID + assert announcement_item_dmr_id(tts, legacy) == DEFAULT_SERVER_VOICE_ID + assert int_id(announcement_item_source_bytes(ann, legacy)) == DEFAULT_SERVER_VOICE_ID + + +def test_legacy_voice_yaml_missing_voice_section_uses_defaults() -> None: + assert server_voice_dmr_id({}) == DEFAULT_SERVER_VOICE_ID + assert announcement_item_dmr_id(None, {}) == DEFAULT_SERVER_VOICE_ID + + +def test_invalid_item_dmr_id_falls_back_without_error() -> None: + config = {"VOICE": {"DMR_ID": "bad"}} + item = {"DMR_ID": [], "TG": 2} + assert announcement_item_dmr_id(item, config) == DEFAULT_SERVER_VOICE_ID diff --git a/tests/harness/voice_helpers.py b/tests/harness/voice_helpers.py index 50a637a..14382f9 100644 --- a/tests/harness/voice_helpers.py +++ b/tests/harness/voice_helpers.py @@ -24,6 +24,10 @@ from __future__ import annotations from typing import Any, Iterator +from adn_server.application.routing.announcement_ptt_inject import ( + announcement_ptt_system, + inject_announcement_ptt, +) from adn_server.application.voice_use_cases import VoiceUseCases from adn_server.domain import bytes_3, bytes_4 from tests.harness.deterministic import DeterministicScenario, FakeHbpProtocol, active_routing_table @@ -43,8 +47,9 @@ class FakeVoiceProvider: ) -> Iterator[bytes]: del phrase ts_bit = 0x80 if slot else 0 + stream_id = bytes_4(0xA0A0A0A0) for seq in range(3): - yield b"DMRD" + bytes([seq]) + rf_src[:3] + dst_id[:3] + peer[:4] + bytes([ts_bit | 0x10]) + bytes_4(0xA0A0A0A0 + seq) + b"\x00" * 33 + b"\x00\x00" + yield b"DMRD" + bytes([seq]) + rf_src[:3] + dst_id[:3] + peer[:4] + bytes([ts_bit | 0x10]) + stream_id + b"\x00" * 33 + b"\x00\x00" def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list: del audio_path, lang, file_number @@ -68,10 +73,11 @@ class FakeMasterForVoice(FakeHbpProtocol): def voice_master_scenario(tg: int = 91) -> tuple[DeterministicScenario, FakeMasterForVoice]: config = DeterministicScenario().config master = FakeMasterForVoice("MASTER-A") + config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} config["SYSTEMS"]["MASTER-A"]["PEERS"] = { "1001": {"CALLSIGN": "TEST", "IP": "127.0.0.1", "PORT": 62032}, } - bridges = active_routing_table(tg, (("MASTER-A", 2),)) + bridges = active_routing_table(tg, (("MASTER-A", 2), ("MASTER-B", 2))) scenario = DeterministicScenario(config=config, routing_table=bridges) scenario.protocols["MASTER-A"] = master master.STATUS[2] = { @@ -108,19 +114,39 @@ def make_voice_uc( audio_path: str = "/tmp/audio", ) -> VoiceUseCases: scheduled: list[tuple[float, tuple]] = [] + ptt_system = announcement_ptt_system(scenario.config) + server_id = scenario.config.get("GLOBAL", {}).get("SERVER_ID", bytes_4(9990)) + if not isinstance(server_id, bytes): + server_id = bytes_4(int(server_id or 0) & 0xFFFFFFFF) def call_later(delay, fn, *args): scheduled.append((delay, (fn, args))) return type("H", (), {"active": lambda self: True, "cancel": lambda self: None})() + def inject_announcement_ptt_cb(pkt: bytes, pkt_time: float) -> bool | None: + if not ptt_system: + return False + accepted = inject_announcement_ptt( + scenario.routing, + ptt_system, + pkt, + pkt_time=pkt_time, + server_id=server_id, + ) + if hasattr(master, "send_system"): + master.send_system(pkt) + return accepted + uc = VoiceUseCases( FakeVoiceProvider(), scenario.config, - get_protocols=lambda: {"MASTER-A": master}, + get_protocols=lambda: {"MASTER-A": master, **{k: v for k, v in scenario.protocols.items() if k != "MASTER-A"}}, routing_table_for_report=scenario.routing.routing_table_for_report, call_later=call_later, audio_path=audio_path, + inject_announcement_ptt=inject_announcement_ptt_cb, ) + uc._announcement_ptt_system = ptt_system uc._scheduled = scheduled return uc diff --git a/tests/infrastructure/test_hbp_repeat_options_filter.py b/tests/infrastructure/test_hbp_repeat_options_filter.py index ba21cf3..3d2d03b 100644 --- a/tests/infrastructure/test_hbp_repeat_options_filter.py +++ b/tests/infrastructure/test_hbp_repeat_options_filter.py @@ -26,6 +26,7 @@ import pytest from tests.harness.deterministic import DeterministicScenario, PacketSpec from tests.support.hbp_repeat_stack import build_hbp_repeat_stack +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID from adn_server.domain import bytes_3, bytes_4 from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM from adn_server.infrastructure.hbp_constants import DMRD @@ -188,7 +189,7 @@ def test_echo_tg_9990_reaches_caller_without_9990_in_options() -> None: def test_server_playback_tg9_reaches_requesting_peer_without_tg9_in_options() -> None: - """On-demand playback (5000 -> TG 9) must reach the peer that keyed 999x.""" + """On-demand playback (server voice ID -> TG 9) must reach the peer that keyed 999x.""" stack = _inject_proxy_stack() caller = bytes_4(730039101) other = bytes_4(730039102) @@ -197,8 +198,8 @@ def test_server_playback_tg9_reaches_requesting_peer_without_tg9_in_options() -> stack.register_peer(caller, addr_caller, options="TS2=730,730444;") stack.register_peer(other, addr_other, options="TS2=91;") spec = PacketSpec( - peer_id=5000, - rf_src=5000, + peer_id=DEFAULT_SERVER_VOICE_ID, + rf_src=DEFAULT_SERVER_VOICE_ID, dst_id=9, slot=2, stream_id=0xCAFEBABE, diff --git a/tests/routing/test_announcement_ptt_inject.py b/tests/routing/test_announcement_ptt_inject.py new file mode 100644 index 0000000..d81b037 --- /dev/null +++ b/tests/routing/test_announcement_ptt_inject.py @@ -0,0 +1,77 @@ +# ADN DMR Peer Server - announcement synthetic PTT inject +# +# 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 +############################################################################### + +"""Announcement synthetic PTT inject routing (proxy MASTER + OBP forward).""" + +from __future__ import annotations + +from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time +from tests.harness.scenarios import obp_bridge_scenario + +from adn_server.application.routing.announcement_ptt_inject import ( + announcement_ptt_system, + inject_announcement_ptt, +) +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID +from adn_server.domain import bytes_4 + + +def test_announcement_ptt_system_prefers_proxy_target() -> None: + config = DeterministicScenario().config + config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + assert announcement_ptt_system(config) == "MASTER-A" + + +def test_inject_forwards_to_obp_and_peer_master() -> None: + scenario = obp_bridge_scenario("OBP-CL", tg=91) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + server_id = scenario.config["GLOBAL"]["SERVER_ID"] + if isinstance(server_id, bytes): + peer_id = int.from_bytes(server_id[:4], "big") + else: + peer_id = int(server_id) + base = PacketSpec( + peer_id=peer_id, + rf_src=DEFAULT_SERVER_VOICE_ID, + dst_id=91, + slot=2, + stream_id=0xA0A0A0A0, + ) + sid = server_id if isinstance(server_id, bytes) else bytes_4(int(server_id)) + with patch_routing_wall_time(scenario.clock): + vhead = DeterministicScenario.voice_head_spec(base) + inject_announcement_ptt( + scenario.routing, + "MASTER-A", + vhead.data(), + pkt_time=scenario.clock.time(), + server_id=sid, + ) + for seq in range(1, 3): + spec = DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=min(seq, 4)) + assert inject_announcement_ptt( + scenario.routing, + "MASTER-A", + spec.data(), + pkt_time=scenario.clock.time(), + server_id=sid, + ) is True + assert len(scenario.capture.for_system("OBP-CL")) == 3 + assert len(scenario.capture.for_system("MASTER-A")) == 0 diff --git a/tests/routing/test_obp_announcement_contention.py b/tests/routing/test_obp_announcement_contention.py index 9390982..cb8ea9e 100644 --- a/tests/routing/test_obp_announcement_contention.py +++ b/tests/routing/test_obp_announcement_contention.py @@ -1,12 +1,31 @@ # ADN DMR Peer Server - OBP vs server announcement slot contention # # 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 +############################################################################### + +"""OBP ingress blocked while server announcement holds MASTER slot TX row.""" from __future__ import annotations from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time from tests.harness.scenarios import obp_bridge_scenario +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4 @@ -17,7 +36,7 @@ def _seed_server_broadcast_slot(scenario: DeterministicScenario, *, ann_tg: int "RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VHEAD, "TX_TIME": t, - "TX_RFS": bytes_3(5000), + "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID), "TX_TGID": bytes_3(ann_tg), "TX_STREAM_ID": bytes_4(0xA0A0A0A0), } diff --git a/tests/voice/test_announcement_anticollision.py b/tests/voice/test_announcement_anticollision.py index 6c865b4..00da6ac 100644 --- a/tests/voice/test_announcement_anticollision.py +++ b/tests/voice/test_announcement_anticollision.py @@ -24,6 +24,7 @@ from __future__ import annotations from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3 @@ -60,7 +61,7 @@ def test_broadcast_aborts_when_qso_starts_mid_transmission() -> None: uc = make_voice_uc(scenario, master) targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} - source = bytes_3(5000) + source = bytes_3(DEFAULT_SERVER_VOICE_ID) dst = bytes_3(91) master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD @@ -80,5 +81,5 @@ def test_mark_slots_busy_stamps_server_voice_tx_row() -> None: uc._mark_slots_busy(targets) assert slot["TX_TYPE"] == HBPF_SLT_VHEAD - assert int.from_bytes(slot["TX_RFS"], "big") == 5000 + assert int.from_bytes(slot["TX_RFS"], "big") == DEFAULT_SERVER_VOICE_ID assert slot["TX_TIME"] > 0 diff --git a/tests/voice/test_broadcast_queue.py b/tests/voice/test_broadcast_queue.py index c9aad74..b2eb720 100644 --- a/tests/voice/test_broadcast_queue.py +++ b/tests/voice/test_broadcast_queue.py @@ -24,6 +24,7 @@ from __future__ import annotations from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario +from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID from adn_server.domain import bytes_3 @@ -32,7 +33,7 @@ def test_enqueue_broadcast_queues_second_same_tg() -> None: uc = make_voice_uc(scenario, master) targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} - source = bytes_3(5000) + source = bytes_3(DEFAULT_SERVER_VOICE_ID) dst = bytes_3(91) uc._broadcast_active_tgs.add("91") @@ -47,7 +48,7 @@ def test_broadcast_queue_drains_after_first_finishes() -> None: uc = make_voice_uc(scenario, master) targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} - source = bytes_3(5000) + source = bytes_3(DEFAULT_SERVER_VOICE_ID) dst = bytes_3(91) uc._broadcast_active_tgs.add("91") diff --git a/tests/voice/test_inject_ptt_peer.py b/tests/voice/test_inject_ptt_peer.py new file mode 100644 index 0000000..a14204a --- /dev/null +++ b/tests/voice/test_inject_ptt_peer.py @@ -0,0 +1,122 @@ +# ADN DMR Peer Server - tests voice inject PTT peer resolution +# +# 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 +############################################################################### + +"""Synthetic announcement PTT uses SERVER_ID, not a hotspot peer id.""" + +from __future__ import annotations + +from tests.harness.scenarios import obp_bridge_scenario +from tests.harness.voice_helpers import FakeMasterForVoice, make_voice_uc + +from adn_server.domain import bytes_3, bytes_4, int_id + + +def test_inject_packet_peer_is_always_server_id() -> None: + scenario = obp_bridge_scenario("OBP-CL", tg=730600) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + peer = bytes_4(730039210) + sys_cfg = scenario.config["SYSTEMS"]["MASTER-A"] + sys_cfg["PEERS"] = {peer: {"CONNECTION": "YES", "OPTIONS": b"TS2=730500;SINGLE=1;"}} + pk = bytes_4(730039210) + sys_cfg.setdefault("_PEER_UA_MULTI_TGS", {}).setdefault(pk, {})[2] = {730600} + master = FakeMasterForVoice("MASTER-A") + master.STATUS[2] = {"RX_TYPE": 2, "TX_TYPE": 2, "RX_STREAM_ID": b"\x00" * 4} + scenario.protocols["MASTER-A"] = master + uc = make_voice_uc(scenario, master) + targets = [{"ts": 2}] + + peer_field = uc._announcement_packet_peer(730600, targets, bytes_4(730039101)) + + assert int_id(peer_field) == int_id(uc._global_server_id_bytes()) + assert int_id(peer_field) != 730039210 + assert int_id(peer_field) != 730039101 + + +def test_inject_emits_master_tx_voice_events_for_monitor() -> None: + scenario = obp_bridge_scenario("OBP-CL", tg=730600) + master = FakeMasterForVoice("MASTER-A") + uc = make_voice_uc(scenario, master) + events: list[str] = [] + uc._send_routing_event = events.append # type: ignore[method-assign] + stream_id = b"\xab\xcd\xef\x01" + rf_src = bytes_3(1000001) + pkts_by_ts = { + 2: [ + b"DMRD" + b"\x00" * 1 + rf_src + bytes_3(730600) + b"\x00" * 3 + b"\x80" + stream_id, + ], + } + targets = [{"name": "MASTER-A", "ts": 2, "slot": master.STATUS[2]}] + + uc._maybe_begin_legacy_voice_report(targets, pkts_by_ts, 730600) + + assert len(events) == 1 + assert events[0].startswith("GROUP VOICE,START,TX,MASTER-A,") + assert ",2,730600" in events[0] + assert uc._voice_report_state is not None + + +def test_inject_ptt_prefers_dynamic_ua_slot() -> None: + scenario = obp_bridge_scenario("OBP-CL", tg=730600) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + peer = bytes_4(730039210) + sys_cfg = scenario.config["SYSTEMS"]["MASTER-A"] + sys_cfg["PEERS"] = {peer: {"CONNECTION": "YES", "OPTIONS": b"TS2=730500;"}} + pk = bytes_4(730039210) + sys_cfg.setdefault("_PEER_UA_MULTI_TGS", {}).setdefault(pk, {})[2] = {730600} + master = FakeMasterForVoice("MASTER-A") + master.STATUS[2] = { + "RX_TYPE": 1, + "TX_TYPE": 1, + "RX_TGID": bytes_3(730600), + "TX_TGID": bytes_3(730600), + "RX_STREAM_ID": b"\x01" * 4, + } + master.STATUS[1] = {"RX_TYPE": 2, "TX_TYPE": 2, "RX_STREAM_ID": b"\x00" * 4} + scenario.protocols["MASTER-A"] = master + uc = make_voice_uc(scenario, master) + + targets, busy = uc._build_inject_ptt_targets("TTS-2", 730600) + + assert busy == 0 + assert len(targets) == 1 + assert targets[0]["ts"] == 2 + + +def test_inject_ptt_prefers_active_bridge_slot_for_static_tg() -> None: + scenario = obp_bridge_scenario("OBP-CL", tg=730500) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + peer = bytes_4(730039210) + sys_cfg = scenario.config["SYSTEMS"]["MASTER-A"] + sys_cfg["PEERS"] = {peer: {"CONNECTION": "YES", "OPTIONS": b"TS2=730500;"}} + master = FakeMasterForVoice("MASTER-A") + master.STATUS[2] = { + "RX_TYPE": 1, + "TX_TYPE": 1, + "TX_TGID": bytes_3(730500), + "RX_STREAM_ID": b"\x01" * 4, + } + master.STATUS[1] = {"RX_TYPE": 2, "TX_TYPE": 2, "RX_STREAM_ID": b"\x00" * 4} + scenario.protocols["MASTER-A"] = master + uc = make_voice_uc(scenario, master) + + targets, busy = uc._build_inject_ptt_targets("TTS-1", 730500) + + assert busy == 0 + assert targets[0]["ts"] == 2 diff --git a/tests/voice/test_scheduled_announcement.py b/tests/voice/test_scheduled_announcement.py index 6d659ec..cca8b87 100644 --- a/tests/voice/test_scheduled_announcement.py +++ b/tests/voice/test_scheduled_announcement.py @@ -47,6 +47,7 @@ def test_scheduled_announcement_starts_broadcast_on_idle_slot() -> None: assert delay == 0.5 drain_call_later(uc) assert len(master.sent) == 3 + assert len(scenario.capture.for_system("MASTER-B")) == 3 def test_scheduled_announcement_retries_when_slot_busy() -> None: @@ -87,4 +88,5 @@ def test_scheduled_announcement_sends_packets_to_master() -> None: assert uc._announcement_running[0] is False assert "91" not in uc._broadcast_active_tgs assert len(master.sent) == 3 + assert len(scenario.capture.for_system("MASTER-B")) == 3 assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM diff --git a/tests/voice/test_scheduled_tts.py b/tests/voice/test_scheduled_tts.py index cb9ae74..1739c69 100644 --- a/tests/voice/test_scheduled_tts.py +++ b/tests/voice/test_scheduled_tts.py @@ -25,7 +25,14 @@ from __future__ import annotations from datetime import datetime from unittest.mock import MagicMock, patch -from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario, voice_tts_config +from tests.harness.scenarios import obp_bridge_scenario +from tests.harness.voice_helpers import ( + FakeMasterForVoice, + drain_call_later, + make_voice_uc, + voice_master_scenario, + voice_tts_config, +) from adn_server.application.voice_use_cases import VoiceUseCases from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD @@ -43,6 +50,25 @@ def test_scheduled_tts_sync_path_enqueues_broadcast() -> None: assert getattr(uc._scheduled[0][1][0], "__name__", "") == "_tts_send_broadcast" +def test_inject_tts_starts_without_active_bridge_and_forwards_obp() -> None: + """Synthetic PTT must not require a pre-armed bridge row (UA relay created on inject).""" + scenario = obp_bridge_scenario("OBP-CL", tg=730500) + scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + scenario.seed_routing_table({}) + master = FakeMasterForVoice("MASTER-A") + master.STATUS[2] = {"RX_TYPE": 2, "TX_TYPE": 2, "RX_STREAM_ID": b"\x00" * 4} + scenario.protocols["MASTER-A"] = master + voice_tts_config(scenario, tg=730500, file_number="730500") + uc = make_voice_uc(scenario, master) + + uc._tts_conversion_done("/tmp/fake.ambe", 0, "730500", 730500, "en_GB", "interval", "TTS-1") + + assert uc._broadcast_active_tgs == {"730500"} + drain_call_later(uc) + assert len(scenario.capture.for_system("OBP-CL")) == 3 + assert len(master.sent) == 3 + + def test_tts_conversion_error_clears_running_flag() -> None: scenario, master = voice_master_scenario() uc = make_voice_uc(scenario, master) @@ -66,6 +92,8 @@ def test_tts_conversion_done_without_ambe_clears_running() -> None: def test_tts_conversion_done_retries_when_slot_busy() -> None: scenario, master = voice_master_scenario() + scenario.seed_routing_table({}) + master.STATUS[1]["RX_TYPE"] = HBPF_SLT_VHEAD master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD uc = make_voice_uc(scenario, master) uc._tts_running[0] = True diff --git a/tests/voice/test_tts_engine_serial.py b/tests/voice/test_tts_engine_serial.py new file mode 100644 index 0000000..742dcc6 --- /dev/null +++ b/tests/voice/test_tts_engine_serial.py @@ -0,0 +1,123 @@ +# ADN DMR Peer Server - tests voice TTS engine serialization +# +# 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 +############################################################################### + +"""TTS engine serial conversion lock (shared AMBEServer).""" + +from __future__ import annotations + +import socket +import threading +import time +import wave + +from adn_server.infrastructure.voice import tts_engine +from adn_server.infrastructure.voice.tts_engine import DV3K_SAMPLES_PER_FRAME + + +def test_text_to_ambe_serializes_parallel_conversions(tmp_path, monkeypatch) -> None: + """Only one full TTS conversion may run at a time (shared AMBEServer).""" + overlap = threading.Event() + active = 0 + peak_active = 0 + state_lock = threading.Lock() + + def fake_encode(*_args, **_kwargs) -> bool: + nonlocal active, peak_active + with state_lock: + active += 1 + peak_active = max(peak_active, active) + time.sleep(0.05) + with state_lock: + active -= 1 + overlap.set() + return True + + monkeypatch.setattr(tts_engine, "_generate_tts_audio", lambda *a, **k: True) + monkeypatch.setattr(tts_engine, "_convert_to_wav", lambda *a, **k: True) + monkeypatch.setattr(tts_engine, "_encode_ambe_ambeserver", fake_encode) + monkeypatch.setattr(tts_engine, "_cleanup", lambda *a, **k: None) + + paths: list[tuple[str, str]] = [] + for name in ("a", "b"): + txt = tmp_path / f"{name}.txt" + ambe = tmp_path / f"{name}.ambe" + txt.write_text(f"hello {name}", encoding="utf-8") + paths.append((str(txt), str(ambe))) + + errors: list[BaseException] = [] + + def run_one(txt_path: str, ambe_path: str) -> None: + try: + assert tts_engine.text_to_ambe( + txt_path, + ambe_path, + "es_ES", + "", + "127.0.0.1", + 2473, + ) + except BaseException as exc: + errors.append(exc) + + threads = [threading.Thread(target=run_one, args=p) for p in paths] + for thread in threads: + thread.start() + for thread in threads: + thread.join(timeout=5) + assert not thread.is_alive() + + assert errors == [] + assert overlap.is_set() + assert peak_active == 1 + + +def test_encode_ambe_ambeserver_aborts_on_consecutive_timeouts(monkeypatch, tmp_path) -> None: + wav_path = tmp_path / "test.wav" + ambe_path = tmp_path / "test.ambe" + + with wave.open(str(wav_path), "wb") as wf: + wf.setnchannels(1) + wf.setsampwidth(2) + wf.setframerate(8000) + wf.writeframes(b"\x00\x00" * (DV3K_SAMPLES_PER_FRAME * 25)) + + class _TimeoutSock: + def __init__(self) -> None: + self._step = 0 + + def sendto(self, *_args, **_kwargs) -> None: + return None + + def recvfrom(self, _size: int) -> tuple[bytes, tuple[str, int]]: + self._step += 1 + if self._step <= 2: + return (DV3K_PRODID_REQ, ("127.0.0.1", 2473)) + raise socket.timeout() + + def close(self) -> None: + return None + + monkeypatch.setattr(tts_engine.socket, "gethostbyname", lambda host: host) + monkeypatch.setattr(tts_engine.socket, "socket", lambda *a, **k: _TimeoutSock()) + monkeypatch.setattr(tts_engine, "_AMBESERVER_MAX_CONSECUTIVE_ERRORS", 3) + monkeypatch.setattr(tts_engine, "_AMBESERVER_ENCODE_TIMEOUT_S", 30.0) + + assert tts_engine._encode_ambe_ambeserver(str(wav_path), str(ambe_path), "127.0.0.1", 2473) is False + assert not ambe_path.exists() diff --git a/tests/voice/test_voice_config_reload.py b/tests/voice/test_voice_config_reload.py index c9f45b6..b172860 100644 --- a/tests/voice/test_voice_config_reload.py +++ b/tests/voice/test_voice_config_reload.py @@ -110,3 +110,29 @@ def test_apply_voice_config_starts_tts_loop() -> None: assert started == [30.0] assert 0 in uc._tts_tasks + + +def test_apply_voice_config_legacy_yaml_without_dmr_id() -> None: + """Legacy adn-voice.yaml without DMR_ID must reload without error.""" + scenario, _ = voice_master_scenario() + scenario.config["VOICE"] = { + "TTS_ANNOUNCEMENTS": [ + { + "ENABLED": True, + "TG": 730500, + "FILE": "730500", + "LANGUAGE": "es_ES", + "MODE": "interval", + "INTERVAL": 60, + } + ], + } + uc = VoiceUseCases( + FakeVoiceProvider(), + scenario.config, + start_looping_call=lambda *_a: MagicMock(running=True, stop=MagicMock()), + audio_path="/tmp/audio", + ) + uc.apply_voice_config() + + assert 0 in uc._tts_tasks