Merge pull request #106 from pyopower/feat/announcements-plugin

refactor(voice): scheduled announcements, TTS and beacons as a plugin
develop
ce5rpy 3 days ago committed by GitHub
commit bec05936f9
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

5
.gitignore vendored

@ -36,7 +36,7 @@ json/*
# Internal
docs-priv/
# Runtime drop-in plugins (see plugins/README.md); example skeleton is versioned
# Runtime drop-in plugins (see plugins/README.md); example skeleton and official plugins are versioned
plugins/*
!plugins/README.md
!plugins/.gitkeep
@ -44,6 +44,9 @@ plugins/*
!plugins/example/**
plugins/example/example-events/
plugins/example/**/__pycache__/
!plugins/voice-announcements/
!plugins/voice-announcements/**
plugins/voice-announcements/**/__pycache__/
# Runtime TTS on-demand cache (generated at runtime)
Audio/es_ES/ondemand/

@ -1,5 +1,7 @@
# Voice, announcements, and TTS
> **Scheduled announcements and TTS are a plugin** — `plugins/voice-announcements/`, enabled by default. Configuration is unchanged (this file); the plugin follows it every 15 s and the server grants it exactly the talkgroups and DMR IDs its items use (override with `PLUGINS.send.voice-announcements`, see [Plugins](plugins.md#sending-unit-data-and-group-voice-opt-in)). Voice ident, on-demand 999x and disconnected prompts stay in the core.
## Configuration files
- **`adn-voice.yaml`** (optional, not committed) — merged into `config["VOICE"]`.
@ -62,7 +64,7 @@ Configure **`TTS_VOCODER_CMD`** or **`TTS_AMBESERVER_HOST`** / **`TTS_AMBESERVER
When scheduling announcements, the server may **wait** if target slots are busy, **drop** targets if a live QSO appears mid-stream, and only mark **hourly** announcement state after a successful target list — this avoids clobbering live traffic.
Scheduled and TTS broadcasts inject each frame as a **synthetic hotspot PTT** on `PROXY.TARGET_SYSTEM` (the same MASTER used by the integrated proxy). Routing fans out to bridged hotspots and OPENBRIDGE legs; local peers on that MASTER still receive frames via `send_system`.
The plugin sends each frame through `send_dmrd`, which the core injects 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`. While it plays, the stream holds that MASTER slot; a radio taking the slot stops it.
## Broadcast queue

@ -1,5 +1,7 @@
# Voz, anuncios y TTS
> **Los anuncios programados y el TTS son un plugin** — `plugins/voice-announcements/`, activado por defecto. La configuración no cambia (este fichero); el plugin la sigue cada 15 s y el servidor le autoriza exactamente los TGs e IDs DMR que usan sus ítems (se puede ajustar con `PLUGINS.send.voice-announcements`, ver [Plugins](plugins.md#envío-de-datos-y-voz-de-grupo-opcional)). El ident de voz, los 999x a demanda y el aviso de desconexión siguen en el núcleo.
## Ficheros de configuración
- **`adn-voice.yaml`** (opcional, no versionada) — fusionada en `config["VOICE"]`.
@ -63,7 +65,7 @@ Configura **`TTS_VOCODER_CMD`** o **`TTS_AMBESERVER_HOST`** / **`TTS_AMBESERVER_
Al programar anuncios, el servidor puede **esperar** si los slots objetivo están ocupados, **descartar** objetivos si aparece un QSO en vivo a mitad, y solo marcar estado de anuncio **horario** tras una lista de objetivos con éxito — evita pisar tráfico en vivo.
Los anuncios programados y TTS inyectan cada trama como **PTT sintético de hotspot** en `PROXY.TARGET_SYSTEM` (el mismo MASTER que usa el proxy integrado). El enrutamiento reparte a hotspots puenteados y piernas OPENBRIDGE; los peers locales en ese MASTER siguen recibiendo tramas vía `send_system`.
El plugin envía cada trama con `send_dmrd`, y el núcleo la inyecta 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`. Mientras suena, la emisión ocupa ese slot del MASTER; una radio que lo toma la corta.
## Cola de emisión

@ -0,0 +1,23 @@
# voice-announcements
Scheduled announcements, TTS announcements and voice beacons on talkgroups. This used to be part of the core (`VoiceUseCases`) and now is a plugin that speaks through `ServerContext.send_dmrd`.
**Nothing changes for sysops.** The items are still configured in `adn-voice.yaml` (`VOICE.ANNOUNCEMENTS`, `VOICE.TTS_ANNOUNCEMENTS`, `AUDIO_PATH`, `TTS_*`), with the same keys and behaviour, and changes are picked up within 15 s. The plugin is enabled by default. The server grants it exactly the talkgroups and DMR IDs those items use, unless `PLUGINS.send.voice-announcements` says otherwise.
## Behaviour
It is the same as the core announcements it replaces:
- **Schedule:** `MODE: interval` plays every `INTERVAL` s. `MODE: hourly` plays once in the first minute of each hour.
- **Talkgroups:** one playback per talkgroup at a time; others wait their turn, 1.5 s apart.
- **Slot:** the slot the server chooses (`voice_slot_for_tg`) is where the TG is a dynamic session or has an active bridge on the MASTER, then TS2, then TS1. While every slot is busy, it retries every 5 s (up to 60 times).
- **Routing:** like any announcement, through the TG's bridges (OpenBridge included) and to the MASTER's hotspots.
- **Pacing and collisions:** 58 ms per frame. It stops when a radio takes the slot.
- **TTS:** `.txt` → speech → `.ambe`, encoded in a worker thread and cached, as before.
## Tests
```bash
cd plugins/voice-announcements
PYTHONPATH=../../src python3 -m pytest tests/ -q
```

@ -0,0 +1,7 @@
# Scheduled announcements, TTS and voice beacons (moved out of the core).
# The items themselves stay where they were: VOICE.ANNOUNCEMENTS and
# VOICE.TTS_ANNOUNCEMENTS in adn-voice.yaml (re-read every 15 s).
enabled: true
# How often to follow changes in adn-voice.yaml, in seconds.
reconcile_s: 15

@ -1,12 +1,8 @@
# ADN DMR Peer Server - TTS to AMBE stub
# ADN DMR Peer Server plugin - voice-announcements factory
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
# 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
@ -22,16 +18,12 @@
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""TTS to AMBE file (legacy tts_engine.ensure_tts_ambe). Stub until ported."""
"""voice-announcements plugin — drop-in factory."""
from __future__ import annotations
from typing import Any
from .application.plugin_impl import VoiceAnnouncementsPlugin
def ensure_tts_ambe(text: str, lang: str, out_path: str, config: dict[str, Any]) -> str | None:
"""
Generate TTS audio and convert to AMBE file; return output path or None.
Legacy: tts_engine.ensure_tts_ambe. Full implementation to be ported when needed.
"""
return None
def create_plugin() -> VoiceAnnouncementsPlugin:
return VoiceAnnouncementsPlugin()

@ -0,0 +1,210 @@
# ADN DMR Peer Server plugin - voice-announcements announcer
#
# 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
###############################################################################
"""Plays scheduled announcements and TTS on their talkgroups through ``send_dmrd``.
Everything runs on the reactor (``call_later``); only TTS encoding goes to a thread.
Behaviour follows the core announcements it replaces: one playback per talkgroup at a
time, retries while the talkgroup or every slot is busy, 58 ms per frame, and the
playback stops as soon as the server refuses a frame (a radio took the slot).
"""
from __future__ import annotations
import logging
import os
import time
from datetime import datetime
from typing import Any, Callable
from adn_server.domain import bytes_3
from adn_server.domain.mesh_engine import server_id_bytes
from ..domain.schedule import FILE, Item, enabled_items, hourly_due
from ..domain.tg_queue import TalkgroupQueue
from ..infrastructure.tts_engine import ensure_tts_ambe
logger = logging.getLogger(__name__)
FRAME_INTERVAL_S = 0.058
START_DELAY_S = 0.5
GAP_BETWEEN_PLAYBACKS_S = 1.5
TG_BUSY_RETRY_S = 3.0
SLOT_BUSY_RETRY_S = 5.0
MAX_RETRIES = 60
class Announcer:
def __init__(
self,
ctx: Any,
voice: Any,
tts: Callable[[dict[str, Any], dict[str, Any], str], str | None] | None = None,
now: Callable[[], datetime] = datetime.now,
) -> None:
self._ctx = ctx
self._voice = voice
self._tts = tts if tts is not None else ensure_tts_ambe
self._now = now
self._items: dict[tuple[str, int], Item] = {}
self._timers: dict[tuple[str, int], Any] = {}
self._running: set[tuple[str, int]] = set()
self._last_hour: dict[tuple[str, int], int] = {}
self._queue: TalkgroupQueue[tuple[Item, list]] = TalkgroupQueue()
self._stopped = False
# --- schedule -------------------------------------------------------------
def reconcile(self) -> None:
"""Start, restart or stop item timers to match the current ``VOICE`` config."""
wanted = {item.key: item for item in enabled_items(self._ctx.config)}
for key in list(self._timers):
if key not in wanted or wanted[key] != self._items.get(key):
self._cancel(key)
logger.info("(VOICE-RELOAD) %s stopped", self._items[key].label)
del self._items[key]
for key, item in wanted.items():
if key in self._timers:
continue
self._items[key] = item
self._timers[key] = self._ctx.call_later(item.tick_s, self._tick, key)
logger.info(
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s", item.label, item.mode, item.file, item.tg
)
def shutdown(self) -> None:
self._stopped = True
for key in list(self._timers):
self._cancel(key)
def _cancel(self, key: tuple[str, int]) -> None:
timer = self._timers.pop(key, None)
try:
if timer is not None and timer.active():
timer.cancel()
except Exception: # already fired
pass
def _tick(self, key: tuple[str, int]) -> None:
item = self._items.get(key)
if item is None or self._stopped:
return
self._timers[key] = self._ctx.call_later(item.tick_s, self._tick, key)
self.fire(item)
# --- one playback ---------------------------------------------------------
def fire(self, item: Item, retry: int = 0) -> None:
if item.key in self._running:
if retry == 0:
logger.debug("(%s) Previous playback still running, skipping", item.label)
return
if item.mode == "hourly" and retry == 0 and not hourly_due(self._now(), self._last_hour.get(item.key)):
return
if self._queue.on_air(item.tg) and retry < MAX_RETRIES:
if retry == 0:
logger.debug("(%s) Same TG %s already broadcasting, deferring", item.label, item.tg)
self._ctx.call_later(TG_BUSY_RETRY_S + item.index * 0.5, self.fire, item, retry + 1)
return
self._running.add(item.key)
if item.mode == "hourly":
self._last_hour[item.key] = self._now().hour
if item.kind == FILE:
self._load(item)
return
logger.info("(%s) Starting TTS conversion in background thread for %s", item.label, item.file)
d = self._ctx.defer_to_thread(self._tts, self._ctx.config, item.raw, self._audio_path())
d.addCallback(self._tts_ready, item)
d.addErrback(self._tts_failed, item)
def _tts_ready(self, ambe_path: str | None, item: Item) -> None:
if not ambe_path:
logger.warning("(%s) No AMBE file available for TTS announcement %s", item.label, item.file)
self._running.discard(item.key)
return
self._load(item)
def _tts_failed(self, failure: Any, item: Item) -> None:
self._running.discard(item.key)
message = failure.getErrorMessage() if hasattr(failure, "getErrorMessage") else str(failure)
logger.error("(%s) TTS conversion error: %s", item.label, message)
def _load(self, item: Item) -> None:
name = item.file[:-5] if item.file.endswith(".ambe") else item.file
try:
ambe = self._voice.read_single_file(self._audio_path(), item.language, name)
except Exception as e:
ambe = None
logger.warning("(%s) Cannot read AMBE file %s/ondemand/%s.ambe: %s", item.label, item.language, name, e)
if not ambe:
logger.warning("(%s) AMBE file empty or not found: %s/ondemand/%s.ambe", item.label, item.language, name)
self._running.discard(item.key)
return
if self._queue.request(item.tg, (item, ambe)):
self._begin(item, ambe)
else:
logger.info("(%s) Same TG %s still on air; deferring next playback (pending %s)", item.label, item.tg, len(self._queue))
def _begin(self, item: Item, ambe: list, retry: int = 0) -> None:
"""Holds the TG in the queue while it waits for a free slot."""
if self._stopped:
return
slot_for_tg = getattr(self._ctx, "voice_slot_for_tg", None)
slot = slot_for_tg(item.tg) if slot_for_tg is not None else None
if slot is None:
if slot_for_tg is not None and retry < MAX_RETRIES:
if retry == 0:
logger.info("(%s) All target slots busy (QSO active), waiting for QSO to finish...", item.label)
self._ctx.call_later(SLOT_BUSY_RETRY_S, self._begin, item, ambe, retry + 1)
return
logger.warning("(%s) No slot to speak TG %s on; not played", item.label, item.tg)
self._finish(item)
return
pkts = list(self._voice.pkt_gen(bytes_3(item.dmr_id), bytes_3(item.tg), self._server_id(), 1 if slot == 2 else 0, [ambe]))
logger.info("(%s) Playing %s to TG %s on TS%s (%s packets)", item.label, item.file, item.tg, slot, len(pkts))
self._ctx.call_later(START_DELAY_S, self._play, item, pkts, 0, None)
def _play(self, item: Item, pkts: list[bytes], idx: int, next_time: float | None) -> None:
if self._stopped:
return
if idx >= len(pkts):
logger.info("(%s) Broadcast complete: %s packets", item.label, len(pkts))
self._finish(item)
return
if not self._ctx.send_dmrd(pkts[idx]):
logger.info("(%s) Broadcast stopped at packet %s/%s: slot taken (QSO)", item.label, idx, len(pkts))
self._finish(item)
return
next_time = (time.time() if next_time is None else next_time) + FRAME_INTERVAL_S
self._ctx.call_later(max(0.001, next_time - time.time()), self._play, item, pkts, idx + 1, next_time)
def _finish(self, item: Item) -> None:
self._running.discard(item.key)
nxt = self._queue.finished(item.tg)
if nxt is not None:
self._ctx.call_later(GAP_BETWEEN_PLAYBACKS_S, self._begin, *nxt)
# --- helpers ----------------------------------------------------------------
def _audio_path(self) -> str:
return os.path.join(self._ctx.project_root, (self._ctx.config.get("VOICE") or {}).get("AUDIO_PATH", "Audio"))
def _server_id(self) -> bytes:
return server_id_bytes(self._ctx.config.get("GLOBAL", {}).get("SERVER_ID"))[:4]

@ -0,0 +1,80 @@
# ADN DMR Peer Server plugin - voice-announcements adapter
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""ServerPlugin adapter: scheduled announcements, TTS and voice beacons."""
from __future__ import annotations
import logging
from typing import Any
from adn_server.infrastructure.voice import DefaultVoiceProvider
from .announcer import Announcer
logger = logging.getLogger(__name__)
class VoiceAnnouncementsPlugin:
name = "voice-announcements"
events = () # listens to nothing: it only speaks
def __init__(self) -> None:
self._announcer: Announcer | None = None
self._ctx: Any = None
self._reconcile_s = 15.0
self._timer: Any = None
def on_load(self, bus: Any, config: dict[str, Any], server_ctx: Any) -> None:
del bus
self._ctx = server_ctx
self._reconcile_s = float(config.get("reconcile_s", 15))
if server_ctx.send_dmrd is None or getattr(server_ctx, "voice_slot_for_tg", None) is None:
logger.warning(
"(VOICE) %s is not allowed to send group voice: add it to PLUGINS.send with "
"allowed_src_ids and group_voice_tgs; announcements will not play",
self.name,
)
return
self._announcer = Announcer(server_ctx, DefaultVoiceProvider())
self._reconcile()
def _reconcile(self) -> None:
if self._announcer is None:
return
self._announcer.reconcile()
# The core re-reads adn-voice.yaml every 15 s into config["VOICE"]; follow it.
self._timer = self._ctx.call_later(self._reconcile_s, self._reconcile)
def on_reload(self, config: dict[str, Any]) -> None:
self._reconcile_s = float(config.get("reconcile_s", self._reconcile_s))
def on_shutdown(self) -> None:
if self._announcer is not None:
self._announcer.shutdown()
self._announcer = None
try:
if self._timer is not None and self._timer.active():
self._timer.cancel()
except Exception:
pass
def on_event(self, event: object) -> None:
pass

@ -0,0 +1,101 @@
# ADN DMR Peer Server plugin - voice-announcements schedule
#
# 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
###############################################################################
"""Scheduled announcements and TTS items, read from the server's ``VOICE`` section.
The same keys the core used (adn-voice.yaml), so existing installs keep working.
"""
from __future__ import annotations
from dataclasses import dataclass
from datetime import datetime
from typing import Any
from adn_server.application.server_voice import announcement_item_dmr_id
FILE = "file"
TTS = "tts"
@dataclass(frozen=True)
class Item:
kind: str # FILE (pre-recorded .ambe) or TTS (.txt spoken, cached as .ambe)
index: int
file: str
tg: int
language: str
mode: str # "interval" or "hourly"
interval_s: float
dmr_id: int
raw: dict[str, Any] # the item as configured: TTS settings live here too
@property
def label(self) -> str:
return f"{'ANNOUNCEMENT' if self.kind == FILE else 'TTS'}-{self.index + 1}"
@property
def key(self) -> tuple[str, int]:
return (self.kind, self.index)
@property
def tick_s(self) -> float:
"""How often the item's timer fires: hourly items check the clock every 30 s."""
return 30.0 if self.mode == "hourly" else self.interval_s
def enabled_items(config: dict[str, Any]) -> list[Item]:
"""Enabled, playable items of ``VOICE.ANNOUNCEMENTS`` and ``VOICE.TTS_ANNOUNCEMENTS``."""
voice = config.get("VOICE") or {}
items: list[Item] = []
for kind, section in ((FILE, "ANNOUNCEMENTS"), (TTS, "TTS_ANNOUNCEMENTS")):
entries = voice.get(section) or []
if not isinstance(entries, list):
continue
for index, raw in enumerate(entries):
if not isinstance(raw, dict) or not raw.get("ENABLED"):
continue
try:
tg = int(raw.get("TG", 0))
interval = float(raw.get("INTERVAL", 60))
except (TypeError, ValueError):
continue
file = str(raw.get("FILE") or "").strip()
if not file or not tg or interval <= 0:
continue
items.append(
Item(
kind=kind,
index=index,
file=file,
tg=tg,
language=str(raw.get("LANGUAGE", "en_GB")),
mode=str(raw.get("MODE", "interval")),
interval_s=interval,
dmr_id=announcement_item_dmr_id(raw, config),
raw=dict(raw),
)
)
return items
def hourly_due(now: datetime, last_hour: int | None) -> bool:
"""Hourly items play once, in the first minute of each hour."""
return now.minute == 0 and last_hour != now.hour

@ -0,0 +1,57 @@
# ADN DMR Peer Server plugin - voice-announcements talkgroup queue
#
# 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
###############################################################################
"""One playback at a time per talkgroup; the rest wait their turn, in order."""
from __future__ import annotations
from typing import Generic, TypeVar
T = TypeVar("T")
class TalkgroupQueue(Generic[T]):
def __init__(self) -> None:
self._on_air: set[int] = set()
self._waiting: list[tuple[int, T]] = []
def on_air(self, tg: int) -> bool:
return tg in self._on_air
def request(self, tg: int, job: T) -> bool:
"""True: start ``job`` now (its TG is free). False: queued behind the one on air."""
if tg in self._on_air:
self._waiting.append((tg, job))
return False
self._on_air.add(tg)
return True
def finished(self, tg: int) -> T | None:
"""Free ``tg``; the next waiting job whose TG is free, now marked on air."""
self._on_air.discard(tg)
for i, (waiting_tg, job) in enumerate(self._waiting):
if waiting_tg not in self._on_air:
del self._waiting[i]
self._on_air.add(waiting_tg)
return job
return None
def __len__(self) -> int:
return len(self._waiting)

@ -0,0 +1,28 @@
# ADN DMR Peer Server plugin - example test path
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
from __future__ import annotations
import sys
from pathlib import Path
_ROOT = Path(__file__).resolve().parents[1]
if str(_ROOT) not in sys.path:
sys.path.insert(0, str(_ROOT))

@ -0,0 +1,258 @@
# ADN DMR Peer Server plugin - voice-announcements tests
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Scheduling, per-TG queue, busy slots, QSO cut, hourly and TTS — on a simulated reactor."""
from __future__ import annotations
import heapq
import itertools
from datetime import datetime
from types import SimpleNamespace
import pytest
from plugin.application import announcer as ann
from plugin.application.announcer import Announcer
from plugin.application.plugin_impl import VoiceAnnouncementsPlugin
class Reactor:
"""call_later on a fake clock; the announcer's time.time follows it."""
def __init__(self) -> None:
self.now = 1000.0
self._heap: list = []
self._seq = itertools.count()
def call_later(self, delay, fn, *args):
handle = SimpleNamespace(cancelled=False, called=False)
handle.active = lambda h=handle: not (h.cancelled or h.called)
handle.cancel = lambda h=handle: setattr(h, "cancelled", True)
heapq.heappush(self._heap, (self.now + delay, next(self._seq), handle, fn, args))
return handle
def advance(self, seconds: float) -> None:
end = self.now + seconds
while self._heap and self._heap[0][0] <= end:
when, _, handle, fn, args = heapq.heappop(self._heap)
self.now = when
if not handle.cancelled:
handle.called = True
fn(*args)
self.now = end
class Voice:
frames = 5
def read_single_file(self, audio_path, lang, name):
return [["ambe"]] if name != "missing" else []
def pkt_gen(self, src, dst, peer, slot_bit, phrase):
return [b"DMRD" + bytes([n]) + src + dst + peer + bytes([slot_bit << 7]) for n in range(self.frames)]
def ensure_tts_ambe(self, config, item, audio_path):
return "/tmp/x.ambe"
class Deferred:
def __init__(self, result):
self.result = result
def addCallback(self, fn, *a):
if not isinstance(self.result, Exception):
fn(self.result, *a)
return self
def addErrback(self, fn, *a):
if isinstance(self.result, Exception):
fn(self.result, *a)
return self
def _ctx(reactor, voice_cfg, *, slots=None, refuse_at=None):
sent: list[bytes] = []
slots = slots if slots is not None else [2]
def send(pkt):
if refuse_at is not None and len(sent) == refuse_at:
return False
sent.append(pkt)
return True
def slot_for_tg(tg):
return slots.pop(0) if len(slots) > 1 else slots[0]
ctx = SimpleNamespace(
config={"GLOBAL": {"SERVER_ID": 2131}, "VOICE": voice_cfg},
project_root="/tmp", call_later=reactor.call_later, send_dmrd=send, voice_slot_for_tg=slot_for_tg,
defer_to_thread=lambda fn, *a: Deferred(fn(*a)),
)
return ctx, sent
@pytest.fixture
def reactor(monkeypatch):
r = Reactor()
monkeypatch.setattr(ann.time, "time", lambda: r.now)
return r
def _ann(file="beacon", tg=213, interval=60, **kw):
return {"ENABLED": True, "FILE": file, "TG": tg, "MODE": "interval", "INTERVAL": interval, "LANGUAGE": "es_ES", **kw}
def test_an_interval_announcement_plays_every_interval(reactor) -> None:
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann(DMR_ID=2130035)]})
a = Announcer(ctx, Voice())
a.reconcile()
reactor.advance(59)
assert sent == []
reactor.advance(2)
assert len(sent) == Voice.frames
assert sent[0][5:8] == (2130035).to_bytes(3, "big") and sent[0][8:11] == (213).to_bytes(3, "big")
assert sent[0][-1] == 0x80 # TS2, as voice_slot_for_tg said
reactor.advance(60)
assert len(sent) == 2 * Voice.frames
def test_frames_are_paced_58_ms_apart(reactor) -> None:
ctx, _ = _ctx(reactor, {"ANNOUNCEMENTS": [_ann()]})
times: list[float] = []
send = ctx.send_dmrd
ctx.send_dmrd = lambda pkt: times.append(reactor.now) or send(pkt)
Announcer(ctx, Voice()).reconcile()
reactor.advance(61)
gaps = [round(b - a, 3) for a, b in zip(times, times[1:])]
assert gaps == [0.058] * (Voice.frames - 1)
def test_a_refused_frame_stops_the_playback(reactor) -> None:
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann()]}, refuse_at=2)
a = Announcer(ctx, Voice())
a.reconcile()
reactor.advance(61)
assert len(sent) == 2
assert not a._running # free to play at the next interval
def test_two_items_on_one_tg_take_turns(reactor) -> None:
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann("one"), _ann("two")]})
Announcer(ctx, Voice()).reconcile()
reactor.advance(60.6) # both fired at 60 s; the first is playing
assert 0 < len(sent) <= Voice.frames
reactor.advance(10)
assert len(sent) == 2 * Voice.frames
def test_busy_slots_are_retried_every_5_s(reactor) -> None:
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann()]}, slots=[None, None, 1])
Announcer(ctx, Voice()).reconcile()
reactor.advance(61)
assert sent == []
reactor.advance(10.6)
assert len(sent) == Voice.frames and sent[0][-1] == 0x00 # TS1 once free
def test_hourly_plays_once_in_the_first_minute_of_the_hour(reactor) -> None:
clock = {"t": datetime(2026, 9, 24, 21, 59, 10)}
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann(MODE="hourly")]})
Announcer(ctx, Voice(), now=lambda: clock["t"]).reconcile()
reactor.advance(31)
assert sent == []
clock["t"] = datetime(2026, 9, 24, 22, 0, 5)
reactor.advance(31)
assert len(sent) == Voice.frames
clock["t"] = datetime(2026, 9, 24, 22, 0, 40)
reactor.advance(31)
assert len(sent) == Voice.frames # not twice in the same hour
def test_tts_is_encoded_in_a_thread_then_played(reactor) -> None:
ctx, sent = _ctx(reactor, {"TTS_ANNOUNCEMENTS": [_ann("texto1")]})
calls: list = []
real = ctx.defer_to_thread
ctx.defer_to_thread = lambda fn, *a: calls.append(fn.__name__) or real(fn, *a)
Announcer(ctx, Voice(), tts=Voice().ensure_tts_ambe).reconcile()
reactor.advance(61)
assert calls == ["ensure_tts_ambe"] and len(sent) == Voice.frames
def test_a_config_change_restarts_and_disabling_stops(reactor) -> None:
voice = {"ANNOUNCEMENTS": [_ann(interval=60)]}
ctx, sent = _ctx(reactor, voice)
a = Announcer(ctx, Voice())
a.reconcile()
voice["ANNOUNCEMENTS"][0] = _ann(interval=10)
a.reconcile()
reactor.advance(11)
assert len(sent) == Voice.frames
voice["ANNOUNCEMENTS"][0]["ENABLED"] = False
a.reconcile()
reactor.advance(100)
assert len(sent) == Voice.frames
def test_a_missing_file_is_skipped(reactor) -> None:
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann("missing")]})
a = Announcer(ctx, Voice())
a.reconcile()
reactor.advance(61)
assert sent == [] and not a._running
def test_without_a_send_grant_the_plugin_does_nothing(reactor) -> None:
ctx, _ = _ctx(reactor, {"ANNOUNCEMENTS": [_ann()]})
ctx.send_dmrd = None
plugin = VoiceAnnouncementsPlugin()
plugin.on_load(None, {}, ctx)
assert plugin._announcer is None
@pytest.mark.parametrize("result", [None, RuntimeError("gTTS down")])
def test_a_failed_tts_conversion_frees_the_item(reactor, result) -> None:
ctx, sent = _ctx(reactor, {"TTS_ANNOUNCEMENTS": [_ann("texto1")]})
ctx.defer_to_thread = lambda fn, *a: Deferred(result)
a = Announcer(ctx, Voice(), tts=lambda *a: None)
a.reconcile()
reactor.advance(61)
assert sent == [] and not a._running
def test_an_item_without_dmr_id_speaks_as_the_server_voice_id(reactor) -> None:
from adn_server.application.server_voice import server_voice_dmr_id
ctx, sent = _ctx(reactor, {"ANNOUNCEMENTS": [_ann()]})
Announcer(ctx, Voice()).reconcile()
reactor.advance(61)
assert int.from_bytes(sent[0][5:8], "big") == server_voice_dmr_id(ctx.config)
def test_the_server_grants_the_plugin_what_voice_configures() -> None:
from adn_server.application.plugins.domain.send import send_permission
config = {"VOICE": {"ANNOUNCEMENTS": [_ann(tg=213, DMR_ID=2130035)], "TTS_ANNOUNCEMENTS": [_ann(tg=9, ENABLED=False)]}}
permission = send_permission(config, "voice-announcements")
assert permission.group_voice_tgs == {213, 9} # disabled items too: enabling needs no restart
assert 2130035 in permission.allowed_src_ids
config["PLUGINS"] = {"send": {"voice-announcements": {"allowed_src_ids": [1], "group_voice_tgs": [9]}}}
assert send_permission(config, "voice-announcements").group_voice_tgs == {9} # an explicit entry wins
config["PLUGINS"] = {"master_kill": True}
assert send_permission(config, "voice-announcements") is None

@ -27,8 +27,8 @@ import threading
import time
import wave
from adn_server.infrastructure.voice import tts_engine
from adn_server.infrastructure.voice.tts_engine import DV3K_PRODID_REQ, DV3K_SAMPLES_PER_FRAME
from plugin.infrastructure import tts_engine
from plugin.infrastructure.tts_engine import DV3K_PRODID_REQ, DV3K_SAMPLES_PER_FRAME
def test_text_to_ambe_serializes_parallel_conversions(tmp_path, monkeypatch) -> None:

@ -35,3 +35,5 @@ class ServerContext:
call_later: Callable[..., Any]
# Present only for a plugin listed in PLUGINS.send (see domain/send.py).
send_dmrd: Callable[[bytes], bool] | None = None
# With group voice granted: the MASTER slot to speak a TG on now, None while all are busy.
voice_slot_for_tg: Callable[[int], int | None] | None = None

@ -26,14 +26,19 @@ import logging
import time
from typing import Any, Callable
from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM
from ....domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, int_id
from ....domain.dmr import decode
from ....domain.dmr.const import LC_OPT
from ....domain.hbp_protocol import STREAM_TO
from ....domain.mesh_engine import server_id_bytes
from ...routing.announcement_ptt_inject import announcement_ptt_system, inject_plugin_dmrd
from ...routing.helpers import slot_voice_held_by_other_stream
from ...routing.helpers import master_dynamic_tg_slots, slot_voice_held_by_other_stream
from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_frame_kind
logger = logging.getLogger(__name__)
_SILENT_STREAM_S = 1.0
class PluginIngress:
"""Routes plugin frames on the reactor thread and reports whether each was accepted.
@ -42,9 +47,11 @@ class PluginIngress:
Group voice is routed like a scheduled announcement: through the bridges (OpenBridge
legs included, as for any local ingress), and out to the hotspots of the MASTER it
enters on. While a plugin's stream plays it holds that MASTER slot (TX_TYPE=VHEAD,
TX_STREAM_ID, TX_RFS), so routed voice finds it busy; a radio or another stream on the
slot makes the frame fail, which tells the plugin to stop. The terminator frees it.
enters on. Each accepted frame records the slot's RX state exactly as ``udp_hbp`` does
for a hotspot's frame (RX_STREAM_ID, RX_LC, RX_RFS, RX_TGID, RX_TYPE, RX_TIME…): routing
then knows the stream it is continuing and the LC it carries, and routed voice finds
the slot busy. A radio or another stream on the slot makes the frame fail, which tells
the plugin to stop. The terminator frees the slot.
"""
def __init__(
@ -54,17 +61,61 @@ class PluginIngress:
get_protocols: Callable[[], dict[str, Any]],
send_local: Callable[[str, bytes], None],
clock: Callable[[], float] = time.time,
send_routing_event: Callable[[str], None] | None = None,
) -> None:
self._routing = routing
self._config = config
self._get_protocols = get_protocols
self._send_local = send_local
self._clock = clock
self._send_routing_event = send_routing_event
# Plugin voice streams on air: stream_id -> [master, slot, tg, rf_src, start, last frame].
self._on_air: dict[bytes, list[Any]] = {}
def set_routing(self, routing: Any) -> None:
"""Bootstrap builds the plugin manager before routing; routing is set before any load."""
self._routing = routing
def voice_slot_for_tg(self, tg: int) -> int | None:
"""The MASTER slot a plugin should speak ``tg`` on now, or None while every slot is busy.
Same choice as scheduled announcements: the slot where ``tg`` is a dynamic UA
session or has an active bridge leg on that MASTER first, then TS2, then TS1.
"""
master = announcement_ptt_system(self._config)
proto = self._get_protocols().get(master) if master else None
status = getattr(proto, "STATUS", None)
sys_cfg = self._config.get("SYSTEMS", {}).get(master or "", {})
if not status or sys_cfg.get("MODE") != "MASTER":
return None
now = self._clock()
self._end_silent_streams(now)
dynamic = master_dynamic_tg_slots(sys_cfg, int(tg))
bridged = self._active_bridge_slots(int(tg), master)
for ts in dict.fromkeys([*sorted(dynamic, reverse=True), *sorted(bridged, reverse=True), 2, 1]):
slot = status.get(ts)
if slot and not self._slot_busy(slot, ts in dynamic or ts in bridged, now):
return ts
return None
@staticmethod
def _slot_busy(slot: dict[str, Any], tg_lives_here: bool, now: float) -> bool:
if slot_voice_held_by_other_stream(slot, b"", now):
return True # live voice on the slot, in or out
if slot.get("RX_TYPE") != HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM:
return True # an outside QSO still in its hang time
if tg_lives_here:
return False
return not (slot.get("RX_TYPE") == HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM)
def _active_bridge_slots(self, tg: int, master: str) -> set[int]:
table = self._routing.routing_table_for_report() if self._routing is not None else {}
return {
int(be["TS"])
for be in table.get(str(tg), [])
if isinstance(be, dict) and be.get("SYSTEM") == master and be.get("ACTIVE") and be.get("TS") is not None
}
def deliver(self, pkt: bytes, plugin: str) -> bool:
header = parse_dmrd_header(pkt)
kind = plugin_frame_kind(header.call_type, header.frame_type, header.dtype_vseq) if header else None
@ -75,6 +126,7 @@ class PluginIngress:
return False
server_id = self._server_id()
now = self._clock()
self._end_silent_streams(now)
if kind == UNIT_DATA:
return self._route(master, pkt, now, server_id, plugin) is not False
return self._group_voice(master, header, pkt[:11] + server_id + pkt[15:], now, server_id, plugin)
@ -84,27 +136,95 @@ class PluginIngress:
slot = getattr(proto, "STATUS", {}).get(header.slot) if proto is not None else None
if slot is None:
return False
if slot_voice_held_by_other_stream(slot, header.stream_id, now):
return False
if self._route(master, pkt, now, server_id, plugin) is not True:
if _slot_taken(slot, header.stream_id, now) or self._route(master, pkt, now, server_id, plugin) is not True:
self._release(slot, header.stream_id)
self._off_air(header.stream_id, now)
return False
is_term = header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VTERM
slot["TX_TYPE"] = HBPF_SLT_VTERM if is_term else HBPF_SLT_VHEAD
slot["TX_STREAM_ID"] = header.stream_id
slot["TX_RFS"] = header.rf_src
slot["TX_TGID"] = header.dst_id
slot["TX_TIME"] = now
_record_rx(slot, header, pkt, server_id, now)
self._send_local(master, pkt)
if header.stream_id not in self._on_air:
self._on_air_start(master, header, now)
self._on_air[header.stream_id][5] = now
if is_term:
self._off_air(header.stream_id, now)
return True
def _on_air_start(self, master: str, header: Any, now: float) -> None:
"""Monitor TX on the MASTER itself: its hotspots hear the stream, which no bridge leg reports."""
tg, rf_src = int_id(header.dst_id), int_id(header.rf_src)
self._on_air[header.stream_id] = [master, header.slot, tg, rf_src, now, now]
self._report("START", master, header.stream_id, header.slot, tg, rf_src)
def _end_silent_streams(self, now: float) -> None:
"""A plugin stream that stopped without a terminator (the plugin gave up, a frame
was refused before routing) is over once it has been silent for a second."""
for stream_id, entry in list(self._on_air.items()):
if now - entry[5] > _SILENT_STREAM_S:
proto = self._get_protocols().get(entry[0])
slot = getattr(proto, "STATUS", {}).get(entry[1]) if proto is not None else None
if slot is not None:
self._release(slot, stream_id)
self._off_air(stream_id, entry[5])
def _off_air(self, stream_id: bytes, now: float) -> None:
entry = self._on_air.pop(stream_id, None)
if entry is not None:
master, slot, tg, rf_src, start, _last = entry
self._report("END", master, stream_id, slot, tg, rf_src, now - start)
def _report(
self, action: str, master: str, stream_id: bytes, slot: int, tg: int, rf_src: int, duration: float | None = None
) -> None:
# Same line announcements send (VoiceUseCases._emit_announcement_voice_event).
if self._send_routing_event is None:
return
parts = ["GROUP VOICE", action, "TX", master, 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}")
parts.append("1")
self._send_routing_event(",".join(parts))
def _route(self, master: str, pkt: bytes, now: float, server_id: bytes, plugin: str) -> bool | None:
return inject_plugin_dmrd(self._routing, master, pkt, pkt_time=now, server_id=server_id, plugin=plugin)
@staticmethod
def _release(slot: dict[str, Any], stream_id: bytes) -> None:
if slot.get("TX_STREAM_ID") == stream_id:
slot["TX_TYPE"] = HBPF_SLT_VTERM
if slot.get("RX_STREAM_ID") == stream_id:
slot["RX_TYPE"] = HBPF_SLT_VTERM
def _server_id(self) -> bytes:
return server_id_bytes(self._config.get("GLOBAL", {}).get("SERVER_ID"))[:4]
def _slot_taken(slot: dict[str, Any], stream_id: bytes, now: float) -> bool:
"""Live voice on the slot that is not this stream: a radio, or another stream."""
for leg in ("RX", "TX"):
leg_type = slot.get(f"{leg}_TYPE")
if leg_type is None or leg_type == HBPF_SLT_VTERM:
continue
if slot.get(f"{leg}_STREAM_ID") == stream_id:
continue
if now - float(slot.get(f"{leg}_TIME", 0) or 0) < STREAM_TO:
return True
return False
def _record_rx(slot: dict[str, Any], header: Any, pkt: bytes, peer_id: bytes, now: float) -> None:
"""The slot RX state ``udp_hbp`` records after routing accepts a hotspot's frame."""
if header.stream_id != slot.get("RX_STREAM_ID"):
slot["RX_START"] = now
lc = LC_OPT + header.dst_id + header.rf_src
if header.frame_type == HBPF_DATA_SYNC and header.dtype_vseq == HBPF_SLT_VHEAD:
try:
lc = decode.voice_head_term(pkt[20:53])["LC"]
except Exception:
pass
slot["RX_LC"] = lc
slot["RX_PEER"] = peer_id
slot["RX_SEQ"] = header.seq
slot["RX_RFS"] = header.rf_src
slot["RX_TYPE"] = header.dtype_vseq
slot["RX_TGID"] = header.dst_id
slot["RX_TIME"] = now
slot["RX_STREAM_ID"] = header.stream_id

@ -136,10 +136,16 @@ class PluginManager:
The sender re-reads the permission on every frame, so revoking it on SIGHUP
is immediate; granting it to an already loaded plugin needs that plugin reloaded.
"""
if self._sender_factory is None or send_permission(server_config, name) is None:
permission = send_permission(server_config, name) if self._sender_factory is not None else None
if permission is None:
return self._ctx
logger.info("(PLUGIN-MANAGER) %s may send unit data (PLUGINS.send)", name)
return replace(self._ctx, send_dmrd=self._sender_factory(name))
sender = self._sender_factory(name)
logger.info(
"(PLUGIN-MANAGER) %s may send unit data%s (PLUGINS.send)",
name, f" and group voice on TG {sorted(permission.group_voice_tgs)}" if permission.group_voice_tgs else "",
)
# voice_slot_for_tg checks the talkgroup grant live, like send_dmrd.
return replace(self._ctx, send_dmrd=sender, voice_slot_for_tg=getattr(sender, "voice_slot_for_tg", None))
def shutdown_all(self) -> None:
for name in list(self._loaded):

@ -54,6 +54,7 @@ class PluginDmrdSender:
call_from_reactor: Callable[..., Any],
clock: Callable[[], float] = time.monotonic,
in_reactor_thread: Callable[[], bool] = lambda: False,
slot_for_tg: Callable[[int], int | None] | None = None,
) -> None:
self._plugin = plugin
self._config = server_config
@ -61,6 +62,7 @@ class PluginDmrdSender:
self._call_from_reactor = call_from_reactor
self._clock = clock
self._in_reactor_thread = in_reactor_thread
self._slot_for_tg = slot_for_tg
self._lock = threading.Lock()
self._tokens = math.inf # starts full: clamped to the rate on first use
self._refilled = clock()
@ -101,6 +103,13 @@ class PluginDmrdSender:
self._call_from_reactor(self._deliver, bytes(pkt), self._plugin)
return True
def voice_slot_for_tg(self, tg: int) -> int | None:
"""``ServerContext.voice_slot_for_tg``: only for granted talkgroups; reactor thread."""
permission = send_permission(self._config, self._plugin)
if permission is None or int(tg) not in permission.group_voice_tgs or self._slot_for_tg is None:
return None
return self._slot_for_tg(int(tg))
def _take_token(self, permission: SendPermission) -> bool:
with self._lock:
now = self._clock()

@ -34,6 +34,8 @@ PLUGIN_SENDABLE_DTYPES = frozenset({3, 6, 7, 8})
UNIT_DATA = "unit_data"
GROUP_VOICE = "group_voice"
DEFAULT_MAX_FRAMES_PER_S = 40.0
# A voice stream is one frame every 60 ms (~17/s), at most one per talkgroup at a time.
VOICE_FRAMES_PER_S_PER_TG = 18.0
class DmrdHeader(NamedTuple):
@ -90,24 +92,56 @@ class SendPermission:
group_voice_tgs: frozenset[int] = frozenset()
ANNOUNCEMENTS_PLUGIN = "voice-announcements"
def announcements_grant(server_config: dict[str, Any]) -> dict[str, Any]:
"""What the official announcements plugin may send, read from ``VOICE``.
The talkgroups and DMR IDs its items use (enabled or not, so enabling one in
adn-voice.yaml needs no restart), plus the server voice ID: announcements kept
working unchanged when they moved out of the core.
"""
from ...server_voice import announcement_item_dmr_id, server_voice_dmr_id
voice = server_config.get("VOICE") or {}
ids = {server_voice_dmr_id(server_config)}
tgs: set[int] = set()
for section in ("ANNOUNCEMENTS", "TTS_ANNOUNCEMENTS"):
for item in voice.get(section) or []:
if not isinstance(item, dict):
continue
ids.add(announcement_item_dmr_id(item, server_config))
try:
if int(item.get("TG", 0)):
tgs.add(int(item["TG"]))
except (TypeError, ValueError):
continue
return {"allowed_src_ids": sorted(ids), "group_voice_tgs": sorted(tgs)}
def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermission | None:
"""The plugin's entry in ``PLUGINS.send``, or None when it may not send.
An entry without source IDs grants nothing. The allowlist and the rate limit
guard against a *buggy* plugin (sending as a radio, flooding the mesh); they are
no sandbox: a plugin runs in-process with the live config and could rewrite its
own entry, so only install plugins you trust.
own entry, so only install plugins you trust. The official announcements plugin
without an entry gets what ``VOICE`` configures (see ``announcements_grant``).
"""
plugins_cfg = server_config.get("PLUGINS") or {}
if plugins_cfg.get("master_kill"):
return None
entry = (plugins_cfg.get("send") or {}).get(plugin)
if entry is None and plugin == ANNOUNCEMENTS_PLUGIN:
entry = announcements_grant(server_config)
if not isinstance(entry, dict):
return None
try:
ids = frozenset(int(i) for i in entry.get("allowed_src_ids") or ())
rate = float(entry.get("max_frames_per_s", DEFAULT_MAX_FRAMES_PER_S))
tgs = frozenset(int(t) for t in entry.get("group_voice_tgs") or ())
default_rate = DEFAULT_MAX_FRAMES_PER_S + VOICE_FRAMES_PER_S_PER_TG * len(tgs)
rate = float(entry.get("max_frames_per_s", default_rate))
except (TypeError, ValueError):
return None
if not ids or not math.isfinite(rate) or rate <= 0:

@ -196,11 +196,6 @@ class VoiceProvider(ABC):
"""Generate HBP voice packets for a phrase (generator). Legacy mk_voice.pkt_gen."""
...
@abstractmethod
def ensure_tts_ambe(self, config: dict[str, Any], item: dict[str, Any], audio_path: str) -> str | None:
"""TTS to AMBE file; return path or None. Legacy tts_engine.ensure_tts_ambe."""
...
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
"""Read one AMBE file (e.g. ondemand/{file_number}.ambe). Legacy readSingleFile for playFileOnRequest."""
return []

@ -39,7 +39,7 @@
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Inject scheduled announcements as a synthetic hotspot PTT on the proxy MASTER."""
"""Plugin frames (announcements, beacons, D-APRS…) as a synthetic ingress on the proxy MASTER."""
from __future__ import annotations
@ -47,7 +47,6 @@ from typing import Any
from ..plugins.domain.send import parse_dmrd_header
from ..proxy.deployment import proxy_target_system
from .helpers import parse_dmrd_burst_fields
def announcement_ptt_system(config: dict[str, Any]) -> str | None:
@ -68,39 +67,6 @@ def announcement_ptt_system(config: dict[str, Any]) -> str | None:
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,
synthetic_announcement=True,
)
def inject_plugin_dmrd(
routing: Any,
master_system: str,
@ -110,10 +76,10 @@ def inject_plugin_dmrd(
server_id: bytes,
plugin: str,
) -> bool | None:
"""Feed one plugin frame through ``dmrd_received`` on the same MASTER as announcements.
"""Feed one plugin frame through ``dmrd_received`` on the announcement MASTER.
The peer field is always the SERVER_ID: a plugin never speaks as a connected peer.
Group voice is marked ``synthetic_announcement`` exactly like announcements are.
Group voice is marked ``synthetic_announcement``, as announcements have always been.
"""
header = parse_dmrd_header(pkt)
if header is None:

@ -22,29 +22,23 @@
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Voice/AMBE/TTS: scheduled announcements, TTS announcements, playback. Orchestrates VoiceProvider."""
"""Server voice prompts (ident, on-demand files, disconnected). Orchestrates VoiceProvider."""
from __future__ import annotations
import logging
import time
from dataclasses import dataclass
from datetime import datetime
from typing import Any, Callable
from ..domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4, int_id
from ..domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3, bytes_4
from .ports import VoiceProvider
from .routing.helpers import slot_voice_held_by_other_stream
from .server_voice import (
announcement_item_source_bytes,
server_voice_rf_src_bytes,
)
from .server_voice import server_voice_rf_src_bytes
logger = logging.getLogger(__name__)
_FRAME_INTERVAL = 0.058
_ANNOUNCEMENT_EXCLUDED = ("ECHO", "D-APRS")
_BROADCAST_GAP = 1.5
@dataclass
@ -60,7 +54,10 @@ class _PromptRun:
class VoiceUseCases:
"""Use cases for voice announcements and TTS."""
"""Server voice prompts: voice ident, on-demand files (TG 9991-9999), disconnected prompt.
Scheduled announcements and TTS live in the voice-announcements plugin.
"""
def __init__(
self,
@ -69,35 +66,12 @@ class VoiceUseCases:
get_protocols: Callable[[], dict[str, Any]] | None = None,
call_from_reactor: Callable[..., None] | None = None,
audio_path: str | None = None,
routing_table_for_report: Callable[[], dict[str, list[dict[str, Any]]]] | None = None,
call_later: Callable[..., Any] | None = None,
start_looping_call: Callable[[Callable[[], None], float, bool], Any] | None = None,
defer_to_thread: Callable[..., Any] | None = None,
inject_announcement_ptt: Callable[[bytes, float], bool | None] | None = None,
send_routing_event: Callable[[str], None] | None = None,
announcement_ptt_system: str | None = None,
) -> None:
self._voice = voice_provider
self._config = config
self._get_protocols = get_protocols
self._call_from_reactor = call_from_reactor
self._audio_path = audio_path or ""
self._routing_table_for_report = routing_table_for_report
self._call_later = call_later
self._start_looping_call = start_looping_call
self._defer_to_thread = defer_to_thread
self._inject_announcement_ptt = inject_announcement_ptt
self._send_routing_event = send_routing_event
self._announcement_ptt_system = announcement_ptt_system
self._voice_report_state: dict[str, Any] | None = None
self._ann_tasks: dict[int, Any] = {}
self._tts_tasks: dict[int, Any] = {}
self._announcement_running: dict[int, bool] = {}
self._tts_running: dict[int, bool] = {}
self._announcement_last_hour: dict[int, int] = {}
self._tts_last_hour: dict[int, int] = {}
self._broadcast_queue: list[dict[str, Any]] = []
self._broadcast_active_tgs: set[str] = set()
def _server_source_id(self) -> bytes:
return server_voice_rf_src_bytes(self._config)
@ -110,795 +84,6 @@ class VoiceUseCases:
"""Generate HBP voice packets for phrase (legacy mk_voice.pkt_gen)."""
return self._voice.pkt_gen(rf_src, dst_id, peer, slot, phrase)
def _global_server_id_bytes(self) -> bytes:
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
return bytes_4(int_id(server_id))
def _active_bridge_slots_for_tg(self, tg: int, system: str) -> set[int]:
bridges = self._routing_table_for_report() if self._routing_table_for_report else {}
entries = bridges.get(str(tg), [])
if not isinstance(entries, list):
return set()
out: set[int] = set()
for be in entries:
if not isinstance(be, dict):
continue
if be.get("SYSTEM") != system or not be.get("ACTIVE"):
continue
ts = be.get("TS")
if ts is not None:
out.add(int(ts))
return out
def _inject_ptt_slot_busy(
self,
slot: dict[str, Any],
tg: int,
sys_cfg: dict[str, Any],
wire_ts: int,
) -> bool:
"""True when MASTER slot cannot accept synthetic PTT for ``tg``."""
from .routing.helpers import master_dynamic_tg_slots, master_slot_holds_server_broadcast
from .server_voice import all_server_voice_ids
if master_slot_holds_server_broadcast(
slot,
time.time(),
server_voice_rf_srcs=all_server_voice_ids(self._config),
):
return False
# External QSO on RX while slot is in hangtime — defer announcement inject.
if slot.get("RX_TYPE") != HBPF_SLT_VTERM and slot.get("TX_TYPE") == HBPF_SLT_VTERM:
return True
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
return True
def _inject_ptt_slot_order(self, tg: int, sys_cfg: dict[str, Any]) -> list[int]:
"""Prefer the RF slot where ``tg`` is already a dynamic UA session."""
from .routing.helpers import master_dynamic_tg_slots
dynamic = sorted(master_dynamic_tg_slots(sys_cfg, int(tg)), reverse=True)
preferred = list(dynamic)
if self._announcement_ptt_system:
for ts in sorted(
self._active_bridge_slots_for_tg(tg, self._announcement_ptt_system),
reverse=True,
):
if ts not in preferred:
preferred.append(ts)
ordered: list[int] = []
for ts in [*preferred, 2, 1]:
if ts not in ordered:
ordered.append(ts)
return ordered
def _build_inject_ptt_targets(
self, label: str, tg: int,
) -> tuple[list[dict[str, Any]], int]:
"""Synthetic PTT on proxy MASTER: routing creates the UA bridge on inject."""
targets: list[dict[str, Any]] = []
busy_count = 0
ptt_system = self._announcement_ptt_system
if not ptt_system or not self._inject_announcement_ptt:
return targets, busy_count
protocols = self._get_protocols() if self._get_protocols else {}
systems_cfg = self._config.get("SYSTEMS", {})
if ptt_system not in protocols or ptt_system not in systems_cfg:
return targets, busy_count
if systems_cfg[ptt_system].get("MODE") != "MASTER":
return targets, busy_count
sys_obj = protocols[ptt_system]
status = getattr(sys_obj, "STATUS", None)
if not status:
return targets, busy_count
sys_cfg = systems_cfg[ptt_system]
for ts in self._inject_ptt_slot_order(tg, sys_cfg):
slot_index = 2 if ts == 2 else 1
slot = status.get(slot_index)
if not slot:
continue
if self._inject_ptt_slot_busy(slot, tg, sys_cfg, ts):
logger.debug("(%s) System %s TS%s busy (QSO active), skipping", label, ptt_system, ts)
busy_count += 1
continue
targets.append({"sys_obj": sys_obj, "name": ptt_system, "slot": slot, "ts": ts})
break
return targets, busy_count
def _build_announcement_targets(
self, tg_int: int, tg_str: str, label: str
) -> tuple[list[dict[str, Any]], int]:
"""MASTER targets for announcement/TTS broadcast.
Inject path: synthetic PTT on the TG's dynamic UA slot when applicable.
Legacy path: MASTER systems with an ACTIVE bridge row for the TG.
"""
if self._inject_announcement_ptt:
return self._build_inject_ptt_targets(label, tg_int)
targets: list[dict[str, Any]] = []
busy_count = 0
protocols = self._get_protocols() if self._get_protocols else {}
systems_cfg = self._config.get("SYSTEMS", {})
bridges = self._routing_table_for_report() if self._routing_table_for_report else {}
bridge_entries = bridges.get(tg_str, [])
for sys_name in list(protocols.keys()):
if sys_name in _ANNOUNCEMENT_EXCLUDED or any(
sys_name.startswith(ex + "-") for ex in _ANNOUNCEMENT_EXCLUDED
):
continue
if sys_name not in systems_cfg or systems_cfg[sys_name].get("MODE") != "MASTER":
continue
if not systems_cfg[sys_name].get("PEERS"):
continue
has_peers = any(
systems_cfg[sys_name]["PEERS"].get(pid, {}).get("CALLSIGN")
for pid in systems_cfg[sys_name]["PEERS"]
)
if not has_peers or sys_name not in protocols:
continue
sys_obj = protocols[sys_name]
if not getattr(sys_obj, "STATUS", None):
continue
active_slots = [
be["TS"] for be in bridge_entries
if be.get("SYSTEM") == sys_name and be.get("ACTIVE") and be.get("TS") is not None
]
active_slots = list(dict.fromkeys(active_slots))
if not active_slots:
continue
for ts in active_slots:
slot_index = 2 if ts == 2 else 1
slot = sys_obj.STATUS.get(slot_index)
if not slot:
continue
rx_type = slot.get("RX_TYPE")
tx_type = slot.get("TX_TYPE")
if (rx_type != HBPF_SLT_VTERM) or (tx_type != HBPF_SLT_VTERM):
logger.debug("(%s) System %s TS%s busy (QSO active), skipping", label, sys_name, ts)
busy_count += 1
continue
targets.append({"sys_obj": sys_obj, "name": sys_name, "slot": slot, "ts": ts})
return targets, busy_count
def _announcement_packet_peer(
self,
tg: int,
targets: list[dict[str, Any]],
server_id: bytes,
) -> bytes:
"""DMRD peer field: always GLOBAL SERVER_ID on inject (never a hotspot radio id)."""
del tg, targets
if self._inject_announcement_ptt:
return self._global_server_id_bytes()
return bytes_4(int_id(server_id))
def _send_filtered_by_tg(
self, sys_obj: Any, pkt: bytes, tg: int, ts: int, bridges: dict[str, list[dict[str, Any]]]
) -> int:
"""Return -1 if sent, 0 if TG/TS not active (legacy _sendFilteredByTG)."""
tg_str = str(tg)
for be in bridges.get(tg_str, []):
if be.get("SYSTEM") == getattr(sys_obj, "_system", None) and be.get("TS") == ts and be.get("ACTIVE"):
sys_obj.send_system(pkt)
return -1
return 0
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}")
parts.append("1")
self._send_routing_event(",".join(parts))
def _maybe_begin_legacy_voice_report(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
tg: int,
) -> None:
# Inject: routing emits START,RX (SERVER_ID) + OBP TX; same-MASTER downlink has no
# bridge leg, so emit START,TX on the proxy MASTER for monitor fan-out (SYSTEM-N).
self._begin_legacy_voice_report(targets, pkts_by_ts, tg)
def _maybe_end_legacy_voice_report(self) -> None:
self._end_legacy_voice_report()
def _begin_legacy_voice_report(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
tg: int,
) -> None:
if not self._send_routing_event or not targets:
return
wire_ts = targets[0]["ts"]
pkts = pkts_by_ts.get(wire_ts) or []
if not pkts:
return
stream_id = pkts[0][16:20]
rf_src = int_id(pkts[0][5:8])
systems: list[tuple[str, int]] = []
seen: set[tuple[str, int]] = set()
for t in targets:
key = (t["name"], t["ts"])
if key in seen:
continue
seen.add(key)
systems.append(key)
self._emit_announcement_voice_event(
"START", "TX", t["name"], stream_id, t["ts"], tg, rf_src
)
self._voice_report_state = {
"stream_id": stream_id,
"rf_src": pkts[0][5:8],
"start": time.time(),
"tg": tg,
"systems": systems,
}
def _end_legacy_voice_report(self) -> None:
state = self._voice_report_state
self._voice_report_state = None
if not state or not self._send_routing_event:
return
duration = time.time() - float(state["start"])
stream_id = state["stream_id"]
tg = int(state["tg"])
for name, ts in state["systems"]:
self._emit_announcement_voice_event(
"END", "TX", name, stream_id, ts, tg, int_id(state["rf_src"]), duration
)
def _send_announcement_packets(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
pkt_idx: int,
source_id: bytes,
dst_id: bytes,
tg: int,
label: str,
) -> None:
"""One frame: synthetic PTT via routing, or legacy per-MASTER send_system."""
if self._inject_announcement_ptt:
wire_ts = targets[0]["ts"] if targets else 2
pkt = pkts_by_ts[wire_ts][pkt_idx]
if self._inject_announcement_ptt(pkt, time.time()) is False:
logger.warning("(%s) Routing rejected announcement frame %s", label, pkt_idx)
return
bridges = self._routing_table_for_report() if self._routing_table_for_report else {}
now = time.time()
for t in targets:
try:
sys_obj = t["sys_obj"]
slot = t["slot"]
t_ts = t["ts"]
pkt = pkts_by_ts[t_ts][pkt_idx]
stream_id = pkt[16:20]
if stream_id not in sys_obj.STATUS:
sys_obj.STATUS[stream_id] = {
"START": now,
"CONTENTION": False,
"RFS": source_id,
"TGID": dst_id,
"LAST": now,
}
slot["TX_TGID"] = dst_id
slot["TX_RFS"] = source_id
else:
sys_obj.STATUS[stream_id]["LAST"] = now
slot["TX_TIME"] = now
self._send_filtered_by_tg(sys_obj, pkt, tg, t_ts, bridges)
except Exception as e:
logger.error(
"(%s) Error sending packet %s to %s/TS%s: %s",
label,
pkt_idx,
t.get("name"),
t.get("ts"),
e,
)
def _mark_slots_busy(
self,
targets: list[dict[str, Any]],
source_id: bytes | None = None,
) -> None:
"""Mark target slots busy (TX_TYPE=VHEAD) to prevent TS conflict."""
server_rfs = source_id if source_id is not None else self._server_source_id()
now = time.time()
for t in targets:
try:
slot = t.get("slot")
if slot is not None:
slot["TX_TYPE"] = HBPF_SLT_VHEAD
slot["TX_TIME"] = now
slot["TX_RFS"] = server_rfs
except (KeyError, TypeError):
pass
def _mark_slots_free(self, targets: list[dict[str, Any]]) -> None:
"""Mark target slots free (TX_TYPE=VTERM) when broadcast done."""
for t in targets:
try:
slot = t.get("slot")
if slot is not None:
slot["TX_TYPE"] = HBPF_SLT_VTERM
except (KeyError, TypeError):
pass
def _enqueue_broadcast(
self, _type: str, targets: list[dict[str, Any]], pkts_by_ts: dict[int, list[bytes]],
source_id: bytes, dst_id: bytes, tg: int, num: int, label: str,
) -> None:
_tg_key = str(tg)
if _tg_key in self._broadcast_active_tgs:
self._broadcast_queue.append({
'type': _type, 'targets': targets, 'pkts_by_ts': pkts_by_ts,
'source_id': source_id, 'dst_id': dst_id, 'tg': tg, 'num': num, 'label': label,
})
_pos = len(self._broadcast_queue)
logger.info('(%s) Same TG %s still on air; deferring next playback (pending %s)', label, tg, _pos)
else:
self._broadcast_active_tgs.add(_tg_key)
self._mark_slots_busy(targets, source_id)
logger.info('(%s) Starting broadcast immediately for TG %s (active TGs: %s)', label, tg, len(self._broadcast_active_tgs))
if self._call_later:
if _type == 'ann':
self._call_later(0.5, self._announcement_send_broadcast, targets, pkts_by_ts, 0, source_id, dst_id, tg, num, label, None)
elif _type == 'tts':
self._call_later(0.5, self._tts_send_broadcast, targets, pkts_by_ts, 0, source_id, dst_id, tg, num, label, None)
def _start_next_broadcast(self) -> None:
if not self._broadcast_queue:
return
_next = None
for i, _item in enumerate(self._broadcast_queue):
_tg_key = str(_item['tg'])
if _tg_key not in self._broadcast_active_tgs:
_next = self._broadcast_queue.pop(i)
break
if not _next:
return
_type = _next['type']
_label = _next['label']
_tg_key = str(_next['tg'])
self._broadcast_active_tgs.add(_tg_key)
self._mark_slots_busy(_next['targets'], _next['source_id'])
logger.info('(%s) Starting deferred same-TG playback for TG %s (%s pending, %s active TG(s))', _label, _next['tg'], len(self._broadcast_queue), len(self._broadcast_active_tgs))
if self._call_later:
if _type == 'ann':
self._call_later(0.5, self._announcement_send_broadcast, _next['targets'], _next['pkts_by_ts'], 0, _next['source_id'], _next['dst_id'], _next['tg'], _next['num'], _label, None)
elif _type == 'tts':
self._call_later(0.5, self._tts_send_broadcast, _next['targets'], _next['pkts_by_ts'], 0, _next['source_id'], _next['dst_id'], _next['tg'], _next['num'], _label, None)
def _broadcast_finished(self, tg: int | None = None) -> None:
if tg is not None:
self._broadcast_active_tgs.discard(str(tg))
if self._broadcast_queue:
logger.info(
'(BROADCAST) TG %s announcement/TTS playback finished; %s same-TG deferred, %s active TG(s)',
tg, len(self._broadcast_queue), len(self._broadcast_active_tgs),
)
if self._call_later:
self._call_later(_BROADCAST_GAP, self._start_next_broadcast)
else:
if not self._broadcast_active_tgs:
logger.info('(BROADCAST) All announcement/TTS playbacks finished')
else:
logger.info(
'(BROADCAST) TG %s announcement/TTS playback finished; %s other TG(s) still on air',
tg, len(self._broadcast_active_tgs),
)
def _announcement_send_broadcast(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
pkt_idx: int,
source_id: bytes,
dst_id: bytes,
tg: int,
ann_idx: int,
label: str,
next_time: float | None = None,
) -> None:
"""Send one batch of packets; schedule next via call_later."""
total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total or not targets:
self._mark_slots_free(targets)
self._maybe_end_legacy_voice_report()
for t in targets:
try:
obj = t.get("sys_obj")
if getattr(obj, "STATUS", None):
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception as e:
logger.warning("(%s) slot STATUS cleanup failed: %s", label, e)
self._announcement_running[ann_idx] = False
if not targets:
logger.info(
"(%s) Broadcast aborted at packet %s/%s: all targets removed (QSO collision)",
label,
pkt_idx,
total,
)
else:
logger.info("(%s) Broadcast complete: %s packets sent to %s targets", label, total, len(targets))
self._broadcast_finished(tg)
return
collided: list[dict[str, Any]] = []
for t in targets:
slot = t["slot"]
if slot.get("RX_TYPE") != HBPF_SLT_VTERM:
logger.info(
"(%s) QSO detected on %s/TS%s during broadcast (packet %s/%s), removing target",
label,
t["name"],
t["ts"],
pkt_idx,
total,
)
slot["TX_TYPE"] = HBPF_SLT_VTERM
collided.append(t)
for t in collided:
targets.remove(t)
if not targets:
self._announcement_running[ann_idx] = False
self._maybe_end_legacy_voice_report()
logger.info(
"(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s",
label,
pkt_idx,
total,
)
self._broadcast_finished(tg)
return
now = time.time()
if pkt_idx == 0:
self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg)
self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label)
if next_time is None:
next_time = now + _FRAME_INTERVAL
else:
next_time = next_time + _FRAME_INTERVAL
delay = max(0.001, next_time - time.time())
if self._call_later:
self._call_later(
delay,
self._announcement_send_broadcast,
targets,
pkts_by_ts,
pkt_idx + 1,
source_id,
dst_id,
tg,
ann_idx,
label,
next_time,
)
def scheduled_announcement(self, ann_idx: int = 0, _retry: int = 0) -> None:
"""Run one scheduled file announcement from ANNOUNCEMENTS[ann_idx]."""
g = self._config.get("VOICE", {})
announcements = g.get("ANNOUNCEMENTS") or []
if ann_idx < 0 or ann_idx >= len(announcements):
return
item = announcements[ann_idx]
if not isinstance(item, dict) or not item.get("ENABLED"):
return
label = "ANNOUNCEMENT-{}".format(ann_idx + 1)
if self._announcement_running.get(ann_idx):
if _retry == 0:
logger.debug("(%s) Previous announcement still running, skipping", label)
return
mode = item.get("MODE", "interval")
if mode == "hourly" and _retry == 0:
now = datetime.now()
if now.minute != 0:
return
if self._announcement_last_hour.get(ann_idx) == now.hour:
return
_tg = int(item.get("TG", 0))
if str(_tg) in self._broadcast_active_tgs and _retry < 60:
if _retry == 0:
logger.debug("(%s) Same TG %s already broadcasting, deferring prep", label, _tg)
if self._call_later:
self._call_later(3.0 + ann_idx * 0.5, self.scheduled_announcement, ann_idx, _retry + 1)
return
_file = str(item.get("FILE") or "").strip()
_lang = item.get("LANGUAGE", "en_GB")
if not _file or not _tg:
return
_dst_id = bytes_3(_tg)
_source_id = announcement_item_source_bytes(item, self._config)
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
logger.info("(%s) Playing file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang)
try:
_say = self.read_single_file(self._audio_path, _lang, str(_file))
except Exception as e:
logger.warning("(%s) Cannot read AMBE file: Audio/%s/ondemand/%s.ambe: %s", label, _lang, _file, e)
return
if not _say:
logger.warning("(%s) AMBE file empty or not found: %s/ondemand/%s.ambe", label, _lang, _file)
return
tg_str = str(_tg)
targets, busy_count = self._build_announcement_targets(_tg, tg_str, label)
if not targets:
if busy_count > 0 and _retry < 60:
if _retry == 0:
logger.info(
"(%s) All %s target slots busy (QSO active), waiting for QSO to finish...",
label,
busy_count,
)
if self._call_later:
self._call_later(5.0, self.scheduled_announcement, ann_idx, _retry + 1)
return
logger.info("(%s) No systems with active bridge for TG %s to send to", label, _tg)
return
if mode == "hourly":
self._announcement_last_hour[ann_idx] = datetime.now().hour
_say_list = [_say]
pkt_peer = self._announcement_packet_peer(_tg, targets, server_id)
pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 0, _say_list)),
2: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 1, _say_list)),
}
ts1_count = sum(1 for t in targets if t["ts"] == 1)
ts2_count = sum(1 for t in targets if t["ts"] == 2)
sys_names = ", ".join("{}/TS{}".format(t["name"], t["ts"]) for t in targets[:8])
if len(targets) > 8:
sys_names += ", ... +{}".format(len(targets) - 8)
logger.info(
"(%s) Broadcasting %s packets to %s targets (TS1:%s TS2:%s): %s",
label, len(pkts_by_ts[1]), len(targets), ts1_count, ts2_count, sys_names,
)
self._announcement_running[ann_idx] = True
self._enqueue_broadcast('ann', targets, pkts_by_ts, _source_id, _dst_id, _tg, ann_idx, label)
def scheduled_tts_announcement(self, tts_idx: int = 0, _retry: int = 0) -> None:
"""Run one scheduled TTS announcement from TTS_ANNOUNCEMENTS[tts_idx]."""
g = self._config.get("VOICE", {})
tts_list = g.get("TTS_ANNOUNCEMENTS") or []
if tts_idx < 0 or tts_idx >= len(tts_list):
return
item = tts_list[tts_idx]
if not isinstance(item, dict) or not item.get("ENABLED", False):
return
label = "TTS-{}".format(tts_idx + 1)
if self._tts_running.get(tts_idx):
if _retry == 0:
logger.debug("(%s) Previous TTS announcement still running, skipping", label)
return
mode = item.get("MODE", "interval")
if mode == "hourly" and _retry == 0:
now = datetime.now()
if now.minute != 0:
return
if self._tts_last_hour.get(tts_idx) == now.hour:
return
_tg = int(item.get("TG", 0))
if str(_tg) in self._broadcast_active_tgs and _retry < 60:
if _retry == 0:
logger.debug("(%s) Same TG %s already broadcasting, deferring TTS prep", label, _tg)
if self._call_later:
self._call_later(3.0 + tts_idx * 0.5, self.scheduled_tts_announcement, tts_idx, _retry + 1)
return
_file = str(item.get("FILE") or "").strip()
_lang = item.get("LANGUAGE", "en_GB")
self._tts_running[tts_idx] = True
logger.info("(%s) Starting TTS conversion in background thread for %s", label, _file)
if self._defer_to_thread:
d = self._defer_to_thread(self._voice.ensure_tts_ambe, self._config, item, self._audio_path)
d.addCallback(self._tts_conversion_done, tts_idx, _file, _tg, _lang, mode, label)
d.addErrback(self._tts_conversion_error, tts_idx, label)
else:
try:
ambe_path = self._voice.ensure_tts_ambe(self._config, item, self._audio_path)
self._tts_conversion_done(ambe_path, tts_idx, _file, _tg, _lang, mode, label)
except Exception as e:
self._tts_conversion_error(e, tts_idx, label)
def _tts_conversion_done(
self, ambe_path: str | None, tts_idx: int, _file: str, _tg: int, _lang: str, mode: str, label: str, _retry: int = 0
) -> None:
"""After TTS conversion: broadcast like scheduled_announcement."""
if not ambe_path:
self._tts_running[tts_idx] = False
logger.warning("(%s) No AMBE file available for TTS announcement %s", label, _file)
return
if str(_tg) in self._broadcast_active_tgs and _retry < 60:
if _retry == 0:
logger.debug("(%s) Same TG %s already broadcasting, deferring TTS packet prep", label, _tg)
if self._call_later:
self._call_later(3.0 + tts_idx * 0.5, self._tts_conversion_done, ambe_path, tts_idx, _file, _tg, _lang, mode, label, _retry + 1)
return
logger.info("(%s) Playing TTS file: %s to TG %s (both TS, mode: %s, lang: %s)", label, _file, _tg, mode, _lang)
_dst_id = bytes_3(_tg)
tts_list = self._config.get("VOICE", {}).get("TTS_ANNOUNCEMENTS") or []
tts_item = tts_list[tts_idx] if 0 <= tts_idx < len(tts_list) else None
_source_id = announcement_item_source_bytes(
tts_item if isinstance(tts_item, dict) else None,
self._config,
)
server_id = self._config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_3(int(server_id))
_file_base = _file.replace(".ambe", "")
_say = self.read_single_file(self._audio_path, _lang, _file_base)
if not _say:
logger.warning("(%s) Cannot read AMBE file: %s", label, ambe_path)
self._tts_running[tts_idx] = False
return
tg_str = str(_tg)
targets, busy_count = self._build_announcement_targets(_tg, tg_str, label)
if not targets:
if busy_count > 0 and _retry < 60:
if _retry == 0:
logger.info(
"(%s) All %s target slots busy (QSO active), waiting for QSO to finish...",
label,
busy_count,
)
if self._call_later:
self._call_later(
5.0,
self._tts_conversion_done,
ambe_path,
tts_idx,
_file,
_tg,
_lang,
mode,
label,
_retry + 1,
)
return
self._tts_running[tts_idx] = False
logger.info("(%s) No systems with active bridge for TG %s to send to", label, _tg)
return
if mode == "hourly":
self._tts_last_hour[tts_idx] = datetime.now().hour
_say_list = [_say]
pkt_peer = self._announcement_packet_peer(_tg, targets, server_id)
pkts_by_ts = {
1: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 0, _say_list)),
2: list(self.pkt_gen(_source_id, _dst_id, pkt_peer, 1, _say_list)),
}
logger.info("(%s) Broadcasting %s packets to %s targets", label, len(pkts_by_ts[1]), len(targets))
self._enqueue_broadcast('tts', targets, pkts_by_ts, _source_id, _dst_id, _tg, tts_idx, label)
def _tts_conversion_error(self, failure: Any, tts_idx: int, label: str) -> None:
self._tts_running[tts_idx] = False
try:
msg = failure.getErrorMessage()
except Exception as e:
logger.warning("(%s) failure.getErrorMessage unavailable: %s", label, e)
msg = str(failure)
logger.error("(%s) TTS conversion error: %s", label, msg)
def _tts_send_broadcast(
self,
targets: list[dict[str, Any]],
pkts_by_ts: dict[int, list[bytes]],
pkt_idx: int,
source_id: bytes,
dst_id: bytes,
tg: int,
tts_idx: int,
label: str,
next_time: float | None = None,
) -> None:
"""Same as _announcement_send_broadcast but clears _tts_running."""
total = len(pkts_by_ts.get(1, []))
if pkt_idx >= total or not targets:
self._mark_slots_free(targets)
self._maybe_end_legacy_voice_report()
for t in targets:
try:
obj = t.get("sys_obj")
if getattr(obj, "STATUS", None):
for sid in list(obj.STATUS.keys()):
if sid not in (1, 2):
del obj.STATUS[sid]
except Exception as e:
logger.warning("(%s) slot STATUS cleanup failed: %s", label, e)
self._tts_running[tts_idx] = False
if not targets:
logger.info(
"(%s) Broadcast aborted at packet %s/%s: all targets removed (QSO collision)",
label,
pkt_idx,
total,
)
else:
logger.info("(%s) Broadcast complete: %s packets sent to %s targets", label, total, len(targets))
self._broadcast_finished(tg)
return
collided: list[dict[str, Any]] = []
for t in targets:
slot = t["slot"]
if slot.get("RX_TYPE") != HBPF_SLT_VTERM:
logger.info(
"(%s) QSO detected on %s/TS%s during broadcast (packet %s/%s), removing target",
label,
t["name"],
t["ts"],
pkt_idx,
total,
)
slot["TX_TYPE"] = HBPF_SLT_VTERM
collided.append(t)
for t in collided:
targets.remove(t)
if not targets:
self._tts_running[tts_idx] = False
self._maybe_end_legacy_voice_report()
logger.info(
"(%s) Broadcast stopped: all targets had QSO collision at packet %s/%s",
label,
pkt_idx,
total,
)
self._broadcast_finished(tg)
return
now = time.time()
if pkt_idx == 0:
self._maybe_begin_legacy_voice_report(targets, pkts_by_ts, tg)
self._send_announcement_packets(targets, pkts_by_ts, pkt_idx, source_id, dst_id, tg, label)
if next_time is None:
next_time = now + _FRAME_INTERVAL
else:
next_time = next_time + _FRAME_INTERVAL
delay = max(0.001, next_time - time.time())
if self._call_later:
self._call_later(
delay,
self._tts_send_broadcast,
targets,
pkts_by_ts,
pkt_idx + 1,
source_id,
dst_id,
tg,
tts_idx,
label,
next_time,
)
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
"""Read one AMBE file (e.g. ondemand/{file_number}.ambe). Legacy readSingleFile."""
return self._voice.read_single_file(audio_path, lang, file_number)
@ -1020,73 +205,3 @@ class VoiceUseCases:
logger.debug("(%s) Sending disconnected voice", system)
self.play_on_slot(protocol, system, speech, _source_id, bytes_3(9))
logger.debug("(%s) disconnected voice thread end", system)
def apply_voice_config(self) -> None:
"""Start/stop announcement and TTS LoopingCalls from ``config["VOICE"]``."""
g = self._config.get("VOICE", {})
if not self._start_looping_call:
return
announcements = g.get("ANNOUNCEMENTS") or []
if not isinstance(announcements, list):
announcements = []
for ann_idx in list(self._ann_tasks.keys()):
if ann_idx >= len(announcements) or not (isinstance(announcements[ann_idx], dict) and announcements[ann_idx].get("ENABLED")):
try:
if getattr(self._ann_tasks[ann_idx], "running", False):
self._ann_tasks[ann_idx].stop()
except Exception as e:
logger.warning("(VOICE-RELOAD) stop announcement task %s failed: %s", ann_idx + 1, e)
del self._ann_tasks[ann_idx]
logger.info("(VOICE-RELOAD) ANNOUNCEMENT-%s stopped", ann_idx + 1)
for ann_idx, item in enumerate(announcements):
if not isinstance(item, dict) or not item.get("ENABLED"):
continue
label = "ANNOUNCEMENT-{}".format(ann_idx + 1)
if ann_idx in self._ann_tasks:
try:
if getattr(self._ann_tasks[ann_idx], "running", False):
self._ann_tasks[ann_idx].stop()
except Exception as e:
logger.warning("(VOICE-RELOAD) stop %s failed: %s", label, e)
del self._ann_tasks[ann_idx]
logger.info("(VOICE-RELOAD) %s stopped", label)
mode = item.get("MODE", "interval")
interval = 30.0 if mode == "hourly" else float(item.get("INTERVAL", 60))
lc = self._start_looping_call(lambda ai=ann_idx: self.scheduled_announcement(ai), interval, False)
self._ann_tasks[ann_idx] = lc
logger.info(
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s",
label, mode, item.get("FILE"), item.get("TG"),
)
tts_list = g.get("TTS_ANNOUNCEMENTS") or []
if not isinstance(tts_list, list):
tts_list = []
for tts_idx in list(self._tts_tasks.keys()):
if tts_idx >= len(tts_list) or not (isinstance(tts_list[tts_idx], dict) and tts_list[tts_idx].get("ENABLED")):
try:
if getattr(self._tts_tasks[tts_idx], "running", False):
self._tts_tasks[tts_idx].stop()
except Exception as e:
logger.warning("(VOICE-RELOAD) stop TTS task %s failed: %s", tts_idx + 1, e)
del self._tts_tasks[tts_idx]
logger.info("(VOICE-RELOAD) TTS-%s stopped", tts_idx + 1)
for tts_idx, item in enumerate(tts_list):
if not isinstance(item, dict) or not item.get("ENABLED"):
continue
label = "TTS-{}".format(tts_idx + 1)
if tts_idx in self._tts_tasks:
try:
if getattr(self._tts_tasks[tts_idx], "running", False):
self._tts_tasks[tts_idx].stop()
except Exception as e:
logger.warning("(VOICE-RELOAD) stop %s failed: %s", label, e)
del self._tts_tasks[tts_idx]
logger.info("(VOICE-RELOAD) %s stopped", label)
mode = item.get("MODE", "interval")
interval = 30.0 if mode == "hourly" else float(item.get("INTERVAL", 60))
lc = self._start_looping_call(lambda ti=tts_idx: self.scheduled_tts_announcement(ti), interval, False)
self._tts_tasks[tts_idx] = lc
logger.info(
"(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s",
label, mode, item.get("FILE"), item.get("TG"),
)

@ -415,24 +415,14 @@ def run_peer_server(
voice_provider = DefaultVoiceProvider()
else:
voice_provider = StubVoiceProvider()
def _start_voice_loop(fn, interval: float, now: bool):
lc = task.LoopingCall(fn)
d = lc.start(interval, now=now)
d.addErrback(_looping_errback, logger)
return lc
# Scheduled announcements and TTS: plugins/voice-announcements (via PLUGINS.send).
voice_use_cases = VoiceUseCases(
voice_provider,
config,
get_protocols=lambda: protocols,
call_from_reactor=reactor.callFromThread,
audio_path=audio_path,
routing_table_for_report=lambda: {},
call_later=reactor.callLater,
start_looping_call=_start_voice_loop,
defer_to_thread=threads.deferToThread,
)
voice_use_cases.apply_voice_config()
ident_use_cases = IdentUseCases(
config,
voice_use_cases,
@ -455,14 +445,17 @@ def run_peer_server(
if callable(send_system):
send_system(pkt)
plugin_ingress = PluginIngress(None, config, lambda: protocols, _send_local) # routing set below
plugin_ingress = PluginIngress( # routing set below
None, config, lambda: protocols, _send_local, send_routing_event=reporting_use_cases.send_routing_event,
)
plugin_manager = PluginManager(
plugin_bus,
server_ctx,
project_root,
sender_factory=lambda name: PluginDmrdSender(
name, config, plugin_ingress.deliver, reactor.callFromThread, in_reactor_thread=isInIOThread,
name, config, plugin_ingress.deliver, reactor.callFromThread,
in_reactor_thread=isInIOThread, slot_for_tg=plugin_ingress.voice_slot_for_tg,
),
)
voice_plugin_bridge = VoicePluginBridge(plugin_bus, config, get_dmra_blocks=get_dmra_blocks)
@ -494,36 +487,6 @@ def run_peer_server(
on_restored=routing_use_cases.sync_restored_dynamic_tgs,
)
routing_table_for_report = routing_use_cases.routing_table_for_report
voice_use_cases._routing_table_for_report = routing_table_for_report
from adn_server.application.routing.announcement_ptt_inject import (
announcement_ptt_system,
inject_announcement_ptt,
)
_ptt_system = announcement_ptt_system(config)
def _inject_announcement_ptt(pkt: bytes, pkt_time: float) -> bool | None:
if not _ptt_system:
return False
server_id = config.get("GLOBAL", {}).get("SERVER_ID", b"\x00\x00\x00\x00")
if not isinstance(server_id, bytes):
server_id = bytes_4(int(server_id or 0) & 0xFFFFFFFF)
accepted = inject_announcement_ptt(
routing_use_cases,
_ptt_system,
pkt,
pkt_time=pkt_time,
server_id=server_id,
)
proto = protocols.get(_ptt_system)
send_system = getattr(proto, "send_system", None) if proto is not None else None
if callable(send_system):
send_system(pkt)
return accepted
voice_use_cases._inject_announcement_ptt = _inject_announcement_ptt
voice_use_cases._announcement_ptt_system = _ptt_system
voice_use_cases._send_routing_event = reporting_use_cases.send_routing_event
report_factory.set_routing_table(routing_table_for_report())
report_factory.set_systems(config.get("SYSTEMS", {}))
@ -882,8 +845,7 @@ def run_peer_server(
return
voice_config_mtime = mtime
logger.info("(VOICE-RELOAD) config file change detected, reloading configuration...")
loader.reload_voice_config(config, voice_config_path)
voice_use_cases.apply_voice_config()
loader.reload_voice_config(config, voice_config_path) # voice-announcements follows config["VOICE"]
if voice_config_path and voice_config_mtime:
logger.info("(VOICE-RELOAD) config reload completed")
except Exception as e:

@ -36,7 +36,6 @@ from bitarray import bitarray
from ...application.ports import VoiceProvider
from .pkt_gen import pkt_gen as _pkt_gen
from .tts_engine import ensure_tts_ambe as _tts_ensure_tts_ambe
from .voice_map import VOICE_MAP
logger = logging.getLogger(__name__)
@ -179,10 +178,6 @@ class DefaultVoiceProvider(VoiceProvider):
"""Generate HBP voice packets for phrase. Legacy mk_voice.pkt_gen."""
return _pkt_gen(rf_src, dst_id, peer, slot, phrase)
def ensure_tts_ambe(self, config: dict[str, Any], item: dict[str, Any], audio_path: str) -> str | None:
"""Delegate to tts_engine.ensure_tts_ambe (full TTS pipeline)."""
return _tts_ensure_tts_ambe(config, item, audio_path)
class StubVoiceProvider(VoiceProvider):
"""Stub: get_ambe_words returns empty; pkt_gen returns empty iterator; read_single_file returns []."""
@ -195,8 +190,5 @@ class StubVoiceProvider(VoiceProvider):
) -> Iterator[bytes]:
return iter([])
def ensure_tts_ambe(self, config: dict[str, Any], item: dict[str, Any], audio_path: str) -> str | None:
return None
def read_single_file(self, audio_path: str, lang: str, file_number: str) -> list:
return []

@ -87,14 +87,12 @@ Mark new stack tests with `@pytest.mark.integration`.
| File | Tests | Topic |
|------|-------|-------|
| `test_announcement_anticollision.py` | 3 | Busy slot skip / abort |
| `test_broadcast_queue.py` | 2 | Same-TG broadcast queue |
| `test_disconnected_voice.py` | 3 | Not-linked / reflector prompts |
| `test_in_band_signalling.py` | 5 | Reflector / single-mode VTERM |
| `test_play_file_on_request.py` | 3 | On-demand file playback |
| `test_scheduled_announcement.py` | 4 | File announcements (AMBE) |
| `test_scheduled_tts.py` | 7 | TTS schedule + conversion |
| `test_voice_config_reload.py` | 3 | Hot reload announcement/TTS loops |
| `test_prompt_holds_slot.py` | 4 | Server prompts hold TS2 while they play |
Scheduled announcements, TTS and beacons are tested with their plugin: `plugins/voice-announcements/tests/`.
### talker_alias/

@ -250,3 +250,21 @@ def test_a_rate_that_is_not_a_positive_number_grants_nothing(rate) -> None:
from adn_server.application.plugins.domain.send import send_permission
assert send_permission(_config(max_frames_per_s=rate), "d-aprs") is None
def test_voice_slot_query_only_for_granted_talkgroups() -> None:
sender = PluginDmrdSender(
"d-aprs", _config(group_voice_tgs=[213]), lambda *a: True, call_from_reactor=lambda *a: None,
slot_for_tg=lambda tg: 2,
)
assert sender.voice_slot_for_tg(213) == 2
assert sender.voice_slot_for_tg(214) is None
def test_the_default_rate_leaves_room_for_one_voice_stream_per_talkgroup() -> None:
from adn_server.application.plugins.domain.send import send_permission
cfg = {"PLUGINS": {"send": {"beacon": {"allowed_src_ids": [1], "group_voice_tgs": [213, 214, 9]}}}}
assert send_permission(cfg, "beacon").max_frames_per_s == 40 + 3 * 18
cfg["PLUGINS"]["send"]["beacon"]["max_frames_per_s"] = 25
assert send_permission(cfg, "beacon").max_frames_per_s == 25

@ -24,11 +24,6 @@ from __future__ import annotations
from typing import Any, Iterator
from adn_server.application.routing.announcement_ptt_inject import (
announcement_ptt_system,
inject_announcement_ptt,
)
from adn_server.application.voice_use_cases import VoiceUseCases
from adn_server.domain import bytes_3, bytes_4
from tests.harness.deterministic import DeterministicScenario, FakeHbpProtocol, active_routing_table
@ -55,11 +50,6 @@ class FakeVoiceProvider:
del audio_path, lang, file_number
return [b"\x00" * 7]
def ensure_tts_ambe(self, config: dict, item: dict, audio_path: str) -> str | None:
del config, item, audio_path
return "/tmp/fake.ambe"
class FakeMasterForVoice(FakeHbpProtocol):
def __init__(self, name: str) -> None:
super().__init__(name)
@ -105,98 +95,3 @@ def reflector_routing_entry(system: str = "MASTER-A", reflector: int = 310) -> d
}
]
}
def make_voice_uc(
scenario: DeterministicScenario,
master: FakeMasterForVoice,
*,
audio_path: str = "/tmp/audio",
) -> VoiceUseCases:
scheduled: list[tuple[float, tuple]] = []
ptt_system = announcement_ptt_system(scenario.config)
server_id = scenario.config.get("GLOBAL", {}).get("SERVER_ID", bytes_4(9990))
if not isinstance(server_id, bytes):
server_id = bytes_4(int(server_id or 0) & 0xFFFFFFFF)
def call_later(delay, fn, *args):
scheduled.append((delay, (fn, args)))
return type("H", (), {"active": lambda self: True, "cancel": lambda self: None})()
def inject_announcement_ptt_cb(pkt: bytes, pkt_time: float) -> bool | None:
if not ptt_system:
return False
accepted = inject_announcement_ptt(
scenario.routing,
ptt_system,
pkt,
pkt_time=pkt_time,
server_id=server_id,
)
if hasattr(master, "send_system"):
master.send_system(pkt)
return accepted
uc = VoiceUseCases(
FakeVoiceProvider(),
scenario.config,
get_protocols=lambda: {"MASTER-A": master, **{k: v for k, v in scenario.protocols.items() if k != "MASTER-A"}},
routing_table_for_report=scenario.routing.routing_table_for_report,
call_later=call_later,
audio_path=audio_path,
inject_announcement_ptt=inject_announcement_ptt_cb,
)
uc._announcement_ptt_system = ptt_system
uc._scheduled = scheduled
return uc
def drain_call_later(uc: VoiceUseCases, max_rounds: int = 200) -> None:
scheduled = getattr(uc, "_scheduled", None)
if scheduled is None:
return
for _ in range(max_rounds):
if not scheduled:
break
_, (fn, args) = scheduled.pop(0)
fn(*args)
def voice_announcement_config(
scenario: DeterministicScenario,
*,
tg: int = 91,
enabled: bool = True,
file_number: str = "test-msg",
) -> None:
scenario.config["VOICE"] = {
"ANNOUNCEMENTS": [
{
"ENABLED": enabled,
"TG": tg,
"FILE": file_number,
"LANGUAGE": "en_GB",
"MODE": "interval",
}
]
}
def voice_tts_config(
scenario: DeterministicScenario,
*,
tg: int = 91,
enabled: bool = True,
file_number: str = "tts-msg.ambe",
) -> None:
scenario.config["VOICE"] = {
"TTS_ANNOUNCEMENTS": [
{
"ENABLED": enabled,
"TG": tg,
"FILE": file_number,
"LANGUAGE": "en_GB",
"MODE": "interval",
}
]
}

@ -35,6 +35,7 @@ from tests.routing.unit_data_helpers import idle_hbp_slot
from adn_server.application.plugins.application.bus import PluginBus
from adn_server.application.plugins.application.data_bridge import DataPluginBridge
from adn_server.application.plugins.application.ingress import PluginIngress
from adn_server.application.routing.helpers import hbp_slot_blocks_group_voice
from adn_server.application.plugins.domain.events import UnitDataFrame
from adn_server.application.routing.announcement_ptt_inject import inject_plugin_dmrd
from adn_server.domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_3, bytes_4
@ -170,9 +171,11 @@ def test_the_master_slot_is_held_while_it_plays_and_freed_by_the_terminator() ->
frames = _beacon()
_play(sc, ingress, frames[:3])
slot = sc.protocols["SYSTEM"].STATUS[2]
assert slot["TX_TYPE"] == HBPF_SLT_VHEAD and slot["TX_RFS"] == bytes_3(BEACON_ID)
assert slot["RX_RFS"] == bytes_3(BEACON_ID) and slot["RX_TYPE"] != HBPF_SLT_VTERM
# routed voice for another TG finds the slot busy while the beacon plays
assert hbp_slot_blocks_group_voice(slot, bytes_3(214), b"\x07" * 4, sc.clock.time(), 0)
_play(sc, ingress, frames[3:])
assert slot["TX_TYPE"] == HBPF_SLT_VTERM
assert slot["RX_TYPE"] == HBPF_SLT_VTERM
def test_a_radio_talking_on_the_slot_stops_the_beacon() -> None:
@ -187,3 +190,94 @@ def test_a_second_plugin_stream_cannot_talk_over_the_first() -> None:
sc, ingress, _ = _voice_scenario()
_play(sc, ingress, _beacon(stream=1)[:2])
assert _play(sc, ingress, _beacon(stream=2)[:1]) == [False]
def test_voice_slot_prefers_the_slot_where_the_tg_is_bridged_on_the_master() -> None:
config = minimal_config(("SYSTEM", "SYSTEM-B"))
config["SYSTEMS"]["SYSTEM"]["PEERS"] = {b"\x00\x00\x03\xe9": {"CALLSIGN": "HOTSPOT"}}
table = active_routing_table(TG, (("SYSTEM", 1), ("SYSTEM-B", 2)), timeout_minutes=10**6)
sc = DeterministicScenario(config=config, routing_table=table)
sc.routing.apply_startup_subscriptions()
for ts in (1, 2):
sc.protocols["SYSTEM"].STATUS[ts] = idle_hbp_slot() | {"RX_TYPE": HBPF_SLT_VTERM}
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None, clock=sc.clock.time)
assert ingress.voice_slot_for_tg(TG) == 1 # TG 213 lives on TS1 of this MASTER
assert ingress.voice_slot_for_tg(9999) == 2 # nowhere yet: TS2 first
def test_voice_slot_is_none_while_every_slot_is_busy() -> None:
sc, ingress, _ = _voice_scenario()
for ts in (1, 2):
sc.protocols["SYSTEM"].STATUS[ts] = {
"RX_TYPE": HBPF_SLT_VHEAD, "RX_TIME": sc.clock.time(), "TX_TYPE": HBPF_SLT_VTERM,
"RX_STREAM_ID": b"\x05\x05\x05\x05",
}
assert ingress.voice_slot_for_tg(TG) is None
def test_the_monitor_sees_the_beacon_on_the_master_it_plays_on() -> None:
sc, _, _ = _voice_scenario()
events: list[str] = []
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None,
clock=sc.clock.time, send_routing_event=events.append)
_play(sc, ingress, _beacon(stream=0x0C0C0C0C))
starts = [e for e in events if e.startswith("GROUP VOICE,START,TX,SYSTEM,")]
ends = [e for e in events if e.startswith("GROUP VOICE,END,TX,SYSTEM,")]
assert len(starts) == 1 and len(ends) == 1
assert starts[0].split(",")[4:9] == [str(0x0C0C0C0C), str(BEACON_ID), str(BEACON_ID), "2", str(TG)]
def test_a_beacon_cut_by_a_radio_still_ends_on_the_monitor() -> None:
sc, _, _ = _voice_scenario()
events: list[str] = []
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None,
clock=sc.clock.time, send_routing_event=events.append)
frames = _beacon()
_play(sc, ingress, frames[:3])
sc.protocols["SYSTEM"].STATUS[2].update(RX_TYPE=HBPF_SLT_VHEAD, RX_TIME=sc.clock.time() + 0.06, RX_STREAM_ID=b"\x09" * 4)
assert _play(sc, ingress, frames[3:4]) == [False]
assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1
slot = sc.protocols["SYSTEM"].STATUS[2]
assert slot["RX_STREAM_ID"] == b"\x09" * 4 and slot["RX_TYPE"] == HBPF_SLT_VHEAD # the radio's state is left alone
def test_a_stream_the_plugin_abandons_ends_after_a_silent_second() -> None:
sc, _, _ = _voice_scenario()
events: list[str] = []
ingress = PluginIngress(sc.routing, sc.config, lambda: sc.protocols, lambda *a: None,
clock=sc.clock.time, send_routing_event=events.append)
_play(sc, ingress, _beacon(stream=1)[:3]) # no terminator
sc.clock.advance(1.5)
assert ingress.voice_slot_for_tg(TG) == 2 # the next query sweeps it
assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1
assert sc.protocols["SYSTEM"].STATUS[2]["RX_TYPE"] == HBPF_SLT_VTERM # released by the sweep
def test_a_beacon_on_a_slot_with_stale_rx_state_is_forwarded_whole_with_its_own_lc() -> None:
"""Regression (2131, 25-sep-2026): the ingress slot still held the last radio's RX
state. From the second frame on, routing took the beacon for a colliding QSO,
suppressed the uplink, and forwarded only 2 frames; the LC came from RX_LC."""
from adn_server.domain.dmr import decode
sc, ingress, _ = _voice_scenario()
sc.protocols["SYSTEM"].STATUS[2].update(
RX_TYPE=HBPF_SLT_VTERM, RX_STREAM_ID=b"\x0e" * 4, RX_RFS=bytes_3(3120001), RX_TGID=bytes_3(214),
RX_PEER=bytes_4(1001), RX_TIME=sc.clock.time() - 30, RX_LC=b"\x00\x00\x20" + bytes_3(214) + bytes_3(3120001),
)
frames = _beacon()
assert _play(sc, ingress, frames) == [True] * len(frames)
to_b = sc.capture.for_system("SYSTEM-B")
assert len(to_b) == len(frames)
lc = decode.voice_head_term(to_b[0].packet[20:53])["LC"]
assert lc[3:6] == bytes_3(TG) and lc[6:9] == bytes_3(BEACON_ID)
# the embedded LC of the voice bursts B-E must be the beacon's too, not the stale RX_LC
from bitarray import bitarray
from adn_server.domain.dmr import bptc
by_vseq = {p.packet[15] & 0x0F: p.packet for p in to_b[1:-1]}
frags = bitarray(endian="big")
for vseq in (1, 2, 3, 4):
frags += decode.voice(by_vseq[vseq][20:53])["EMBED"]
emb = bptc.decode_emblc(frags)
assert emb[3:6] == bytes_3(TG) and emb[6:9] == bytes_3(BEACON_ID)

@ -1,4 +1,4 @@
# ADN DMR Peer Server - announcement synthetic PTT inject
# ADN DMR Peer Server - tests plugin group voice through the synthetic ingress
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
@ -18,19 +18,21 @@
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Announcement synthetic PTT inject routing (proxy MASTER + OBP forward)."""
"""Plugin group voice (announcements, beacons) through the synthetic ingress on the proxy MASTER."""
from __future__ import annotations
from tests.harness.deterministic import DeterministicScenario, PacketSpec, patch_routing_wall_time
from tests.harness.scenarios import obp_bridge_scenario
from tests.harness.voice_helpers import FakeMasterForVoice
from adn_server.application.plugins.application.ingress import PluginIngress
from adn_server.application.routing.announcement_ptt_inject import (
announcement_ptt_system,
inject_announcement_ptt,
inject_plugin_dmrd,
)
from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID
from adn_server.domain import bytes_4, int_id
from adn_server.domain import bytes_3, bytes_4, int_id
def test_announcement_ptt_system_prefers_proxy_target() -> None:
@ -57,28 +59,30 @@ def test_inject_forwards_to_obp_and_peer_master() -> None:
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(
inject_plugin_dmrd(
scenario.routing,
"MASTER-A",
vhead.data(),
pkt_time=scenario.clock.time(),
server_id=sid,
plugin="beacon",
)
for seq in range(1, 3):
spec = DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=min(seq, 4))
assert inject_announcement_ptt(
assert inject_plugin_dmrd(
scenario.routing,
"MASTER-A",
spec.data(),
pkt_time=scenario.clock.time(),
server_id=sid,
plugin="beacon",
) is True
assert len(scenario.capture.for_system("OBP-CL")) == 3
assert len(scenario.capture.for_system("MASTER-A")) == 0
def test_inject_does_not_attribute_call_to_real_peer_matching_rf_src() -> None:
"""An announcement whose rf_src (e.g. 1000001) matches a real connected peer's
"""A plugin voice stream (announcement, beacon) whose rf_src (e.g. 1000001) matches a real connected peer's
login id must never be reported as that peer's call: no RX/TX attribution,
no dynamic TG learned for it. Regression for the misattribution bug where
@ -103,21 +107,23 @@ def test_inject_does_not_attribute_call_to_real_peer_matching_rf_src() -> None:
)
with patch_routing_wall_time(scenario.clock):
vhead = DeterministicScenario.voice_head_spec(base)
inject_announcement_ptt(
inject_plugin_dmrd(
scenario.routing,
"MASTER-A",
vhead.data(),
pkt_time=scenario.clock.time(),
server_id=sid,
plugin="beacon",
)
for seq in range(1, 3):
spec = DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=min(seq, 4))
inject_announcement_ptt(
inject_plugin_dmrd(
scenario.routing,
"MASTER-A",
spec.data(),
pkt_time=scenario.clock.time(),
server_id=sid,
plugin="beacon",
)
assert events, "expected at least one GROUP VOICE report"
@ -129,3 +135,56 @@ def test_inject_does_not_attribute_call_to_real_peer_matching_rf_src() -> None:
assert reported_peer != real_peer_id
assert reported_peer == int_id(sid)
assert fields[-1] == "1"
def _ingress(scenario, master: FakeMasterForVoice) -> PluginIngress:
scenario.protocols["MASTER-A"] = master
return PluginIngress(
scenario.routing, scenario.config, lambda: scenario.protocols,
lambda system, pkt: scenario.protocols[system].send_system(pkt), clock=scenario.clock.time,
)
def test_voice_slot_prefers_the_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;"}}
sys_cfg.setdefault("_PEER_UA_MULTI_TGS", {}).setdefault(peer, {})[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}
assert _ingress(scenario, master).voice_slot_for_tg(730600) == 2
def test_voice_slot_prefers_the_active_bridge_slot_for_a_static_tg() -> None:
scenario = obp_bridge_scenario("OBP-CL", tg=730500)
scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
scenario.config["SYSTEMS"]["MASTER-A"]["PEERS"] = {
bytes_4(730039210): {"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}
assert _ingress(scenario, master).voice_slot_for_tg(730500) == 2
def test_plugin_voice_needs_no_pre_armed_bridge_and_reaches_obp() -> None:
"""The synthetic PTT creates the UA relay itself, as announcements always did."""
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}
ingress = _ingress(scenario, master)
base = PacketSpec(rf_src=DEFAULT_SERVER_VOICE_ID, dst_id=730500, slot=2, stream_id=0xA0A0A0A0)
frames = [DeterministicScenario.voice_head_spec(base)] + [
DeterministicScenario.voice_burst_spec(base, seq=n, dtype_vseq=n) for n in (1, 2)
]
with patch_routing_wall_time(scenario.clock):
assert [ingress.deliver(f.data(), "voice-announcements") for f in frames] == [True, True, True]
assert len(scenario.capture.for_system("OBP-CL")) == 3
assert len(master.sent) == 3

@ -1,85 +0,0 @@
# ADN DMR Peer Server - tests voice announcement anticollision
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Voice announcement anti-collision."""
from __future__ import annotations
from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario
from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID
from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3
def test_build_targets_skips_busy_qso_slot() -> None:
scenario, master = voice_master_scenario()
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD
master.STATUS[2]["TX_TYPE"] = HBPF_SLT_VTERM
uc = make_voice_uc(scenario, master)
targets, busy = uc._build_announcement_targets(91, "91", "ANN-TEST")
assert targets == []
assert busy == 1
def test_build_targets_includes_idle_slot() -> None:
scenario, master = voice_master_scenario()
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VTERM
master.STATUS[2]["TX_TYPE"] = HBPF_SLT_VTERM
uc = make_voice_uc(scenario, master)
targets, busy = uc._build_announcement_targets(91, "91", "ANN-TEST")
assert busy == 0
assert len(targets) == 1
assert targets[0]["name"] == "MASTER-A"
assert targets[0]["ts"] == 2
def test_broadcast_aborts_when_qso_starts_mid_transmission() -> None:
scenario, master = voice_master_scenario()
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VTERM
master.STATUS[2]["TX_TYPE"] = HBPF_SLT_VTERM
uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91)
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD
uc._announcement_send_broadcast(targets, pkts, 0, source, dst, 91, 0, "ANN-TEST", None)
assert uc._announcement_running[0] is False
assert master.sent == []
def test_mark_slots_busy_stamps_server_voice_tx_row() -> None:
scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master)
slot = master.STATUS[2]
slot["TX_TYPE"] = HBPF_SLT_VTERM
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": slot, "ts": 2}]
uc._mark_slots_busy(targets)
assert slot["TX_TYPE"] == HBPF_SLT_VHEAD
assert int.from_bytes(slot["TX_RFS"], "big") == DEFAULT_SERVER_VOICE_ID
assert slot["TX_TIME"] > 0

@ -1,74 +0,0 @@
# ADN DMR Peer Server - tests voice broadcast queue
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Voice broadcast queue (same-TG serialization)."""
from __future__ import annotations
from tests.harness.voice_helpers import make_voice_uc, voice_master_scenario
from adn_server.application.server_voice import DEFAULT_SERVER_VOICE_ID
from adn_server.domain import bytes_3
def test_enqueue_broadcast_queues_second_same_tg() -> None:
scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91)
uc._broadcast_active_tgs.add("91")
uc._enqueue_broadcast("ann", targets, pkts, source, dst, 91, 0, "ANN-2")
assert len(uc._broadcast_queue) == 1
assert uc._scheduled == []
def test_broadcast_queue_drains_after_first_finishes() -> None:
scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master)
targets = [{"sys_obj": master, "name": "MASTER-A", "slot": master.STATUS[2], "ts": 2}]
pkts = {1: [b"\x00" * 55], 2: [b"\x00" * 55]}
source = bytes_3(DEFAULT_SERVER_VOICE_ID)
dst = bytes_3(91)
uc._broadcast_active_tgs.add("91")
uc._broadcast_queue.append(
{
"type": "ann",
"targets": targets,
"pkts_by_ts": pkts,
"source_id": source,
"dst_id": dst,
"tg": 91,
"num": 1,
"label": "ANN-2",
}
)
uc._broadcast_finished(91)
assert len(uc._scheduled) == 1
uc._scheduled.pop(0)[1][0]()
assert "91" in uc._broadcast_active_tgs
assert uc._broadcast_queue == []
assert len(uc._scheduled) == 1
assert uc._scheduled[0][0] == 0.5

@ -1,122 +0,0 @@
# 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

@ -1,92 +0,0 @@
# ADN DMR Peer Server - tests voice scheduled announcement
#
# 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
###############################################################################
"""Scheduled file announcements (AMBE on disk)."""
from __future__ import annotations
from tests.harness.voice_helpers import (
drain_call_later,
make_voice_uc,
voice_announcement_config,
voice_master_scenario,
)
from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD, HBPF_SLT_VTERM
def test_scheduled_announcement_starts_broadcast_on_idle_slot() -> None:
scenario, master = voice_master_scenario()
voice_announcement_config(scenario)
uc = make_voice_uc(scenario, master)
uc.scheduled_announcement(0)
assert uc._announcement_running[0] is True
assert "91" in uc._broadcast_active_tgs
assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VHEAD
assert len(uc._scheduled) == 1
delay, _scheduled = uc._scheduled[0]
assert delay == 0.5
drain_call_later(uc)
assert len(master.sent) == 3
assert len(scenario.capture.for_system("MASTER-B")) == 3
def test_scheduled_announcement_retries_when_slot_busy() -> None:
scenario, master = voice_master_scenario()
voice_announcement_config(scenario)
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD
uc = make_voice_uc(scenario, master)
uc.scheduled_announcement(0)
assert uc._announcement_running.get(0) is not True
assert len(uc._scheduled) == 1
assert master.sent == []
delay, (_fn, args) = uc._scheduled[0]
assert delay == 5.0
assert args == (0, 1)
def test_scheduled_announcement_skips_when_disabled() -> None:
scenario, master = voice_master_scenario()
voice_announcement_config(scenario, enabled=False)
uc = make_voice_uc(scenario, master)
uc.scheduled_announcement(0)
assert uc._announcement_running.get(0) is not True
assert uc._scheduled == []
def test_scheduled_announcement_sends_packets_to_master() -> None:
scenario, master = voice_master_scenario()
voice_announcement_config(scenario)
uc = make_voice_uc(scenario, master)
uc.scheduled_announcement(0)
drain_call_later(uc)
assert uc._announcement_running[0] is False
assert "91" not in uc._broadcast_active_tgs
assert len(master.sent) == 3
assert len(scenario.capture.for_system("MASTER-B")) == 3
assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM

@ -1,185 +0,0 @@
# ADN DMR Peer Server - tests voice scheduled tts
#
# 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
###############################################################################
"""Scheduled TTS announcements and conversion callbacks."""
from __future__ import annotations
from datetime import datetime
from unittest.mock import MagicMock, patch
from tests.harness.scenarios import obp_bridge_scenario
from tests.harness.voice_helpers import (
FakeMasterForVoice,
drain_call_later,
make_voice_uc,
voice_master_scenario,
voice_tts_config,
)
from adn_server.application.voice_use_cases import VoiceUseCases
from adn_server.domain.hbp_protocol import HBPF_SLT_VHEAD
def test_scheduled_tts_sync_path_enqueues_broadcast() -> None:
scenario, master = voice_master_scenario()
voice_tts_config(scenario)
uc = make_voice_uc(scenario, master)
uc.scheduled_tts_announcement(0)
assert uc._broadcast_active_tgs == {"91"}
assert len(uc._scheduled) == 1
assert getattr(uc._scheduled[0][1][0], "__name__", "") == "_tts_send_broadcast"
def test_inject_tts_starts_without_active_bridge_and_forwards_obp() -> None:
"""Synthetic PTT must not require a pre-armed bridge row (UA relay created on inject)."""
scenario = obp_bridge_scenario("OBP-CL", tg=730500)
scenario.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"}
scenario.seed_routing_table({})
master = FakeMasterForVoice("MASTER-A")
master.STATUS[2] = {"RX_TYPE": 2, "TX_TYPE": 2, "RX_STREAM_ID": b"\x00" * 4}
scenario.protocols["MASTER-A"] = master
voice_tts_config(scenario, tg=730500, file_number="730500")
uc = make_voice_uc(scenario, master)
uc._tts_conversion_done("/tmp/fake.ambe", 0, "730500", 730500, "en_GB", "interval", "TTS-1")
assert uc._broadcast_active_tgs == {"730500"}
drain_call_later(uc)
assert len(scenario.capture.for_system("OBP-CL")) == 3
assert len(master.sent) == 3
def test_tts_conversion_error_clears_running_flag() -> None:
scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master)
uc._tts_running[0] = True
uc._tts_conversion_error(RuntimeError("TTS failed"), 0, "TTS-1")
assert uc._tts_running[0] is False
def test_tts_conversion_done_without_ambe_clears_running() -> None:
scenario, master = voice_master_scenario()
uc = make_voice_uc(scenario, master)
uc._tts_running[0] = True
uc._tts_conversion_done(None, 0, "msg.ambe", 91, "en_GB", "interval", "TTS-1")
assert uc._tts_running[0] is False
assert uc._scheduled == []
def test_tts_conversion_done_retries_when_slot_busy() -> None:
scenario, master = voice_master_scenario()
scenario.seed_routing_table({})
master.STATUS[1]["RX_TYPE"] = HBPF_SLT_VHEAD
master.STATUS[2]["RX_TYPE"] = HBPF_SLT_VHEAD
uc = make_voice_uc(scenario, master)
uc._tts_running[0] = True
uc._tts_conversion_done("/tmp/fake.ambe", 0, "msg.ambe", 91, "en_GB", "interval", "TTS-1")
assert uc._tts_running[0] is True
assert len(uc._scheduled) == 1
delay, (fn, args) = uc._scheduled[0]
assert delay == 5.0
assert getattr(fn, "__name__", "") == "_tts_conversion_done"
assert args[-1] == 1
def test_scheduled_tts_skips_outside_top_of_hour() -> None:
scenario, master = voice_master_scenario()
scenario.config["VOICE"] = {
"TTS_ANNOUNCEMENTS": [
{
"ENABLED": True,
"TG": 91,
"FILE": "hourly.ambe",
"LANGUAGE": "en_GB",
"MODE": "hourly",
}
]
}
uc = make_voice_uc(scenario, master)
with patch("adn_server.application.voice_use_cases.datetime") as mock_dt:
mock_dt.now.return_value = datetime(2026, 5, 24, 10, 30)
uc.scheduled_tts_announcement(0)
assert uc._tts_running.get(0) is not True
assert uc._scheduled == []
def test_scheduled_tts_defers_when_same_tg_already_broadcasting() -> None:
scenario, master = voice_master_scenario()
voice_tts_config(scenario)
uc = make_voice_uc(scenario, master)
uc._broadcast_active_tgs.add("91")
uc.scheduled_tts_announcement(0)
assert uc._tts_running.get(0) is not True
assert len(uc._scheduled) == 1
delay, (fn, args) = uc._scheduled[0]
assert delay == 3.0
assert getattr(fn, "__name__", "") == "scheduled_tts_announcement"
assert args == (0, 1)
def test_scheduled_tts_sync_exception_clears_running() -> None:
scenario, master = voice_master_scenario()
voice_tts_config(scenario)
class FailingProvider:
def ensure_tts_ambe(self, config, item, audio_path):
del config, item, audio_path
raise OSError("disk full")
def read_single_file(self, *args):
return [b"\x00" * 7]
def pkt_gen(self, *args, **kwargs):
return iter([])
def get_ambe_words(self, *args):
return {}
scheduled: list[tuple[float, tuple]] = []
def call_later(delay, fn, *args):
scheduled.append((delay, (fn, args)))
return MagicMock()
uc = VoiceUseCases(
FailingProvider(),
scenario.config,
get_protocols=lambda: {"MASTER-A": master},
routing_table_for_report=scenario.routing.routing_table_for_report,
call_later=call_later,
audio_path="/tmp/audio",
)
uc.scheduled_tts_announcement(0)
assert uc._tts_running[0] is False

@ -1,138 +0,0 @@
# ADN DMR Peer Server - tests voice voice config reload
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation; either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program; if not, write to the Free Software Foundation,
# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA
###############################################################################
"""Voice config reload and announcement LoopingCall management."""
from __future__ import annotations
from unittest.mock import MagicMock
from tests.harness.voice_helpers import FakeVoiceProvider, voice_master_scenario
from adn_server.application.voice_use_cases import VoiceUseCases
def _reload_uc(scenario, *, start_looping_call) -> VoiceUseCases:
return VoiceUseCases(
FakeVoiceProvider(),
scenario.config,
start_looping_call=start_looping_call,
audio_path="/tmp/audio",
)
def test_apply_voice_config_starts_enabled_announcement_loop() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {
"ANNOUNCEMENTS": [
{
"ENABLED": True,
"TG": 91,
"FILE": "welcome",
"LANGUAGE": "en_GB",
"MODE": "interval",
"INTERVAL": 120,
}
],
}
started: list[tuple[float, object]] = []
def start_looping_call(fn, interval, _now):
started.append((interval, fn))
handle = MagicMock(running=True)
handle.stop = MagicMock()
return handle
uc = _reload_uc(scenario, start_looping_call=start_looping_call)
uc.apply_voice_config()
assert len(started) == 1
assert started[0][0] == 120.0
assert 0 in uc._ann_tasks
def test_apply_voice_config_stops_removed_announcement() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {"ANNOUNCEMENTS": [{"ENABLED": False, "TG": 91, "FILE": "x"}]}
stop_mock = MagicMock()
uc = VoiceUseCases(
FakeVoiceProvider(),
scenario.config,
start_looping_call=lambda *_a: MagicMock(running=True, stop=stop_mock),
audio_path="/tmp/audio",
)
uc._ann_tasks[0] = MagicMock(running=True, stop=stop_mock)
uc.apply_voice_config()
assert 0 not in uc._ann_tasks
stop_mock.assert_called_once()
def test_apply_voice_config_starts_tts_loop() -> None:
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {
"TTS_ANNOUNCEMENTS": [
{
"ENABLED": True,
"TG": 91,
"FILE": "hourly.ambe",
"LANGUAGE": "en_GB",
"MODE": "hourly",
}
],
}
started: list[float] = []
def start_looping_call(_fn, interval, _now):
started.append(interval)
return MagicMock(running=True, stop=MagicMock())
uc = _reload_uc(scenario, start_looping_call=start_looping_call)
uc.apply_voice_config()
assert started == [30.0]
assert 0 in uc._tts_tasks
def test_apply_voice_config_legacy_yaml_without_dmr_id() -> None:
"""Legacy adn-voice.yaml without DMR_ID must reload without error."""
scenario, _ = voice_master_scenario()
scenario.config["VOICE"] = {
"TTS_ANNOUNCEMENTS": [
{
"ENABLED": True,
"TG": 730500,
"FILE": "730500",
"LANGUAGE": "es_ES",
"MODE": "interval",
"INTERVAL": 60,
}
],
}
uc = VoiceUseCases(
FakeVoiceProvider(),
scenario.config,
start_looping_call=lambda *_a: MagicMock(running=True, stop=MagicMock()),
audio_path="/tmp/audio",
)
uc.apply_voice_config()
assert 0 in uc._tts_tasks
Loading…
Cancel
Save

Powered by TurnKey Linux.