Merge pull request #102 from pyopower/fix/threaded-voice-slot-busy

fix(voice): server prompts hold TS2 while they play
pull/107/head
ce5rpy 4 days ago committed by GitHub
commit 114ad296f5
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

@ -153,13 +153,6 @@ class IdentUseCases:
speech = self._voice.pkt_gen(_source_id, _dst_id, server_id, 1, _say)
time.sleep(1)
_slot = protocol.STATUS.get(2)
if not _slot:
if not protocol.STATUS.get(2):
continue
_next_time = time.time()
for pkt in speech:
_next_time += 0.058
_delay = _next_time - time.time()
if _delay > 0.001:
time.sleep(_delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, _dst_id, _slot)
self._voice.play_on_slot(protocol, system, speech, _source_id, _dst_id)

@ -112,17 +112,24 @@ def peer_is_simplex(peer: dict[str, Any]) -> bool:
return peer_rf_mode(peer) == RF_MODE_SIMPLEX
def _slot_leg_active(slot_st: dict[str, Any], leg: str, pkt_time: float) -> bool:
"""True when the slot's ``RX`` or ``TX`` leg carries voice (within STREAM_TO)."""
leg_type = slot_st.get(f"{leg}_TYPE")
if leg_type is None or leg_type == HBPF_SLT_VTERM:
return False
return (pkt_time - float(slot_st.get(f"{leg}_TIME", 0) or 0)) < STREAM_TO
def slot_has_active_voice(slot_st: dict[str, Any], pkt_time: float) -> bool:
"""True when the slot has an active group-voice RX or TX leg (within STREAM_TO)."""
rx_type = slot_st.get("RX_TYPE")
if rx_type is not None and rx_type != HBPF_SLT_VTERM:
if (pkt_time - float(slot_st.get("RX_TIME", 0))) < STREAM_TO:
return True
tx_type = slot_st.get("TX_TYPE")
if tx_type is not None and tx_type != HBPF_SLT_VTERM:
if (pkt_time - float(slot_st.get("TX_TIME", 0))) < STREAM_TO:
return True
return False
return _slot_leg_active(slot_st, "RX", pkt_time) or _slot_leg_active(slot_st, "TX", pkt_time)
def slot_voice_held_by_other_stream(slot_st: dict[str, Any], stream_id: bytes, pkt_time: float) -> bool:
"""Like ``slot_has_active_voice``, but ``stream_id``'s own TX leg does not count."""
if _slot_leg_active(slot_st, "RX", pkt_time):
return True
return slot_st.get("TX_STREAM_ID") != stream_id and _slot_leg_active(slot_st, "TX", pkt_time)
def _slot_last_voice_activity(slot_st: dict[str, Any]) -> tuple[bytes, float]:

@ -28,11 +28,13 @@ 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 .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,
@ -45,6 +47,18 @@ _ANNOUNCEMENT_EXCLUDED = ("ECHO", "D-APRS")
_BROADCAST_GAP = 1.5
@dataclass
class _PromptRun:
"""State of one prompt, shared by its worker thread and the reactor.
Only the reactor writes it; the thread reads ``stopped`` between frames.
"""
stopped: bool = False
sent: int = 0
stream_id: bytes | None = None
class VoiceUseCases:
"""Use cases for voice announcements and TTS."""
@ -889,6 +903,60 @@ class VoiceUseCases:
"""Read one AMBE file (e.g. ondemand/{file_number}.ambe). Legacy readSingleFile."""
return self._voice.read_single_file(audio_path, lang, file_number)
def play_on_slot(
self, protocol: Any, system: str, speech: Any, source_id: bytes, dst_id: bytes
) -> int:
"""Play a prompt on TS2 of an HBP system from a worker thread; frames sent.
The thread only paces the frames: each one is sent on the reactor, which
holds the slot while the prompt plays the way scheduled broadcasts do
(TX_TYPE=VHEAD), so routed group voice finds it busy instead of going out
as a second stream on the same slot. A radio keying up, or a call already
on the slot, stops the prompt; the slot is released when it ends.
"""
run = _PromptRun()
_next_time = time.time()
for pkt in speech:
if run.stopped:
break
_next_time += _FRAME_INTERVAL
delay = _next_time - time.time()
if delay > 0.001:
time.sleep(delay)
self._call_from_reactor(self._prompt_frame, protocol, system, pkt, source_id, dst_id, run)
self._call_from_reactor(self._prompt_end, protocol, run)
return run.sent
def _prompt_frame(
self, protocol: Any, system: str, pkt: bytes, source_id: bytes, dst_id: bytes, run: _PromptRun
) -> None:
"""Reactor side of ``play_on_slot``: one frame, unless the slot is taken."""
if run.stopped:
return
slot = protocol.STATUS.get(2) if getattr(protocol, "STATUS", None) else None
if not slot:
run.stopped = True
return
stream_id = pkt[16:20]
now = time.time()
if slot_voice_held_by_other_stream(slot, stream_id, now):
run.stopped = True
logger.info("(%s) Voice on TS2, stopping server prompt after %s frames", system, run.sent)
return
slot["TX_TYPE"] = HBPF_SLT_VHEAD
slot["TX_STREAM_ID"] = stream_id
slot["TX_RFS"] = source_id
slot["TX_TIME"] = now
run.stream_id = stream_id
protocol.send_voice_packet(pkt, source_id, dst_id, slot)
run.sent += 1
def _prompt_end(self, protocol: Any, run: _PromptRun) -> None:
"""Reactor side of ``play_on_slot``: free the slot if the prompt still holds it."""
slot = protocol.STATUS.get(2) if getattr(protocol, "STATUS", None) else None
if slot and run.stream_id is not None and slot.get("TX_STREAM_ID") == run.stream_id:
slot["TX_TYPE"] = HBPF_SLT_VTERM
def play_file_on_request(self, file_number: str, system: str) -> None:
"""Play AMBE file on request (legacy playFileOnRequest). TG 9991-9999 triggers this."""
if not self._get_protocols or not self._call_from_reactor or not self._audio_path:
@ -908,19 +976,9 @@ class VoiceUseCases:
_source_id = self._server_source_id()
speech = self.pkt_gen(_source_id, bytes_3(9), bytes_4(9), 1, _say)
time.sleep(1)
_slot = protocol.STATUS.get(2)
if not _slot:
if not protocol.STATUS.get(2):
return
_dst_id = bytes_3(9)
_next_time = time.time()
_pkt_count = 0
for pkt in speech:
_next_time += 0.058
delay = _next_time - time.time()
if delay > 0.001:
time.sleep(delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, _dst_id, _slot)
_pkt_count += 1
_pkt_count = self.play_on_slot(protocol, system, speech, _source_id, bytes_3(9))
logger.info("(%s) On-demand playback complete: %s (%d packets)", system, file_number, _pkt_count)
def disconnected_voice(self, system: str) -> None:
@ -957,17 +1015,10 @@ class VoiceUseCases:
_source_id = self._server_source_id()
speech = self.pkt_gen(_source_id, bytes_3(9), bytes_4(9), 1, _say)
time.sleep(1)
_slot = protocol.STATUS.get(2)
if not _slot:
if not protocol.STATUS.get(2):
return
logger.debug("(%s) Sending disconnected voice", system)
_next_time = time.time()
for pkt in speech:
_next_time += 0.058
_delay = _next_time - time.time()
if _delay > 0.001:
time.sleep(_delay)
self._call_from_reactor(protocol.send_voice_packet, pkt, _source_id, bytes_3(9), _slot)
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:

@ -0,0 +1,131 @@
# ADN DMR Peer Server - tests server prompts hold the slot they play on
#
# 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
###############################################################################
"""Ident, on-demand and disconnected prompts hold TS2 while they play, like broadcasts do."""
from __future__ import annotations
import time
from unittest.mock import MagicMock, patch
from tests.harness.voice_helpers import FakeMasterForVoice, FakeVoiceProvider, voice_master_scenario
from adn_server.application.routing.helpers import hbp_slot_blocks_group_voice
from adn_server.application.voice_use_cases import VoiceUseCases
from adn_server.domain import HBPF_SLT_VHEAD, HBPF_SLT_VTERM, bytes_3
PROMPT_STREAM = b"\xaa\xbb\xcc\xdd"
OTHER_STREAM = b"\x01\x02\x03\x04"
FRAMES = 6
def _speech(n: int = FRAMES) -> list[bytes]:
return [b"DMRD" + b"\x00" * 12 + PROMPT_STREAM + bytes([i]) * 33 for i in range(n)]
def _idle_slot() -> dict:
return {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM, "RX_TIME": 0.0, "TX_TIME": 0.0}
class _Master(FakeMasterForVoice):
"""Records, per frame sent, whether routing would let another TG onto the slot."""
def __init__(self, name: str, on_frame=None) -> None:
super().__init__(name)
self.voice_packets: list[bytes] = []
self.other_tg_admitted: list[bool] = []
self._on_frame = on_frame
def send_voice_packet(self, packet: bytes, _source_id: bytes, _dst_id: bytes, slot: dict) -> None:
self.voice_packets.append(packet)
blocked = hbp_slot_blocks_group_voice(slot, bytes_3(214), OTHER_STREAM, time.time(), 0)
self.other_tg_admitted.append(not blocked)
if self._on_frame:
self._on_frame(len(self.voice_packets), slot)
def _uc(master: _Master) -> VoiceUseCases:
scenario, _ = voice_master_scenario()
return VoiceUseCases(
FakeVoiceProvider(),
scenario.config,
get_protocols=lambda: {"MASTER-A": master},
call_from_reactor=lambda fn, *args: fn(*args),
audio_path="/tmp/audio",
)
def _play(uc: VoiceUseCases, master: _Master) -> int:
with patch("adn_server.application.voice_use_cases.time") as mock_time:
mock_time.time.side_effect = time.time
mock_time.sleep = MagicMock()
return uc.play_on_slot(master, "MASTER-A", _speech(), bytes_3(9990), bytes_3(9))
def test_routed_group_voice_is_kept_off_the_slot_while_a_prompt_plays() -> None:
master = _Master("MASTER-A")
master.STATUS[2] = _idle_slot()
uc = _uc(master)
with patch("adn_server.application.voice_use_cases.time") as mock_time:
mock_time.time.side_effect = time.time
mock_time.sleep = MagicMock()
uc.play_file_on_request("9991", "MASTER-A")
assert master.voice_packets
# With GROUP_HANGTIME 0 only the slot state keeps a second stream off TS2: the
# prompt used to leave it at VTERM, so routing admitted another TG on every frame.
assert not any(master.other_tg_admitted)
def test_the_slot_is_released_when_the_prompt_ends() -> None:
master = _Master("MASTER-A")
master.STATUS[2] = _idle_slot()
_play(_uc(master), master)
assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM
def test_a_radio_keying_up_stops_the_prompt() -> None:
def key_up(n: int, slot: dict) -> None:
if n == 2:
slot["RX_TYPE"] = HBPF_SLT_VHEAD
slot["RX_TIME"] = time.time()
master = _Master("MASTER-A", on_frame=key_up)
master.STATUS[2] = _idle_slot()
sent = _play(_uc(master), master)
assert sent == 2
assert master.STATUS[2]["TX_TYPE"] == HBPF_SLT_VTERM
def test_a_call_already_on_the_slot_is_left_alone() -> None:
master = _Master("MASTER-A")
routed = {"TX_TYPE": HBPF_SLT_VHEAD, "TX_STREAM_ID": OTHER_STREAM, "TX_TIME": time.time()}
master.STATUS[2] = {**_idle_slot(), **routed}
sent = _play(_uc(master), master)
assert sent == 0
assert master.voice_packets == []
assert {k: master.STATUS[2][k] for k in routed} == routed
Loading…
Cancel
Save

Powered by TurnKey Linux.