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.
pull/45/head
Rodrigo Pérez 3 months ago
parent 80823c59ff
commit 62f429aad9

@ -7,6 +7,11 @@
# Each announcement/TTS has its own LANGUAGE. ENABLED: true activates that item. # Each announcement/TTS has its own LANGUAGE. ENABLED: true activates that item.
VOICE: 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). # 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. # 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. # 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/<LANGUAGE>/ondemand/ # FILE: filename without .ambe, looked up in Audio/<LANGUAGE>/ondemand/
# MODE: "hourly" (once per hour) or "interval" (every INTERVAL seconds) # MODE: "hourly" (once per hour) or "interval" (every INTERVAL seconds)
# ENABLED: true = plays; false = configured but inactive (edit to true when ready) # 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: ANNOUNCEMENTS:
- ENABLED: true - ENABLED: true
FILE: announcement1 FILE: announcement1
TG: 2 TG: 2
# DMR_ID: 1000001
MODE: interval MODE: interval
INTERVAL: 60 INTERVAL: 60
LANGUAGE: es_ES LANGUAGE: es_ES
@ -77,12 +84,14 @@ VOICE:
- ENABLED: true - ENABLED: true
FILE: texto1 FILE: texto1
TG: 2 TG: 2
# DMR_ID: 1000001
MODE: interval MODE: interval
INTERVAL: 60 INTERVAL: 60
LANGUAGE: es_ES LANGUAGE: es_ES
- ENABLED: false - ENABLED: false
FILE: texto2 FILE: texto2
TG: 2 TG: 2
# DMR_ID: 1000001
MODE: hourly MODE: hourly
INTERVAL: 3600 INTERVAL: 3600
LANGUAGE: es_ES LANGUAGE: es_ES

@ -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. 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 | | Traffic | Destination in the DMR packet | Notes |
|---------|----------------------------------|--------| |---------|----------------------------------|--------|
| **Scheduled AMBE** (`ANNOUNCEMENTS`) | Whatever **`TG`** you set in `adn-voice.yaml` | Source ID **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`** | Source ID **5000**. | | **TTS** (`TTS_ANNOUNCEMENTS`) | Same — configured **`TG`** | Same. |
| **On-demand** (TG **9991–9999**) | **TG 9** | Short info clips; source **5000** (see [Voice, announcements, and TTS](voice-and-tts.md)). | | **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** | Source **5000**. | | **Disconnected / reflector prompts** | **TG 9** | Same. |
| **Voice ident** | **All-call** (`16777215`) or **`OVERRIDE_IDENT_TG`** if set | Source **5000**. | | **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) ### 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: 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** - **Destination TG 9**
- **Timeslot 2** (the code drives the **TS2** slot for that hotspot) - **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. - Trigger path is **private VTERM** for destination **9991–9999**, then async playback generation.
- Works from **MASTER** and **PEER** paths. - 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) ## 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 | | 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` | | **5000** (destination) | No auto user-activated bridge if missing from `BRIDGES` |
| **4000** (group) | Deactivate dynamic bridges | | **4000** (group) | Deactivate dynamic bridges |
| **4000** (unit) | Disconnect dynamics; not routed as PC | | **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 | | **9** | Service/prompt lane (short server audio); reserved for auto-bridges; internal **TGID** on TS2 legs |
| **9990** | Echo bridge TG (with ECHO system) | | **9990** | Echo bridge TG (with ECHO system) |
| **16777215** | All-call (default voice-ident destination unless overridden) | | **16777215** | All-call (default voice-ident destination unless overridden) |

@ -7,9 +7,26 @@
Template: `adn-voice.example.yaml`. 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). - 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). - **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. 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`**). - **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 ## 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. 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 ## Broadcast queue
Parallel **broadcasts** on **different** TGs may run concurrently; **same-TG** broadcasts are **serialised** so one announcement completes before another on that TG. Parallel **broadcasts** on **different** TGs may run concurrently; **same-TG** broadcasts are **serialised** so one announcement completes before another on that TG.

@ -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. 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 | | Tráfico | Destino en el paquete DMR | Notas |
|---------|---------------------------|--------| |---------|---------------------------|--------|
| **AMBE programado** (`ANNOUNCEMENTS`) | El **`TG`** que definas en `adn-voice.yaml` | ID de 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 | ID de fuente **5000**. | | **TTS** (`TTS_ANNOUNCEMENTS`) | Igual — **`TG`** configurado | Igual. |
| **Bajo demanda** (TG **9991–9999**) | **TG 9** | Clips informativos cortos; fuente **5000** (ver [Voz, anuncios y TTS](voice-and-tts.md)). | | **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** | Fuente **5000**. | | **Desconectado / reflector** | **TG 9** | Igual. |
| **Ident por voz** | **All-call** (`16777215`) o **`OVERRIDE_IDENT_TG`** si está definido | Fuente **5000**. | | **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) ### 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: 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** - **TG de destino 9**
- **Timeslot 2** (el código usa el slot **TS2** para ese hotspot) - **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. - 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**. - 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) ## 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 | | 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` | | **5000** (destino) | Sin bridge UA automático si falta en `BRIDGES` |
| **4000** (grupo) | Desactivar bridges dinámicos | | **4000** (grupo) | Desactivar bridges dinámicos |
| **4000** (unitaria) | Desconectar dinámicos; no enrutada como PC | | **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 | | **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) | | **9990** | TG de bridge de eco (con sistema ECHO) |
| **16777215** | All-call (destino por defecto de ident por voz salvo sobrescritura) | | **16777215** | All-call (destino por defecto de ident por voz salvo sobrescritura) |

@ -7,9 +7,26 @@
Plantilla: `adn-voice.example.yaml`. 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). - **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). - 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`**). - **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 ## 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. 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 ## 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. 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.

@ -31,6 +31,7 @@ import time
from typing import Any, Callable from typing import Any, Callable
from ..domain import HBPF_SLT_VTERM, bytes_3, int_id from ..domain import HBPF_SLT_VTERM, bytes_3, int_id
from .server_voice import server_voice_rf_src_bytes
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -119,7 +120,7 @@ class IdentUseCases:
continue continue
_all_call = bytes_3(16777215) _all_call = bytes_3(16777215)
_source_id = bytes_3(5000) _source_id = server_voice_rf_src_bytes(self._config)
_dst_id = b"" _dst_id = b""
override_tg = sys_cfg.get("OVERRIDE_IDENT_TG") override_tg = sys_cfg.get("OVERRIDE_IDENT_TG")
if override_tg is not None and int(override_tg) > 0 and int(override_tg) < 16777215: if override_tg is not None and int(override_tg) > 0 and int(override_tg) < 16777215:

@ -0,0 +1,99 @@
# ADN DMR Peer Server - announcement synthetic PTT ingress
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
#
# Derived from ADN DMR Server / FreeDMR / HBlink. Original license:
###############################################################################
# Copyright (C) 2026 Joaquin Madrid Belando, EA5GVK <ea5gvk@gmail.com>
# Copyright (C) 2020 Simon Adlem, G7RZU <g7rzu@gb7fr.org.uk>
# Copyright (C) 2016-2019 Cortney T. Buffington, N0MJS <n0mjs@me.com>
#
# 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,
)

@ -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 import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, bytes_3, bytes_4, int_id
from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO from ...domain.hbp_protocol import HBPF_SLT_VTERM, STREAM_TO
from ..server_voice import DEFAULT_SERVER_VOICE_ID
PeerVoiceSlotRow = dict[str, Any] PeerVoiceSlotRow = dict[str, Any]
PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow] PeerVoiceSlotMap = dict[int, PeerVoiceSlotRow]
@ -461,17 +462,30 @@ def inject_only_defer_obp_hbp_slot_contention(
return source_is_hbp and connected_count > 1 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. """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 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). 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 return False
tx_type = slot_st.get("TX_TYPE") tx_type = slot_st.get("TX_TYPE")
if tx_type is None or tx_type == HBPF_SLT_VTERM: 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 return 9991 <= int(dst_id) <= 9999
def is_server_originated_voice(packet: bytes) -> bool: def is_server_originated_voice(
"""True for server group playback (src 5000 -> TG 9 TS2), e.g. on-demand / disconnected.""" 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: if len(packet) < 11:
return False return False
burst = parse_dmrd_burst_fields(packet) burst = parse_dmrd_burst_fields(packet)
@ -1062,7 +1081,16 @@ def is_server_originated_voice(packet: bytes) -> bool:
slot, _, _, _, dst_id, call_type = burst slot, _, _, _, dst_id, call_type = burst
if call_type not in ("group", "vcsbk"): if call_type not in ("group", "vcsbk"):
return False 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: def parse_dmrd_route_fields(packet: bytes) -> tuple[int, int, str] | None:

@ -64,6 +64,7 @@ from .routing.store_authority_mixin import StoreAuthorityMixin
from .routing.subscription_table import SubscriptionTableMixin from .routing.subscription_table import SubscriptionTableMixin
from .routing.timers import RoutingTimerMixin from .routing.timers import RoutingTimerMixin
from .routing.voice_subscription import VoiceSubscriptionMixin from .routing.voice_subscription import VoiceSubscriptionMixin
from .server_voice import all_server_voice_ids
from .talker_alias_use_cases import TalkerAliasUseCases from .talker_alias_use_cases import TalkerAliasUseCases
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -692,7 +693,11 @@ class RoutingUseCases(
_ts_st["TX_TYPE"] = HBPF_SLT_VTERM _ts_st["TX_TYPE"] = HBPF_SLT_VTERM
_apply_master_slot_contention = ( _apply_master_slot_contention = (
not _defer_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 ( if (
_apply_master_slot_contention _apply_master_slot_contention

@ -0,0 +1,114 @@
# ADN DMR Peer Server - server-originated voice identity
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""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)

@ -31,8 +31,13 @@ import time
from datetime import datetime from datetime import datetime
from typing import Any, Callable 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 .ports import VoiceProvider
from .server_voice import (
announcement_item_source_bytes,
server_voice_id,
server_voice_rf_src_bytes,
)
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
@ -55,6 +60,9 @@ class VoiceUseCases:
call_later: Callable[..., Any] | None = None, call_later: Callable[..., Any] | None = None,
start_looping_call: Callable[[Callable[[], None], float, bool], Any] | None = None, start_looping_call: Callable[[Callable[[], None], float, bool], Any] | None = None,
defer_to_thread: Callable[..., 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: ) -> None:
self._voice = voice_provider self._voice = voice_provider
self._config = config self._config = config
@ -65,6 +73,10 @@ class VoiceUseCases:
self._call_later = call_later self._call_later = call_later
self._start_looping_call = start_looping_call self._start_looping_call = start_looping_call
self._defer_to_thread = defer_to_thread 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._ann_tasks: dict[int, Any] = {}
self._tts_tasks: dict[int, Any] = {} self._tts_tasks: dict[int, Any] = {}
self._announcement_running: dict[int, bool] = {} self._announcement_running: dict[int, bool] = {}
@ -74,6 +86,9 @@ class VoiceUseCases:
self._broadcast_queue: list[dict[str, Any]] = [] self._broadcast_queue: list[dict[str, Any]] = []
self._broadcast_active_tgs: set[str] = set() 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]]: def get_ambe_words(self, languages: str, audio_path: str) -> dict[str, dict[str, Any]]:
"""Load AMBE words for given languages (legacy readAMBE.readfiles).""" """Load AMBE words for given languages (legacy readAMBE.readfiles)."""
return self._voice.get_ambe_words(languages, audio_path) 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).""" """Generate HBP voice packets for phrase (legacy mk_voice.pkt_gen)."""
return self._voice.pkt_gen(rf_src, dst_id, peer, slot, phrase) 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( def _build_announcement_targets(
self, tg_int: int, tg_str: str, label: str self, tg_int: int, tg_str: str, label: str
) -> tuple[list[dict[str, Any]], int]: ) -> 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 Inject path: synthetic PTT on the TG's dynamic UA slot when applicable.
because a QSO is active (RX/TX not VTERM), for anti-collision retry scheduling. 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]] = [] targets: list[dict[str, Any]] = []
busy_count = 0 busy_count = 0
protocols = self._get_protocols() if self._get_protocols else {} 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}) targets.append({"sys_obj": sys_obj, "name": sys_name, "slot": slot, "ts": ts})
return targets, busy_count 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( def _send_filtered_by_tg(
self, sys_obj: Any, pkt: bytes, tg: int, ts: int, bridges: dict[str, list[dict[str, Any]]] self, sys_obj: Any, pkt: bytes, tg: int, ts: int, bridges: dict[str, list[dict[str, Any]]]
) -> int: ) -> int:
@ -146,9 +273,150 @@ class VoiceUseCases:
return -1 return -1
return 0 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.""" """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() now = time.time()
for t in targets: for t in targets:
try: try:
@ -181,10 +449,10 @@ class VoiceUseCases:
'source_id': source_id, 'dst_id': dst_id, 'tg': tg, 'num': num, 'label': label, 'source_id': source_id, 'dst_id': dst_id, 'tg': tg, 'num': num, 'label': label,
}) })
_pos = len(self._broadcast_queue) _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: else:
self._broadcast_active_tgs.add(_tg_key) 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)) logger.info('(%s) Starting broadcast immediately for TG %s (active TGs: %s)', label, tg, len(self._broadcast_active_tgs))
if self._call_later: if self._call_later:
if _type == 'ann': if _type == 'ann':
@ -207,8 +475,8 @@ class VoiceUseCases:
_label = _next['label'] _label = _next['label']
_tg_key = str(_next['tg']) _tg_key = str(_next['tg'])
self._broadcast_active_tgs.add(_tg_key) self._broadcast_active_tgs.add(_tg_key)
self._mark_slots_busy(_next['targets']) self._mark_slots_busy(_next['targets'], _next['source_id'])
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)) 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 self._call_later:
if _type == 'ann': 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) 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: if tg is not None:
self._broadcast_active_tgs.discard(str(tg)) self._broadcast_active_tgs.discard(str(tg))
if self._broadcast_queue: 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: if self._call_later:
self._call_later(_BROADCAST_GAP, self._start_next_broadcast) self._call_later(_BROADCAST_GAP, self._start_next_broadcast)
else: else:
if not self._broadcast_active_tgs: if not self._broadcast_active_tgs:
logger.info('(QUEUE) All broadcasts finished, queue empty') logger.info('(BROADCAST) All announcement/TTS playbacks finished')
else: 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( def _announcement_send_broadcast(
self, self,
@ -244,6 +518,7 @@ class VoiceUseCases:
total = len(pkts_by_ts.get(1, [])) total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total or not targets: if pkt_idx >= total or not targets:
self._mark_slots_free(targets) self._mark_slots_free(targets)
self._maybe_end_legacy_voice_report()
for t in targets: for t in targets:
try: try:
obj = t.get("sys_obj") obj = t.get("sys_obj")
@ -283,6 +558,7 @@ class VoiceUseCases:
targets.remove(t) targets.remove(t)
if not targets: if not targets:
self._announcement_running[ann_idx] = False self._announcement_running[ann_idx] = False
self._maybe_end_legacy_voice_report()
logger.info( logger.info(
"(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s", "(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s",
label, label,
@ -291,31 +567,10 @@ class VoiceUseCases:
) )
self._broadcast_finished(tg) self._broadcast_finished(tg)
return return
bridges = self._routing_table_for_report() if self._routing_table_for_report else {}
now = time.time() now = time.time()
for t in targets: if pkt_idx == 0:
try: self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg)
sys_obj = t["sys_obj"] self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label)
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 next_time is None: if next_time is None:
next_time = now + _FRAME_INTERVAL next_time = now + _FRAME_INTERVAL
else: else:
@ -369,7 +624,7 @@ class VoiceUseCases:
if not _file or not _tg: if not _file or not _tg:
return return
_dst_id = bytes_3(_tg) _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") server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes): if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id)) server_id = bytes_3(int(server_id))
@ -400,9 +655,10 @@ class VoiceUseCases:
if mode == "hourly": if mode == "hourly":
self._announcement_last_hour[ann_idx] = datetime.now().hour self._announcement_last_hour[ann_idx] = datetime.now().hour
_say_list = [_say] _say_list = [_say]
pkt_peer = self._announcement_packet_peer(_tg, targets, server_id)
pkts_by_ts = { pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _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, server_id, 1, _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) ts1_count = sum(1 for t in targets if t["ts"] == 1)
ts2_count = sum(1 for t in targets if t["ts"] == 2) ts2_count = sum(1 for t in targets if t["ts"] == 2)
@ -475,7 +731,12 @@ class VoiceUseCases:
return return
logger.info("(%s) Playing TTS file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang) 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) _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") server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes): if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id)) server_id = bytes_3(int(server_id))
@ -515,9 +776,10 @@ class VoiceUseCases:
if mode == "hourly": if mode == "hourly":
self._tts_last_hour[tts_idx] = datetime.now().hour self._tts_last_hour[tts_idx] = datetime.now().hour
_say_list = [_say] _say_list = [_say]
pkt_peer = self._announcement_packet_peer(_tg, targets, server_id)
pkts_by_ts = { pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, server_id, 0, _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, server_id, 1, _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)) 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) 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, [])) total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total or not targets: if pkt_idx >= total or not targets:
self._mark_slots_free(targets) self._mark_slots_free(targets)
self._maybe_end_legacy_voice_report()
for t in targets: for t in targets:
try: try:
obj = t.get("sys_obj") obj = t.get("sys_obj")
@ -586,6 +849,7 @@ class VoiceUseCases:
targets.remove(t) targets.remove(t)
if not targets: if not targets:
self._tts_running[tts_idx] = False self._tts_running[tts_idx] = False
self._maybe_end_legacy_voice_report()
logger.info( logger.info(
"(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s", "(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s",
label, label,
@ -594,31 +858,10 @@ class VoiceUseCases:
) )
self._broadcast_finished(tg) self._broadcast_finished(tg)
return return
bridges = self._routing_table_for_report() if self._routing_table_for_report else {}
now = time.time() now = time.time()
for t in targets: if pkt_idx == 0:
try: self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg)
sys_obj = t["sys_obj"] self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label)
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 next_time is None: if next_time is None:
next_time = now + _FRAME_INTERVAL next_time = now + _FRAME_INTERVAL
else: else:
@ -659,12 +902,12 @@ class VoiceUseCases:
logger.info("(%s) Playing on-demand AMBE file: %s (ID: %s)", system, file_number, file_number) logger.info("(%s) Playing on-demand AMBE file: %s (ID: %s)", system, file_number, file_number)
time.sleep(1) time.sleep(1)
_say = [pairs] _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) time.sleep(1)
_slot = protocol.STATUS.get(2) _slot = protocol.STATUS.get(2)
if not _slot: if not _slot:
return return
_source_id = bytes_3(5000)
_dst_id = bytes_3(9) _dst_id = bytes_3(9)
_next_time = time.time() _next_time = time.time()
_pkt_count = 0 _pkt_count = 0
@ -708,7 +951,8 @@ class VoiceUseCases:
else: else:
_say.append(words.get("notlinked") or silence) _say.append(words.get("notlinked") or silence)
_say.append(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) time.sleep(1)
_slot = protocol.STATUS.get(2) _slot = protocol.STATUS.get(2)
if not _slot: if not _slot:
@ -720,7 +964,7 @@ class VoiceUseCases:
_delay = _next_time - time.time() _delay = _next_time - time.time()
if _delay > 0.001: if _delay > 0.001:
time.sleep(_delay) 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) logger.debug("(%s) disconnected voice thread end", system)
def apply_voice_config(self) -> None: def apply_voice_config(self) -> None:

@ -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.echo_seed import seed_echo_routing_table
from adn_server.application.subscription.store_sync import replace_store_from_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.domain.dmr.bptc import encode_emblc
from adn_server.infrastructure.acl_router import InMemoryAclRouter from adn_server.infrastructure.acl_router import InMemoryAclRouter
from adn_server.infrastructure.config_normalizer import ( 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 routing_table_for_report = routing_use_cases.routing_table_for_report
voice_use_cases._routing_table_for_report = 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_routing_table(routing_table_for_report())
report_factory.set_systems(config.get("SYSTEMS", {})) report_factory.set_systems(config.get("SYSTEMS", {}))

@ -39,6 +39,7 @@ from twisted.internet import reactor, task
from twisted.internet.protocol import DatagramProtocol from twisted.internet.protocol import DatagramProtocol
from ...application.proxy.deployment import is_proxy_inject_only from ...application.proxy.deployment import is_proxy_inject_only
from ...application.server_voice import all_server_voice_ids
from ...application.routing.downlink import ( from ...application.routing.downlink import (
DownlinkContext, DownlinkContext,
iter_downlink_voice_slots, iter_downlink_voice_slots,
@ -398,7 +399,9 @@ class HBPProtocol(DatagramProtocol):
connected = self._cached_connected_peer_count() connected = self._cached_connected_peer_count()
if connected <= 0: if connected <= 0:
return () 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, {}) slot_st = self.STATUS.get(2, {})
rx_peer = slot_st.get("RX_PEER", b"") rx_peer = slot_st.get("RX_PEER", b"")
if rx_peer and rx_peer in self._peers: if rx_peer and rx_peer in self._peers:
@ -575,8 +578,10 @@ class HBPProtocol(DatagramProtocol):
slot, tgid, call_type = parsed slot, tgid, call_type = parsed
if call_type not in ("group", "vcsbk"): if call_type not in ("group", "vcsbk"):
return True return True
# Server playback (5000 -> TG 9): not in OPTIONS; deliver to requesting RX peer on TS2. # Server playback (server voice ID -> TG 9): not in OPTIONS; deliver to requesting RX peer on TS2.
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, {}) slot_st = self.STATUS.get(2, {})
rx_peer = slot_st.get("RX_PEER", b"") rx_peer = slot_st.get("RX_PEER", b"")
if rx_peer and rx_peer != b"\x00\x00\x00\x00": if rx_peer and rx_peer != b"\x00\x00\x00\x00":

@ -34,11 +34,21 @@ import os
import socket import socket
import struct import struct
import subprocess import subprocess
import threading
import time
import wave import wave
from typing import Any from typing import Any
logger = logging.getLogger(__name__) 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 = { _LANG_MAP = {
"es_ES": "es", "en_GB": "en", "en_US": "en", "fr_FR": "fr", "es_ES": "es", "en_GB": "en", "en_US": "en", "fr_FR": "fr",
"de_DE": "de", "it_IT": "it", "pt_PT": "pt", "pt_BR": "pt", "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 return False
try: try:
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
sock.settimeout(5.0) sock.settimeout(_AMBESERVER_FRAME_TIMEOUT_S)
except Exception as e: except Exception as e:
logger.error("(TTS-AMBESERVER) Error creating UDP socket: %s", e) logger.error("(TTS-AMBESERVER) Error creating UDP socket: %s", e)
return False return False
@ -212,12 +222,33 @@ def _encode_ambe_ambeserver(wav_path: str, ambe_path: str, host: str, port: int)
return False return False
_total_frames = wf.getnframes() _total_frames = wf.getnframes()
_sample_rate = wf.getframerate() _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) _raw_frames = wf.readframes(_total_frames)
wf.close() wf.close()
_samples = list(struct.unpack("<" + "h" * _total_frames, _raw_frames)) _samples = list(struct.unpack("<" + "h" * _total_frames, _raw_frames))
_ambe_frames: list[bytes] = [] _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): 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] _chunk = _samples[i : i + DV3K_SAMPLES_PER_FRAME]
if len(_chunk) < DV3K_SAMPLES_PER_FRAME: if len(_chunk) < DV3K_SAMPLES_PER_FRAME:
_chunk = _chunk + [0] * (DV3K_SAMPLES_PER_FRAME - len(_chunk)) _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) _ambe_data = _parse_ambe_response(data)
if _ambe_data is not None: if _ambe_data is not None:
_ambe_frames.append(_ambe_data) _ambe_frames.append(_ambe_data)
except (socket.timeout, Exception): _frames_sent += 1
pass _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() sock.close()
if not _ambe_frames: if not _ambe_frames:
logger.error("(TTS-AMBESERVER) No AMBE frames received") 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: except Exception as e:
logger.error("(TTS-AMBESERVER) Error writing AMBE: %s", e) logger.error("(TTS-AMBESERVER) Error writing AMBE: %s", e)
return False 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 return True
@ -254,24 +320,33 @@ def _cleanup(files: list[str]) -> None:
logger.warning("(TTS) cleanup remove %s failed: %s", f, e) 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, txt_path: str,
ambe_path: str, ambe_path: str,
language: str, language: str,
vocoder_cmd: str, vocoder_cmd: str,
ambeserver_host: str = "", ambeserver_host: str,
ambeserver_port: int = 2460, ambeserver_port: int,
volume_db: int = 0, volume_db: int,
speed: float = 1.0, speed: float,
) -> bool: ) -> 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: with open(txt_path, "r", encoding="utf-8") as f:
text = f.read().strip() text = f.read().strip()
if not text: if not text:
@ -308,6 +383,39 @@ def text_to_ambe(
return True 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: 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.""" """Ensure .ambe exists for TTS item; create from .txt if needed. Returns path or None."""
if not item.get("ENABLED", False): if not item.get("ENABLED", False):

@ -1,6 +1,24 @@
# ADN DMR Peer Server - server broadcast slot contention # ADN DMR Peer Server - server broadcast slot contention
# #
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl> # Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""MASTER slot contention while server scheduled voice holds TX row."""
from __future__ import annotations from __future__ import annotations
@ -8,6 +26,7 @@ from adn_server.application.routing.helpers import (
hbp_slot_blocks_group_voice, hbp_slot_blocks_group_voice,
master_slot_holds_server_broadcast, 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 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 = { slot = {
"TX_TYPE": HBPF_SLT_VHEAD, "TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": 100.0, "TX_TIME": 100.0,
"TX_RFS": bytes_3(5000), "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID),
"TX_TGID": bytes_3(91), "TX_TGID": bytes_3(91),
} }
assert master_slot_holds_server_broadcast(slot, 100.1) is True 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, "RX_TYPE": HBPF_SLT_VTERM,
"TX_TYPE": HBPF_SLT_VHEAD, "TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": 200.0, "TX_TIME": 200.0,
"TX_RFS": bytes_3(5000), "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID),
"TX_TGID": bytes_3(91), "TX_TGID": bytes_3(91),
} }
blocked = hbp_slot_blocks_group_voice( blocked = hbp_slot_blocks_group_voice(

@ -367,6 +367,23 @@ def test_remap_voice_event_for_single_peer() -> None:
assert mapped.startswith("GROUP VOICE,START,TX,SYSTEM-7,") 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: def test_non_proxy_systems_pass_through_unchanged() -> None:
systems = { systems = {
"ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}},

@ -0,0 +1,105 @@
# ADN DMR Peer Server - tests application server voice identity
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""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

@ -24,6 +24,10 @@ from __future__ import annotations
from typing import Any, Iterator 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.application.voice_use_cases import VoiceUseCases
from adn_server.domain import bytes_3, bytes_4 from adn_server.domain import bytes_3, bytes_4
from tests.harness.deterministic import DeterministicScenario, FakeHbpProtocol, active_routing_table from tests.harness.deterministic import DeterministicScenario, FakeHbpProtocol, active_routing_table
@ -43,8 +47,9 @@ class FakeVoiceProvider:
) -> Iterator[bytes]: ) -> Iterator[bytes]:
del phrase del phrase
ts_bit = 0x80 if slot else 0 ts_bit = 0x80 if slot else 0
stream_id = bytes_4(0xA0A0A0A0)
for seq in range(3): 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: def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
del audio_path, lang, file_number del audio_path, lang, file_number
@ -68,10 +73,11 @@ class FakeMasterForVoice(FakeHbpProtocol):
def voice_master_scenario(tg: int = 91) -> tuple[DeterministicScenario, FakeMasterForVoice]: def voice_master_scenario(tg: int = 91) -> tuple[DeterministicScenario, FakeMasterForVoice]:
config = DeterministicScenario().config config = DeterministicScenario().config
master = FakeMasterForVoice("MASTER-A") master = FakeMasterForVoice("MASTER-A")
config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
config["SYSTEMS"]["MASTER-A"]["PEERS"] = { config["SYSTEMS"]["MASTER-A"]["PEERS"] = {
"1001": {"CALLSIGN": "TEST", "IP": "127.0.0.1", "PORT": 62032}, "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 = DeterministicScenario(config=config, routing_table=bridges)
scenario.protocols["MASTER-A"] = master scenario.protocols["MASTER-A"] = master
master.STATUS[2] = { master.STATUS[2] = {
@ -108,19 +114,39 @@ def make_voice_uc(
audio_path: str = "/tmp/audio", audio_path: str = "/tmp/audio",
) -> VoiceUseCases: ) -> VoiceUseCases:
scheduled: list[tuple[float, tuple]] = [] 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): def call_later(delay, fn, *args):
scheduled.append((delay, (fn, args))) scheduled.append((delay, (fn, args)))
return type("H", (), {"active": lambda self: True, "cancel": lambda self: None})() 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( uc = VoiceUseCases(
FakeVoiceProvider(), FakeVoiceProvider(),
scenario.config, 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, routing_table_for_report=scenario.routing.routing_table_for_report,
call_later=call_later, call_later=call_later,
audio_path=audio_path, audio_path=audio_path,
inject_announcement_ptt=inject_announcement_ptt_cb,
) )
uc._announcement_ptt_system = ptt_system
uc._scheduled = scheduled uc._scheduled = scheduled
return uc return uc

@ -26,6 +26,7 @@ import pytest
from tests.harness.deterministic import DeterministicScenario, PacketSpec from tests.harness.deterministic import DeterministicScenario, PacketSpec
from tests.support.hbp_repeat_stack import build_hbp_repeat_stack 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 import bytes_3, bytes_4
from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM
from adn_server.infrastructure.hbp_constants import DMRD 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: 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() stack = _inject_proxy_stack()
caller = bytes_4(730039101) caller = bytes_4(730039101)
other = bytes_4(730039102) 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(caller, addr_caller, options="TS2=730,730444;")
stack.register_peer(other, addr_other, options="TS2=91;") stack.register_peer(other, addr_other, options="TS2=91;")
spec = PacketSpec( spec = PacketSpec(
peer_id=5000, peer_id=DEFAULT_SERVER_VOICE_ID,
rf_src=5000, rf_src=DEFAULT_SERVER_VOICE_ID,
dst_id=9, dst_id=9,
slot=2, slot=2,
stream_id=0xCAFEBABE, stream_id=0xCAFEBABE,

@ -0,0 +1,77 @@
# ADN DMR Peer Server - announcement synthetic PTT inject
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""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

@ -1,12 +1,31 @@
# ADN DMR Peer Server - OBP vs server announcement slot contention # ADN DMR Peer Server - OBP vs server announcement slot contention
# #
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl> # Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""OBP ingress blocked while server announcement holds MASTER slot TX row."""
from __future__ import annotations from __future__ import annotations
from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time
from tests.harness.scenarios import obp_bridge_scenario 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 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, "RX_TYPE": HBPF_SLT_VTERM,
"TX_TYPE": HBPF_SLT_VHEAD, "TX_TYPE": HBPF_SLT_VHEAD,
"TX_TIME": t, "TX_TIME": t,
"TX_RFS": bytes_3(5000), "TX_RFS": bytes_3(DEFAULT_SERVER_VOICE_ID),
"TX_TGID": bytes_3(ann_tg), "TX_TGID": bytes_3(ann_tg),
"TX_STREAM_ID": bytes_4(0xA0A0A0A0), "TX_STREAM_ID": bytes_4(0xA0A0A0A0),
} }

@ -24,6 +24,7 @@ from __future__ import annotations
from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario 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 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) uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(5000) source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91) dst = bytes_3(91)
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD 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) uc._mark_slots_busy(targets)
assert slot["TX_TYPE"] == HBPF_SLT_VHEAD 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 assert slot["TX_TIME"] > 0

@ -24,6 +24,7 @@ from __future__ import annotations
from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario 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 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) uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(5000) source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91) dst = bytes_3(91)
uc._broadcast_active_tgs.add("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) uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}] targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]} pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(5000) source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91) dst = bytes_3(91)
uc._broadcast_active_tgs.add("91") uc._broadcast_active_tgs.add("91")

@ -0,0 +1,122 @@
# ADN DMR Peer Server - tests voice inject PTT peer resolution
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""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

@ -47,6 +47,7 @@ def test_scheduled_announcement_starts_broadcast_on_idle_slot() -> None:
assert delay == 0.5 assert delay == 0.5
drain_call_later(uc) drain_call_later(uc)
assert len(master.sent) == 3 assert len(master.sent) == 3
assert len(scenario.capture.for_system("MASTER-B")) == 3
def test_scheduled_announcement_retries_when_slot_busy() -> None: 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 uc._announcement_running[0] is False
assert "91" not in uc._broadcast_active_tgs assert "91" not in uc._broadcast_active_tgs
assert len(master.sent) == 3 assert len(master.sent) == 3
assert len(scenario.capture.for_system("MASTER-B")) == 3
assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM

@ -25,7 +25,14 @@ from __future__ import annotations
from datetime import datetime from datetime import datetime
from unittest.mock import MagicMock, patch 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.application.voice_use_cases import VoiceUseCases
from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD 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" 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: def test_tts_conversion_error_clears_running_flag() -> None:
scenario, master = voice_master_scenario() scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master) 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: def test_tts_conversion_done_retries_when_slot_busy() -> None:
scenario, master = voice_master_scenario() 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 master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD
uc = make_voice_uc(scenario, master) uc = make_voice_uc(scenario, master)
uc._tts_running[0] = True uc._tts_running[0] = True

@ -0,0 +1,123 @@
# ADN DMR Peer Server - tests voice TTS engine serialization
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""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()

@ -110,3 +110,29 @@ def test_apply_voice_config_starts_tts_loop() -> None:
assert started == [30.0] assert started == [30.0]
assert 0 in uc._tts_tasks 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

Loading…
Cancel
Save

Powered by TurnKey Linux.