fix: audit wave 1 server hygiene refactors

Inject call_later into PlaybackUseCases, move voice config mtime watch to
bootstrap, complete SubscriptionStore port methods, and relocate echo
routing seed to the application layer.
pull/36/head
Rodrigo Pérez 3 months ago
parent 0b9e536279
commit 574b100807

@ -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():

@ -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)."""

@ -0,0 +1,87 @@
# ADN DMR Peer Server - application subscription echo seed
#
# 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
###############################################################################
"""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

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

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

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

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

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

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

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

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

@ -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()]

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

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

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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.