diff --git a/src/adn_server/application/playback_use_cases.py b/src/adn_server/application/playback_use_cases.py index 8b8d769..2ca6255 100644 --- a/src/adn_server/application/playback_use_cases.py +++ b/src/adn_server/application/playback_use_cases.py @@ -26,13 +26,11 @@ from __future__ import annotations import logging +from collections.abc import Callable from random import randint from time import time from typing import Any -from twisted.internet import reactor -from twisted.internet.base import DelayedCall - from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, HBPF_VOICE, bytes_4, int_id logger = logging.getLogger(__name__) @@ -49,8 +47,15 @@ _SOURCE_MAX_S = 180.0 class PlaybackUseCases: """Legacy playback class behaviour; playback is scheduled on the reactor (non-blocking).""" - def __init__(self, system_name: str, get_protocol: Any = None) -> None: + def __init__( + self, + system_name: str, + *, + call_later: Callable[..., Any], + get_protocol: Any = None, + ) -> None: self._system = system_name + self._call_later = call_later self._get_protocol = get_protocol self.STATUS: dict[str, Any] = {} self.CALL_DATA: list[bytes] = [] @@ -66,10 +71,10 @@ class PlaybackUseCases: self._playback_index = 0 self._playback_stream_id = b"" self._playback_ta_from_stream = b"" - self._delay_call: DelayedCall | None = None - self._packet_call: DelayedCall | None = None - self._idle_call: DelayedCall | None = None - self._max_call: DelayedCall | None = None + self._delay_call: Any = None + self._packet_call: Any = None + self._idle_call: Any = None + self._max_call: Any = None def dmrd_received( self, @@ -219,7 +224,7 @@ class PlaybackUseCases: self._cancel_record_idle() if not self._recording_active: return - self._idle_call = reactor.callLater(_RECORD_IDLE_S, self._on_record_idle, proto) + self._idle_call = self._call_later(_RECORD_IDLE_S, self._on_record_idle, proto) def _cancel_record_idle(self) -> None: if self._idle_call is not None and self._idle_call.active(): @@ -277,7 +282,7 @@ class PlaybackUseCases: if not recorded: return self._playback_busy = True - self._delay_call = reactor.callLater( + self._delay_call = self._call_later( _PLAYBACK_DELAY_S, self._start_playback, proto, @@ -327,7 +332,7 @@ class PlaybackUseCases: self._cancel_record_max() if not self._recording_active: return - self._max_call = reactor.callLater(_SOURCE_MAX_S, self._on_record_max_duration, proto) + self._max_call = self._call_later(_SOURCE_MAX_S, self._on_record_max_duration, proto) def _cancel_record_max(self) -> None: if self._max_call is not None and self._max_call.active(): @@ -454,7 +459,7 @@ class PlaybackUseCases: else: logger.warning("(%s) Playback packet %d dropped (no protocol/send_system)", self._system, self._playback_index) self._playback_index += 1 - self._packet_call = reactor.callLater(_PACKET_INTERVAL_S, self._send_next_packet, proto) + self._packet_call = self._call_later(_PACKET_INTERVAL_S, self._send_next_packet, proto) def _finish_playback(self) -> None: if self._delay_call is not None and self._delay_call.active(): diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index 8f30e32..c365556 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -285,6 +285,23 @@ class SubscriptionStore(ABC): def list_by_phase(self, phase: "SubscriptionPhase") -> tuple["Subscription", ...]: ... + @abstractmethod + def relay_tables_with_active_source( + self, + system: str, + slot: int, + dst_tgid: int, + ) -> tuple[str, ...]: + ... + + @abstractmethod + def legs_in_table(self, table_key: str) -> tuple["Subscription", ...]: + ... + + @abstractmethod + def has_active_target_leg(self, system: str, slot: int, tgid: int) -> bool: + ... + class ProxySlotStore(ABC): """Hotspot session registry keyed by peer_id (Phase 3).""" diff --git a/src/adn_server/application/subscription/echo_seed.py b/src/adn_server/application/subscription/echo_seed.py new file mode 100644 index 0000000..e5bb6eb --- /dev/null +++ b/src/adn_server/application/subscription/echo_seed.py @@ -0,0 +1,87 @@ +# ADN DMR Peer Server - application subscription echo seed +# +# 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 +############################################################################### + +"""Initial BRIDGES snapshot for the ECHO parrot system (bootstrap only).""" + +from __future__ import annotations + +import time +from typing import Any + +from adn_server.domain import bytes_3 + + +def seed_echo_routing_table(config: dict[str, Any]) -> dict[str, list[dict[str, Any]]]: + """Initial BRIDGES for ECHO system (legacy make_bridges 9990 + MASTER expansion).""" + now = time.time() + timeout_sec = 2 * 60 + tgid_b = bytes_3(9990) + bridges: dict[str, list[dict[str, Any]]] = { + "9990": [ + { + "SYSTEM": "ECHO", + "TS": 2, + "TGID": tgid_b, + "ACTIVE": True, + "TIMEOUT": timeout_sec, + "TO_TYPE": "NONE", + "ON": [], + "OFF": [], + "RESET": [], + "TIMER": now + timeout_sec, + } + ] + } + systems_cfg = config.get("SYSTEMS", {}) + for _system, sys_cfg in systems_cfg.items(): + if _system == "ECHO": + continue + if sys_cfg.get("MODE") != "MASTER": + continue + _tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10)) + bridges["9990"].append( + { + "SYSTEM": _system, + "TS": 1, + "TGID": tgid_b, + "ACTIVE": False, + "TIMEOUT": _tmout * 60, + "TO_TYPE": "ON", + "OFF": [], + "ON": [tgid_b], + "RESET": [], + "TIMER": now, + } + ) + bridges["9990"].append( + { + "SYSTEM": _system, + "TS": 2, + "TGID": tgid_b, + "ACTIVE": False, + "TIMEOUT": _tmout * 60, + "TO_TYPE": "ON", + "OFF": [], + "ON": [tgid_b], + "RESET": [], + "TIMER": now, + } + ) + return bridges diff --git a/src/adn_server/application/subscription/router.py b/src/adn_server/application/subscription/router.py index 46d3e2f..8cfe8aa 100644 --- a/src/adn_server/application/subscription/router.py +++ b/src/adn_server/application/subscription/router.py @@ -31,9 +31,6 @@ class SubscriptionRouter: def __init__(self, store: SubscriptionStore) -> None: self._store = store - self._indexed = hasattr(store, "relay_tables_with_active_source") and hasattr( - store, "legs_in_table" - ) def resolve(self, ingress: VoiceIngress) -> tuple[ForwardLeg, ...]: """Return active forward legs when the source has an ACTIVE row on the dst TG (legacy to_target).""" @@ -50,16 +47,7 @@ class SubscriptionRouter: legs: list[ForwardLeg] = [] seen_obp: set[tuple[str, int]] = set() for table_key in tables: - subs = ( - self._store.legs_in_table(table_key) - if self._indexed - else tuple( - sub - for sub in self._store.snapshot() - if sub.table_key() == table_key - ) - ) - for sub in subs: + for sub in self._store.legs_in_table(table_key): if sub.system.value == ingress.source_system: continue if not sub.is_active(): @@ -80,21 +68,4 @@ class SubscriptionRouter: def relay_tables_with_active_source(self, system: str, slot: int, dst_tgid: int) -> tuple[str, ...]: """Mirror legacy ``relay_tables_with_active_source`` on subscription rows.""" - if self._indexed: - return self._store.relay_tables_with_active_source(system, slot, dst_tgid) - tables: list[str] = [] - seen: set[str] = set() - for sub in self._store.snapshot(): - if sub.system.value != system: - continue - if int(sub.channel.slot) != int(slot): - continue - if not sub.is_active(): - continue - if int(sub.target_tgid) != int(dst_tgid): - continue - key = sub.table_key() - if key not in seen: - seen.add(key) - tables.append(key) - return tuple(sorted(tables)) + return self._store.relay_tables_with_active_source(system, slot, dst_tgid) diff --git a/src/adn_server/application/subscription/store_sync.py b/src/adn_server/application/subscription/store_sync.py index b6c3009..d0683d1 100644 --- a/src/adn_server/application/subscription/store_sync.py +++ b/src/adn_server/application/subscription/store_sync.py @@ -22,7 +22,7 @@ Runtime hot paths mutate the store directly and publish via ``export_routing_table``; they must not call ``replace_store_from_routing_table`` (would overwrite store authority with the shim). -Bootstrap and tests may import once — e.g. ``_seed_echo_routing_table`` in ``peer_server.py``. +Bootstrap and tests may import once — e.g. ``seed_echo_routing_table`` in ``application/subscription/echo_seed.py``. """ from __future__ import annotations diff --git a/src/adn_server/application/subscription/subscription_queries.py b/src/adn_server/application/subscription/subscription_queries.py index f1d6f32..8ea4573 100644 --- a/src/adn_server/application/subscription/subscription_queries.py +++ b/src/adn_server/application/subscription/subscription_queries.py @@ -37,19 +37,7 @@ def system_has_active_leg_in_store( tgid: int, ) -> bool: """True when an ACTIVE subscription leg exists for ``(system, slot, tgid)``.""" - indexed = getattr(store, "has_active_target_leg", None) - if callable(indexed): - return bool(indexed(system, slot, tgid)) - for sub in store.snapshot(): - if sub.system.value != system: - continue - if int(sub.channel.slot) != int(slot): - continue - if int(sub.target_tgid) != int(tgid): - continue - if sub.is_active(): - return True - return False + return store.has_active_target_leg(system, slot, tgid) def active_system_slots_for_tg_in_store( diff --git a/src/adn_server/application/voice_use_cases.py b/src/adn_server/application/voice_use_cases.py index 7017d31..735f706 100644 --- a/src/adn_server/application/voice_use_cases.py +++ b/src/adn_server/application/voice_use_cases.py @@ -27,7 +27,6 @@ from __future__ import annotations import logging -import os import time from datetime import datetime from typing import Any, Callable @@ -72,7 +71,6 @@ class VoiceUseCases: self._tts_running: dict[int, bool] = {} self._announcement_last_hour: dict[int, int] = {} self._tts_last_hour: dict[int, int] = {} - self._config_file_mtime: float = 0.0 self._broadcast_queue: list[dict[str, Any]] = [] self._broadcast_active_tgs: set[str] = set() @@ -718,17 +716,8 @@ class VoiceUseCases: self._call_from_reactor(protocol.send_voice_packet, pkt, bytes_3(5000), bytes_3(9), _slot) logger.debug("(%s) disconnected voice thread end", system) - def check_voice_config_reload(self, config_file_path: str | None = None) -> None: - """Check adn-voice.yaml mtime and reload if changed (15s loop). Start/stop announcement LoopingCalls.""" - if config_file_path and os.path.isfile(config_file_path): - try: - mtime = os.path.getmtime(config_file_path) - except OSError: - return - if mtime == self._config_file_mtime: - return - self._config_file_mtime = mtime - logger.info("(VOICE-RELOAD) config file change detected, reloading configuration...") + 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 @@ -796,5 +785,3 @@ class VoiceUseCases: "(VOICE-RELOAD) %s enabled - mode: %s, file: %s, TG: %s", label, mode, item.get("FILE"), item.get("TG"), ) - if config_file_path and self._config_file_mtime: - logger.info("(VOICE-RELOAD) config reload completed") diff --git a/src/adn_server/infrastructure/bootstrap/peer_server.py b/src/adn_server/infrastructure/bootstrap/peer_server.py index d4436d5..201b2f6 100644 --- a/src/adn_server/infrastructure/bootstrap/peer_server.py +++ b/src/adn_server/infrastructure/bootstrap/peer_server.py @@ -44,7 +44,6 @@ from adn_server.application.proxy.deployment import ( proxy_target_system, ) from adn_server.application.report.queue import BoundedReportQueue, QueuedReportSender -from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes from adn_server.application.runtime_context import ( ConfigProxy, RuntimeContext, @@ -52,8 +51,8 @@ from adn_server.application.runtime_context import ( prepare_reload_config, swap_runtime_config, ) +from adn_server.application.subscription.echo_seed import seed_echo_routing_table from adn_server.application.subscription.store_sync import replace_store_from_routing_table -from adn_server.domain import bytes_3 from adn_server.domain.dmr.bptc import encode_emblc from adn_server.infrastructure.acl_router import InMemoryAclRouter from adn_server.infrastructure.config_normalizer import ( @@ -92,6 +91,7 @@ from adn_server.infrastructure.twisted_adapters.report.mqtt_publisher import ( from adn_server.infrastructure.twisted_adapters.report.worker import start_report_queue_worker from adn_server.infrastructure.twisted_adapters.report_server import ReportServerFactory from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocolFactory +from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes from adn_server.infrastructure.voice import DefaultVoiceProvider, StubVoiceProvider from adn_server.infrastructure.voice.recording import RecordingHandler @@ -163,39 +163,6 @@ def _wire_monitor_downlink_ctx( report_factory.set_downlink_ctx_for_system(_ctx_for) -def _seed_echo_routing_table(config: dict) -> dict: - """Initial BRIDGES for ECHO system (legacy make_bridges 9990 + MASTER expansion).""" - now = time.time() - timeout_sec = 2 * 60 - tgid_b = bytes_3(9990) - bridges: dict = { - "9990": [ - { - "SYSTEM": "ECHO", - "TS": 2, - "TGID": tgid_b, - "ACTIVE": True, - "TIMEOUT": timeout_sec, - "TO_TYPE": "NONE", - "ON": [], - "OFF": [], - "RESET": [], - "TIMER": now + timeout_sec, - } - ] - } - systems_cfg = config.get("SYSTEMS", {}) - for _system, sys_cfg in systems_cfg.items(): - if _system == "ECHO": - continue - if sys_cfg.get("MODE") != "MASTER": - continue - _tmout = float(sys_cfg.get("DEFAULT_UA_TIMER", 10)) - bridges["9990"].append({"SYSTEM": _system, "TS": 1, "TGID": tgid_b, "ACTIVE": False, "TIMEOUT": _tmout * 60, "TO_TYPE": "ON", "OFF": [], "ON": [tgid_b], "RESET": [], "TIMER": now}) - bridges["9990"].append({"SYSTEM": _system, "TS": 2, "TGID": tgid_b, "ACTIVE": False, "TIMEOUT": _tmout * 60, "TO_TYPE": "ON", "OFF": [], "ON": [tgid_b], "RESET": [], "TIMER": now}) - return bridges - - def _looping_errback(logger: logging.Logger, failure): """Errback for LoopingCalls (legacy loopingErrHandle). Stops reactor to avoid memory leaks.""" try: @@ -279,7 +246,7 @@ def run_peer_server( systems_cfg = config.get("SYSTEMS", {}) subscription_store = InMemorySubscriptionStore() if "ECHO" in systems_cfg and systems_cfg["ECHO"].get("MODE") in ("PEER", "MASTER"): - replace_store_from_routing_table(subscription_store, _seed_echo_routing_table(config)) + replace_store_from_routing_table(subscription_store, seed_echo_routing_table(config)) _ctx = runtime_holder.get() runtime_holder.swap( RuntimeContext( @@ -376,7 +343,7 @@ def run_peer_server( start_looping_call=_start_voice_loop, defer_to_thread=threads.deferToThread, ) - voice_use_cases.check_voice_config_reload(voice_config_path) + voice_use_cases.apply_voice_config() ident_use_cases = IdentUseCases( config, voice_use_cases, @@ -670,10 +637,29 @@ def run_peer_server( reactor.addSystemEventTrigger("before", "shutdown", shutdown_handler) # Voice config reload (15s): re-read adn-voice.yaml, start/stop announcement LoopingCalls on change + voice_config_mtime = 0.0 + if voice_config_path and os.path.isfile(voice_config_path): + try: + voice_config_mtime = os.path.getmtime(voice_config_path) + except OSError: + voice_config_mtime = 0.0 + def voice_reload_loop(): + nonlocal voice_config_mtime try: + if voice_config_path and os.path.isfile(voice_config_path): + try: + mtime = os.path.getmtime(voice_config_path) + except OSError: + return + if mtime == voice_config_mtime: + 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.check_voice_config_reload(voice_config_path) + voice_use_cases.apply_voice_config() + if voice_config_path and voice_config_mtime: + logger.info("(VOICE-RELOAD) config reload completed") except Exception as e: logger.debug("(VOICE-RELOAD) %s", e) diff --git a/src/adn_server/infrastructure/echo/runtime.py b/src/adn_server/infrastructure/echo/runtime.py index b8f9dcf..a1b4711 100644 --- a/src/adn_server/infrastructure/echo/runtime.py +++ b/src/adn_server/infrastructure/echo/runtime.py @@ -100,7 +100,11 @@ def run_echo(config: dict[str, Any], *, logger: logging.Logger) -> None: logger.debug("(ECHO) skip %s (MODE=%s; echo only runs PEER systems)", system_name, sys_cfg.get("MODE")) continue - pb = PlaybackUseCases(system_name, get_protocol=lambda sn=system_name: protocols.get(sn)) + pb = PlaybackUseCases( + system_name, + call_later=reactor.callLater, + get_protocol=lambda sn=system_name: protocols.get(sn), + ) protocol = HBPProtocolFactory( system_name, diff --git a/tests/application/test_store_sync.py b/tests/application/test_store_sync.py index c28944a..0ffdfc8 100644 --- a/tests/application/test_store_sync.py +++ b/tests/application/test_store_sync.py @@ -27,7 +27,7 @@ from typing import Any from adn_server.application.subscription.routing_table_export import export_routing_table from adn_server.application.subscription.store_sync import replace_store_from_routing_table from adn_server.domain import bytes_3, int_id -from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table +from adn_server.application.subscription.echo_seed import seed_echo_routing_table from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore @@ -59,7 +59,7 @@ def test_replace_store_from_echo_bridges_round_trip(): "MASTER-B": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 15}, } } - bridges = _seed_echo_routing_table(config) + bridges = seed_echo_routing_table(config) store = InMemorySubscriptionStore() replace_store_from_routing_table(store, bridges) diff --git a/tests/echo/test_playback_ingress.py b/tests/echo/test_playback_ingress.py index 3dd0207..c62dc8e 100644 --- a/tests/echo/test_playback_ingress.py +++ b/tests/echo/test_playback_ingress.py @@ -26,7 +26,12 @@ import logging from unittest.mock import patch from tests.harness.deterministic import DeterministicScenario, PacketSpec -from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, send_playback +from tests.harness.playback_helpers import ( + FakePlaybackProtocol, + make_capture_call_later, + noop_call_later, + send_playback, +) from adn_server.application.playback_use_cases import _PLAYBACK_DELAY_S, _RECORD_IDLE_S, PlaybackUseCases @@ -34,7 +39,7 @@ from adn_server.application.playback_use_cases import _PLAYBACK_DELAY_S, _RECORD def test_dmrd_received_accepts_ingress_pkt_time_kwarg() -> None: """Regression: PEER udp_hbp passes ingress_pkt_time (echo must not TypeError).""" proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + pb = PlaybackUseCases("ECHO", call_later=noop_call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x12121212, slot=2) args = DeterministicScenario.voice_head_spec(base).decoded_hbp_args() @@ -60,23 +65,22 @@ def test_dmrd_received_accepts_ingress_pkt_time_kwarg() -> None: def test_ingress_pkt_time_enables_record_to_playback(caplog) -> None: """Regression: PEER path (ingress_pkt_time) records voice and schedules playback.""" proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x34343434, slot=2) - mock_reactor, scheduled = install_reactor_capture() t0 = 1_700_000_000.0 - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - with patch("adn_server.application.playback_use_cases.time", return_value=t0): - send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) - send_playback( - pb, - "ECHO", - DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), - ingress_pkt_time=t0 + 2.0, - ) - with patch("adn_server.application.playback_use_cases.time", return_value=t0 + _RECORD_IDLE_S + 1): - with caplog.at_level(logging.INFO, logger="adn_server.application.playback_use_cases"): - pb._on_record_idle(proto) + with patch("adn_server.application.playback_use_cases.time", return_value=t0): + send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) + send_playback( + pb, + "ECHO", + DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), + ingress_pkt_time=t0 + 2.0, + ) + with patch("adn_server.application.playback_use_cases.time", return_value=t0 + _RECORD_IDLE_S + 1): + with caplog.at_level(logging.INFO, logger="adn_server.application.playback_use_cases"): + pb._on_record_idle(proto) assert any(item[0] == _PLAYBACK_DELAY_S for item in scheduled) assert pb._playback_busy is True diff --git a/tests/echo/test_playback_logging.py b/tests/echo/test_playback_logging.py index 6b63601..5703a05 100644 --- a/tests/echo/test_playback_logging.py +++ b/tests/echo/test_playback_logging.py @@ -26,7 +26,7 @@ import logging import re from tests.harness.deterministic import DeterministicScenario, PacketSpec -from tests.harness.playback_helpers import FakePlaybackProtocol +from tests.harness.playback_helpers import FakePlaybackProtocol, noop_call_later from adn_server.application.playback_use_cases import PlaybackUseCases from adn_server.domain import bytes_3, bytes_4 @@ -35,7 +35,7 @@ from adn_server.domain import bytes_3, bytes_4 def test_start_playback_logs_duration_with_two_decimals(caplog) -> None: """PLAYBACK duration matches bridge-style %.2f (no float noise in logs).""" proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + pb = PlaybackUseCases("ECHO", call_later=noop_call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x88888888, slot=2) recorded = [DeterministicScenario.voice_head_spec(base).data()] diff --git a/tests/echo/test_playback_send_loop.py b/tests/echo/test_playback_send_loop.py index cc25cda..657e263 100644 --- a/tests/echo/test_playback_send_loop.py +++ b/tests/echo/test_playback_send_loop.py @@ -25,7 +25,7 @@ from __future__ import annotations from unittest.mock import patch from tests.harness.deterministic import DeterministicScenario, PacketSpec -from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, run_scheduled +from tests.harness.playback_helpers import FakePlaybackProtocol, make_capture_call_later, run_scheduled from adn_server.application.playback_use_cases import ( _PACKET_INTERVAL_S, @@ -38,17 +38,16 @@ from adn_server.domain import bytes_3, bytes_4 def test_send_next_packet_emits_all_packets_then_finishes() -> None: proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) pb._playback_busy = True pb._playback_stream_id = bytes_4(0x77777777) pb._playback_packets = [bytes([i]) * 55 for i in range(1, 5)] pb._playback_index = 0 - mock_reactor, scheduled = install_reactor_capture() - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - pb._send_next_packet(proto) - while scheduled: - run_scheduled(scheduled) + pb._send_next_packet(proto) + while scheduled: + run_scheduled(scheduled) assert len(proto.sent) == 4 assert pb._playback_busy is False @@ -58,31 +57,29 @@ def test_send_next_packet_emits_all_packets_then_finishes() -> None: def test_start_playback_sends_first_packet_and_schedules_rest() -> None: proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x88888888, slot=2) recorded = [ DeterministicScenario.voice_head_spec(base).data(), DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1).data(), ] - mock_reactor, scheduled = install_reactor_capture() - - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - pb._start_playback( - proto, - recorded, - bytes_3(base.rf_src), - bytes_4(base.peer_id), - bytes_3(base.dst_id), - 2, - 1.5, - ) + + pb._start_playback( + proto, + recorded, + bytes_3(base.rf_src), + bytes_4(base.peer_id), + bytes_3(base.dst_id), + 2, + 1.5, + ) assert len(proto.sent) == 1 assert pb._playback_index == 1 assert any(item[0] == _PACKET_INTERVAL_S for item in scheduled) - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - while scheduled: - run_scheduled(scheduled) + while scheduled: + run_scheduled(scheduled) assert len(proto.sent) == 2 assert pb._playback_busy is False @@ -90,7 +87,8 @@ def test_start_playback_sends_first_packet_and_schedules_rest() -> None: def test_max_duration_commits_recording_when_no_vterm() -> None: proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x99999999, slot=2) pb._recording_active = True pb.CALL_DATA = [DeterministicScenario.voice_head_spec(base).data()] @@ -102,11 +100,9 @@ def test_max_duration_commits_recording_when_no_vterm() -> None: "peer_id": bytes_4(base.peer_id), "dst_id": bytes_3(base.dst_id), } - mock_reactor, scheduled = install_reactor_capture() - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - with patch("adn_server.application.playback_use_cases.time", return_value=100.0 + _SOURCE_MAX_S): - pb._on_record_max_duration(proto) + with patch("adn_server.application.playback_use_cases.time", return_value=100.0 + _SOURCE_MAX_S): + pb._on_record_max_duration(proto) assert not pb._recording_active assert pb._playback_busy is True @@ -115,14 +111,13 @@ def test_max_duration_commits_recording_when_no_vterm() -> None: def test_packet_interval_matches_expected() -> None: proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) pb._playback_busy = True pb._playback_stream_id = bytes_4(0xAAAAAAAA) pb._playback_packets = [b"\x01" * 55, b"\x02" * 55] pb._playback_index = 0 - mock_reactor, scheduled = install_reactor_capture() - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - pb._send_next_packet(proto) + pb._send_next_packet(proto) assert scheduled[0][0] == _PACKET_INTERVAL_S diff --git a/tests/echo/test_recording_timers.py b/tests/echo/test_recording_timers.py index 24d97ce..969510e 100644 --- a/tests/echo/test_recording_timers.py +++ b/tests/echo/test_recording_timers.py @@ -27,7 +27,8 @@ from unittest.mock import patch from tests.harness.deterministic import DeterministicScenario, PacketSpec from tests.harness.playback_helpers import ( FakePlaybackProtocol, - install_reactor_capture, + make_capture_call_later, + noop_call_later, run_scheduled, send_playback, ) @@ -37,22 +38,21 @@ from adn_server.application.playback_use_cases import _RECORD_IDLE_S, PlaybackUs def test_idle_timeout_appends_synthetic_vterm_and_schedules_playback() -> None: proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x55555555, slot=2) - mock_reactor, scheduled = install_reactor_capture() - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - with patch("adn_server.application.playback_use_cases.time", return_value=200.0): - send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) - send_playback( - pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), - ) + with patch("adn_server.application.playback_use_cases.time", return_value=200.0): + send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) + send_playback( + pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), + ) - idle_calls = [item for item in scheduled if item[0] == _RECORD_IDLE_S] - assert idle_calls + idle_calls = [item for item in scheduled if item[0] == _RECORD_IDLE_S] + assert idle_calls - with patch("adn_server.application.playback_use_cases.time", return_value=200.0 + _RECORD_IDLE_S): - run_scheduled(scheduled, delay=_RECORD_IDLE_S) + with patch("adn_server.application.playback_use_cases.time", return_value=200.0 + _RECORD_IDLE_S): + run_scheduled(scheduled, delay=_RECORD_IDLE_S) assert not pb._recording_active assert pb._playback_busy is True @@ -60,7 +60,7 @@ def test_idle_timeout_appends_synthetic_vterm_and_schedules_playback() -> None: def test_vterm_commit_does_not_require_synthetic_vterm() -> None: - pb = PlaybackUseCases("ECHO") + pb = PlaybackUseCases("ECHO", call_later=noop_call_later) base = PacketSpec(dst_id=9990, stream_id=0x66666666, slot=2) recorded = [ DeterministicScenario.voice_head_spec(base).data(), diff --git a/tests/echo/test_rekey_playback.py b/tests/echo/test_rekey_playback.py index 0e3129d..056cc56 100644 --- a/tests/echo/test_rekey_playback.py +++ b/tests/echo/test_rekey_playback.py @@ -26,7 +26,12 @@ from unittest.mock import patch import pytest from tests.harness.deterministic import DeterministicScenario, PacketSpec -from tests.harness.playback_helpers import FakePlaybackProtocol, install_reactor_capture, send_playback +from tests.harness.playback_helpers import ( + FakePlaybackProtocol, + make_capture_call_later, + noop_call_later, + send_playback, +) from adn_server.application.playback_use_cases import PlaybackUseCases from adn_server.domain import bytes_4 @@ -54,7 +59,7 @@ def _long_voice_recording( @pytest.mark.behavior def test_prepare_playback_preserves_source_seq_past_255_packets() -> None: """Regression: long QSO keeps source seq; new stream segment when seq byte wraps.""" - pb = PlaybackUseCases("ECHO") + pb = PlaybackUseCases("ECHO", call_later=noop_call_later) base = PacketSpec(dst_id=9990, stream_id=0x55555555, slot=2) recorded = _long_voice_recording(base, burst_count=500) pb._playback_stream_id = bytes_4(0x77777777) @@ -83,45 +88,47 @@ def test_prepare_playback_preserves_source_seq_past_255_packets() -> None: def test_start_playback_sends_preserved_seq_over_30s_recording() -> None: """End-to-end: ~500 bursts @ 60ms ≈ 30s; replay seq on wire must match recording.""" proto = FakePlaybackProtocol() - pb = PlaybackUseCases("ECHO", get_protocol=lambda: proto) + call_later, scheduled = make_capture_call_later() + pb = PlaybackUseCases("ECHO", call_later=call_later, get_protocol=lambda: proto) base = PacketSpec(dst_id=9990, stream_id=0x66666666, slot=2) burst_count = 500 recorded = _long_voice_recording(base, burst_count=burst_count) - mock_reactor, scheduled = install_reactor_capture() pb._playback_stream_id = bytes_4(0x77777777) out = pb._prepare_playback_packets(recorded) assert len(out) == len(recorded) + 1 - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - pb._playback_packets = out - pb._playback_index = 0 - pb._playback_busy = True - pb._send_next_packet(proto) - while scheduled: - _delay, fn, args = scheduled.pop(0) - fn(*args) + pb._playback_packets = out + pb._playback_index = 0 + pb._playback_busy = True + pb._send_next_packet(proto) + while scheduled: + _delay, fn, args = scheduled.pop(0) + fn(*args) assert len(proto.sent) == len(out) def test_recording_rekey_does_not_store_second_vhead() -> None: - pb = PlaybackUseCases("ECHO", get_protocol=lambda: FakePlaybackProtocol()) + call_later, _scheduled = make_capture_call_later() + pb = PlaybackUseCases( + "ECHO", + call_later=call_later, + get_protocol=lambda: FakePlaybackProtocol(), + ) base = PacketSpec(dst_id=9990, stream_id=0x11111111, slot=2) rekey = PacketSpec(dst_id=9990, stream_id=0x22222222, slot=2) - mock_reactor, _scheduled = install_reactor_capture() - - with patch("adn_server.application.playback_use_cases.reactor", mock_reactor): - with patch("adn_server.application.playback_use_cases.time") as mock_time: - mock_time.side_effect = [100.0, 100.1, 100.2, 100.3] - send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) - send_playback( - pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), - ) - send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(rekey)) - send_playback( - pb, "ECHO", DeterministicScenario.voice_burst_spec(rekey, seq=2, dtype_vseq=2), - ) + + with patch("adn_server.application.playback_use_cases.time") as mock_time: + mock_time.side_effect = [100.0, 100.1, 100.2, 100.3] + send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(base)) + send_playback( + pb, "ECHO", DeterministicScenario.voice_burst_spec(base, seq=1, dtype_vseq=1), + ) + send_playback(pb, "ECHO", DeterministicScenario.voice_head_spec(rekey)) + send_playback( + pb, "ECHO", DeterministicScenario.voice_burst_spec(rekey, seq=2, dtype_vseq=2), + ) vheads = sum( 1 @@ -133,7 +140,7 @@ def test_recording_rekey_does_not_store_second_vhead() -> None: def test_prepare_playback_skips_mid_call_vhead_and_preserves_seq() -> None: - pb = PlaybackUseCases("ECHO") + pb = PlaybackUseCases("ECHO", call_later=noop_call_later) base = PacketSpec(dst_id=9990, stream_id=0x33333333, slot=2) rekey = PacketSpec(dst_id=9990, stream_id=0x44444444, slot=2) recorded = [ diff --git a/tests/harness/playback_helpers.py b/tests/harness/playback_helpers.py index d3d4720..ba7d1c6 100644 --- a/tests/harness/playback_helpers.py +++ b/tests/harness/playback_helpers.py @@ -22,6 +22,7 @@ from __future__ import annotations +from collections.abc import Callable from typing import Any from unittest.mock import MagicMock @@ -41,10 +42,17 @@ def packet_bytes(spec: PacketSpec) -> bytes: return spec.data() -def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[Any, ...]]]]: - """Patch reactor.callLater; returns mock and scheduled (delay, fn, args) list.""" +def noop_call_later(*_args: Any, **_kwargs: Any) -> MagicMock: + """No-op scheduler for tests that never fire timers.""" + handle = MagicMock() + handle.active.return_value = False + handle.cancel = MagicMock() + return handle + + +def make_capture_call_later() -> tuple[Callable[..., Any], list[tuple[float, Any, tuple[Any, ...]]]]: + """Build a call_later stub; returns (callable, scheduled (delay, fn, args) list).""" scheduled: list[tuple[float, Any, tuple[Any, ...]]] = [] - mock_reactor = MagicMock() def call_later(delay: float, fn: Any, *args: Any) -> MagicMock: scheduled.append((delay, fn, args)) @@ -53,6 +61,13 @@ def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[A handle.cancel = MagicMock() return handle + return call_later, scheduled + + +def install_reactor_capture() -> tuple[MagicMock, list[tuple[float, Any, tuple[Any, ...]]]]: + """Backward-compatible alias; prefer make_capture_call_later + call_later= kwarg.""" + call_later, scheduled = make_capture_call_later() + mock_reactor = MagicMock() mock_reactor.callLater = call_later return mock_reactor, scheduled diff --git a/tests/routing/test_echo_bootstrap_seed.py b/tests/routing/test_echo_bootstrap_seed.py index 076abd9..6eb5a0e 100644 --- a/tests/routing/test_echo_bootstrap_seed.py +++ b/tests/routing/test_echo_bootstrap_seed.py @@ -26,7 +26,7 @@ from tests.harness.assertions import assert_forwarded from tests.harness.deterministic import DeterministicScenario, PacketSpec from tests.routing.test_echo_subscription_reset import _echo_scenario_config -from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table +from adn_server.application.subscription.echo_seed import seed_echo_routing_table def _prod_like_config() -> dict: @@ -49,7 +49,7 @@ def test_echo_routing_missing_without_echo_store_seed() -> None: def test_echo_routing_works_after_echo_seed_on_vhead() -> None: config = _prod_like_config() - scenario = DeterministicScenario(config=config, routing_table=_seed_echo_routing_table(config)) + scenario = DeterministicScenario(config=config, routing_table=seed_echo_routing_table(config)) scenario.routing.apply_startup_subscriptions() base = PacketSpec(dst_id=9990, stream_id=0xABCDEF02, slot=2) diff --git a/tests/routing/test_echo_subscription_reset.py b/tests/routing/test_echo_subscription_reset.py index 5dbde6f..f46d84f 100644 --- a/tests/routing/test_echo_subscription_reset.py +++ b/tests/routing/test_echo_subscription_reset.py @@ -27,7 +27,7 @@ from tests.harness.assertions import assert_forwarded from tests.harness.deterministic import DeterministicScenario, PacketSpec, minimal_config from adn_server.domain import ID_MAX, PEER_MAX -from adn_server.infrastructure.bootstrap.peer_server import _seed_echo_routing_table +from adn_server.application.subscription.echo_seed import seed_echo_routing_table from adn_server.infrastructure.config_loader import acl_build @@ -67,7 +67,7 @@ def _echo_scenario_config() -> dict: def _bridges_after_bridgereset(config: dict) -> dict: - bridges = _seed_echo_routing_table(config) + bridges = seed_echo_routing_table(config) for entry in bridges["9990"]: if entry["SYSTEM"] == "ECHO": entry["ACTIVE"] = False diff --git a/tests/voice/test_voice_config_reload.py b/tests/voice/test_voice_config_reload.py index d2f8caf..c9f45b6 100644 --- a/tests/voice/test_voice_config_reload.py +++ b/tests/voice/test_voice_config_reload.py @@ -38,7 +38,7 @@ def _reload_uc(scenario, *, start_looping_call) -> VoiceUseCases: ) -def test_check_voice_config_reload_starts_enabled_announcement_loop() -> None: +def test_apply_voice_config_starts_enabled_announcement_loop() -> None: scenario, _ = voice_master_scenario() scenario.config["VOICE"] = { "ANNOUNCEMENTS": [ @@ -61,14 +61,14 @@ def test_check_voice_config_reload_starts_enabled_announcement_loop() -> None: return handle uc = _reload_uc(scenario, start_looping_call=start_looping_call) - uc.check_voice_config_reload() + uc.apply_voice_config() assert len(started) == 1 assert started[0][0] == 120.0 assert 0 in uc._ann_tasks -def test_check_voice_config_reload_stops_removed_announcement() -> None: +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() @@ -80,13 +80,13 @@ def test_check_voice_config_reload_stops_removed_announcement() -> None: ) uc._ann_tasks[0] = MagicMock(running=True, stop=stop_mock) - uc.check_voice_config_reload() + uc.apply_voice_config() assert 0 not in uc._ann_tasks stop_mock.assert_called_once() -def test_check_voice_config_reload_starts_tts_loop() -> None: +def test_apply_voice_config_starts_tts_loop() -> None: scenario, _ = voice_master_scenario() scenario.config["VOICE"] = { "TTS_ANNOUNCEMENTS": [ @@ -106,7 +106,7 @@ def test_check_voice_config_reload_starts_tts_loop() -> None: return MagicMock(running=True, stop=MagicMock()) uc = _reload_uc(scenario, start_looping_call=start_looping_call) - uc.check_voice_config_reload() + uc.apply_voice_config() assert started == [30.0] assert 0 in uc._tts_tasks