From 620a2b7b43d0927667d121a5a388c9ff4073c7c5 Mon Sep 17 00:00:00 2001 From: yo Date: Thu, 24 Sep 2026 23:47:25 +0200 Subject: [PATCH] feat(plugins): voice-announcements plugin (scheduled, TTS, beacons) The scheduled announcements and TTS announcements of VoiceUseCases, as an official plugin that speaks through send_dmrd. Same configuration (VOICE.ANNOUNCEMENTS / TTS_ANNOUNCEMENTS in adn-voice.yaml, followed every 15 s) and the same behaviour: interval or hourly, one playback per TG at a time 1.5 s apart, retries while the TG or every slot is busy, 58 ms per frame, stop on a refused frame, TTS encoded in a thread. Versioned next to plugins/example. Core side, found by an end-to-end run of the plugin: - PluginIngress ends a plugin voice stream that has been silent for a second without a terminator (frees the slot, monitor END); - the default max_frames_per_s adds 18 frames/s per group voice TG, so one stream per TG fits (a voice stream is ~17 frames/s). Co-Authored-By: Claude Opus 5.5 --- .gitignore | 5 +- plugins/voice-announcements/README.md | 23 ++ plugins/voice-announcements/config.yaml | 7 + .../voice-announcements/plugin/__init__.py | 29 +++ .../plugin/application/__init__.py | 0 .../plugin/application/announcer.py | 204 ++++++++++++++++ .../plugin/application/plugin_impl.py | 80 +++++++ .../plugin/domain/__init__.py | 0 .../plugin/domain/schedule.py | 101 ++++++++ .../plugin/domain/tg_queue.py | 57 +++++ plugins/voice-announcements/tests/conftest.py | 28 +++ .../tests/test_announcer.py | 226 ++++++++++++++++++ .../plugins/application/ingress.py | 29 ++- .../application/plugins/domain/send.py | 5 +- tests/application/test_plugin_sender.py | 9 + tests/routing/test_plugin_send_routing.py | 12 + 16 files changed, 805 insertions(+), 10 deletions(-) create mode 100644 plugins/voice-announcements/README.md create mode 100644 plugins/voice-announcements/config.yaml create mode 100644 plugins/voice-announcements/plugin/__init__.py create mode 100644 plugins/voice-announcements/plugin/application/__init__.py create mode 100644 plugins/voice-announcements/plugin/application/announcer.py create mode 100644 plugins/voice-announcements/plugin/application/plugin_impl.py create mode 100644 plugins/voice-announcements/plugin/domain/__init__.py create mode 100644 plugins/voice-announcements/plugin/domain/schedule.py create mode 100644 plugins/voice-announcements/plugin/domain/tg_queue.py create mode 100644 plugins/voice-announcements/tests/conftest.py create mode 100644 plugins/voice-announcements/tests/test_announcer.py diff --git a/.gitignore b/.gitignore index 1620cb1..27337c9 100644 --- a/.gitignore +++ b/.gitignore @@ -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/ diff --git a/plugins/voice-announcements/README.md b/plugins/voice-announcements/README.md new file mode 100644 index 0000000..9233904 --- /dev/null +++ b/plugins/voice-announcements/README.md @@ -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 +``` diff --git a/plugins/voice-announcements/config.yaml b/plugins/voice-announcements/config.yaml new file mode 100644 index 0000000..59c982d --- /dev/null +++ b/plugins/voice-announcements/config.yaml @@ -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 diff --git a/plugins/voice-announcements/plugin/__init__.py b/plugins/voice-announcements/plugin/__init__.py new file mode 100644 index 0000000..bb5b9c7 --- /dev/null +++ b/plugins/voice-announcements/plugin/__init__.py @@ -0,0 +1,29 @@ +# ADN DMR Peer Server plugin - voice-announcements factory +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""voice-announcements plugin — drop-in factory.""" + +from __future__ import annotations + +from .application.plugin_impl import VoiceAnnouncementsPlugin + + +def create_plugin() -> VoiceAnnouncementsPlugin: + return VoiceAnnouncementsPlugin() diff --git a/plugins/voice-announcements/plugin/application/__init__.py b/plugins/voice-announcements/plugin/application/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/plugins/voice-announcements/plugin/application/announcer.py b/plugins/voice-announcements/plugin/application/announcer.py new file mode 100644 index 0000000..4ca380a --- /dev/null +++ b/plugins/voice-announcements/plugin/application/announcer.py @@ -0,0 +1,204 @@ +# ADN DMR Peer Server plugin - voice-announcements announcer +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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, bytes_4 + +from ..domain.schedule import FILE, Item, enabled_items, hourly_due +from ..domain.tg_queue import TalkgroupQueue + +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, now: Callable[[], datetime] = datetime.now) -> None: + self._ctx = ctx + self._voice = voice + 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._voice.ensure_tts_ambe, 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: + server_id = self._ctx.config.get("GLOBAL", {}).get("SERVER_ID", 0) + if isinstance(server_id, bytes): + return server_id[:4].rjust(4, b"\x00") + return bytes_4(int(server_id or 0) & 0xFFFFFFFF) diff --git a/plugins/voice-announcements/plugin/application/plugin_impl.py b/plugins/voice-announcements/plugin/application/plugin_impl.py new file mode 100644 index 0000000..fd34a33 --- /dev/null +++ b/plugins/voice-announcements/plugin/application/plugin_impl.py @@ -0,0 +1,80 @@ +# ADN DMR Peer Server plugin - voice-announcements adapter +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""ServerPlugin adapter: 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 diff --git a/plugins/voice-announcements/plugin/domain/__init__.py b/plugins/voice-announcements/plugin/domain/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/plugins/voice-announcements/plugin/domain/schedule.py b/plugins/voice-announcements/plugin/domain/schedule.py new file mode 100644 index 0000000..5dd9a8a --- /dev/null +++ b/plugins/voice-announcements/plugin/domain/schedule.py @@ -0,0 +1,101 @@ +# ADN DMR Peer Server plugin - voice-announcements schedule +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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 diff --git a/plugins/voice-announcements/plugin/domain/tg_queue.py b/plugins/voice-announcements/plugin/domain/tg_queue.py new file mode 100644 index 0000000..d38c255 --- /dev/null +++ b/plugins/voice-announcements/plugin/domain/tg_queue.py @@ -0,0 +1,57 @@ +# ADN DMR Peer Server plugin - voice-announcements talkgroup queue +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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) diff --git a/plugins/voice-announcements/tests/conftest.py b/plugins/voice-announcements/tests/conftest.py new file mode 100644 index 0000000..a9c1590 --- /dev/null +++ b/plugins/voice-announcements/tests/conftest.py @@ -0,0 +1,28 @@ +# ADN DMR Peer Server plugin - example test path +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +from __future__ import annotations + +import sys +from pathlib import Path + +_ROOT = Path(__file__).resolve().parents[1] +if str(_ROOT) not in sys.path: + sys.path.insert(0, str(_ROOT)) diff --git a/plugins/voice-announcements/tests/test_announcer.py b/plugins/voice-announcements/tests/test_announcer.py new file mode 100644 index 0000000..e2b5780 --- /dev/null +++ b/plugins/voice-announcements/tests/test_announcer.py @@ -0,0 +1,226 @@ +# ADN DMR Peer Server plugin - voice-announcements tests +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""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()).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 diff --git a/src/adn_server/application/plugins/application/ingress.py b/src/adn_server/application/plugins/application/ingress.py index 01878fe..2d9e02c 100644 --- a/src/adn_server/application/plugins/application/ingress.py +++ b/src/adn_server/application/plugins/application/ingress.py @@ -34,6 +34,8 @@ from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_fram logger = logging.getLogger(__name__) +_SILENT_STREAM_S = 1.0 + class PluginIngress: """Routes plugin frames on the reactor thread and reports whether each was accepted. @@ -62,8 +64,8 @@ class PluginIngress: 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). - self._on_air: dict[bytes, tuple[str, int, int, int, float]] = {} + # 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.""" @@ -81,9 +83,10 @@ class PluginIngress: 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) - now = self._clock() 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): @@ -118,6 +121,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) @@ -142,23 +146,32 @@ class PluginIngress: 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.""" - if len(self._on_air) >= 64: # plugins that never sent a terminator - for stream_id in list(self._on_air)[:32]: - self._off_air(stream_id, now) 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) + 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 = entry + master, slot, tg, rf_src, start, _last = entry self._report("END", master, stream_id, slot, tg, rf_src, now - start) def _report( diff --git a/src/adn_server/application/plugins/domain/send.py b/src/adn_server/application/plugins/domain/send.py index 383ab1e..2b2ea66 100644 --- a/src/adn_server/application/plugins/domain/send.py +++ b/src/adn_server/application/plugins/domain/send.py @@ -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): @@ -106,8 +108,9 @@ def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermissio 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: diff --git a/tests/application/test_plugin_sender.py b/tests/application/test_plugin_sender.py index 011c544..11abe5a 100644 --- a/tests/application/test_plugin_sender.py +++ b/tests/application/test_plugin_sender.py @@ -259,3 +259,12 @@ def test_voice_slot_query_only_for_granted_talkgroups() -> None: ) 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 diff --git a/tests/routing/test_plugin_send_routing.py b/tests/routing/test_plugin_send_routing.py index e751acf..cf68e3a 100644 --- a/tests/routing/test_plugin_send_routing.py +++ b/tests/routing/test_plugin_send_routing.py @@ -235,3 +235,15 @@ def test_a_beacon_cut_by_a_radio_still_ends_on_the_monitor() -> None: assert _play(sc, ingress, frames[3:4]) == [False] assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1 assert sc.protocols["SYSTEM"].STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM + + +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]["TX_TYPE"] == HBPF_SLT_VTERM