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 <noreply@anthropic.com>
pull/106/head
yo 4 days ago
parent d76de25cbc
commit 620a2b7b43

5
.gitignore vendored

@ -36,7 +36,7 @@ json/*
# Internal # Internal
docs-priv/ 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/*
!plugins/README.md !plugins/README.md
!plugins/.gitkeep !plugins/.gitkeep
@ -44,6 +44,9 @@ plugins/*
!plugins/example/** !plugins/example/**
plugins/example/example-events/ plugins/example/example-events/
plugins/example/**/__pycache__/ plugins/example/**/__pycache__/
!plugins/voice-announcements/
!plugins/voice-announcements/**
plugins/voice-announcements/**/__pycache__/
# Runtime TTS on-demand cache (generated at runtime) # Runtime TTS on-demand cache (generated at runtime)
Audio/es_ES/ondemand/ Audio/es_ES/ondemand/

@ -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

@ -0,0 +1,29 @@
# ADN DMR Peer Server plugin - voice-announcements factory
#
# 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-announcements plugin — drop-in factory."""
from __future__ import annotations
from .application.plugin_impl import VoiceAnnouncementsPlugin
def create_plugin() -> VoiceAnnouncementsPlugin:
return VoiceAnnouncementsPlugin()

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

@ -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,226 @@
# 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()).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

@ -34,6 +34,8 @@ from ..domain.send import GROUP_VOICE, UNIT_DATA, parse_dmrd_header, plugin_fram
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
_SILENT_STREAM_S = 1.0
class PluginIngress: class PluginIngress:
"""Routes plugin frames on the reactor thread and reports whether each was accepted. """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._send_local = send_local
self._clock = clock self._clock = clock
self._send_routing_event = send_routing_event self._send_routing_event = send_routing_event
# Plugin voice streams on air: stream_id -> (master, slot, tg, rf_src, start). # Plugin voice streams on air: stream_id -> [master, slot, tg, rf_src, start, last frame].
self._on_air: dict[bytes, tuple[str, int, int, int, float]] = {} self._on_air: dict[bytes, list[Any]] = {}
def set_routing(self, routing: Any) -> None: def set_routing(self, routing: Any) -> None:
"""Bootstrap builds the plugin manager before routing; routing is set before any load.""" """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 "", {}) sys_cfg = self._config.get("SYSTEMS", {}).get(master or "", {})
if not status or sys_cfg.get("MODE") != "MASTER": if not status or sys_cfg.get("MODE") != "MASTER":
return None return None
now = self._clock()
self._end_silent_streams(now)
dynamic = master_dynamic_tg_slots(sys_cfg, int(tg)) dynamic = master_dynamic_tg_slots(sys_cfg, int(tg))
bridged = self._active_bridge_slots(int(tg), master) 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]): for ts in dict.fromkeys([*sorted(dynamic, reverse=True), *sorted(bridged, reverse=True), 2, 1]):
slot = status.get(ts) slot = status.get(ts)
if slot and not self._slot_busy(slot, ts in dynamic or ts in bridged, now): if slot and not self._slot_busy(slot, ts in dynamic or ts in bridged, now):
@ -118,6 +121,7 @@ class PluginIngress:
return False return False
server_id = self._server_id() server_id = self._server_id()
now = self._clock() now = self._clock()
self._end_silent_streams(now)
if kind == UNIT_DATA: if kind == UNIT_DATA:
return self._route(master, pkt, now, server_id, plugin) is not False 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) 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) self._send_local(master, pkt)
if header.stream_id not in self._on_air: if header.stream_id not in self._on_air:
self._on_air_start(master, header, now) self._on_air_start(master, header, now)
self._on_air[header.stream_id][5] = now
if is_term: if is_term:
self._off_air(header.stream_id, now) self._off_air(header.stream_id, now)
return True return True
def _on_air_start(self, master: str, header: Any, now: float) -> None: 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.""" """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) 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) 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: def _off_air(self, stream_id: bytes, now: float) -> None:
entry = self._on_air.pop(stream_id, None) entry = self._on_air.pop(stream_id, None)
if entry is not 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) self._report("END", master, stream_id, slot, tg, rf_src, now - start)
def _report( def _report(

@ -34,6 +34,8 @@ PLUGIN_SENDABLE_DTYPES = frozenset({3, 6, 7, 8})
UNIT_DATA = "unit_data" UNIT_DATA = "unit_data"
GROUP_VOICE = "group_voice" GROUP_VOICE = "group_voice"
DEFAULT_MAX_FRAMES_PER_S = 40.0 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): class DmrdHeader(NamedTuple):
@ -106,8 +108,9 @@ def send_permission(server_config: dict[str, Any], plugin: str) -> SendPermissio
return None return None
try: try:
ids = frozenset(int(i) for i in entry.get("allowed_src_ids") or ()) 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 ()) 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): except (TypeError, ValueError):
return None return None
if not ids or not math.isfinite(rate) or rate <= 0: if not ids or not math.isfinite(rate) or rate <= 0:

@ -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(213) == 2
assert sender.voice_slot_for_tg(214) is None 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

@ -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 _play(sc, ingress, frames[3:4]) == [False]
assert sum(e.startswith("GROUP VOICE,END,TX,SYSTEM,") for e in events) == 1 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 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

Loading…
Cancel
Save

Powered by TurnKey Linux.