diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index 7ab34ae..92327e1 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -115,6 +115,16 @@ class ReportWireEncoder(ABC): class ReportSender(ABC): """Send config and bridge state to report TCP clients (CONFIG_SND, BRIDGE_SND, BRDG_EVENT).""" + @abstractmethod + def set_systems(self, systems: dict[str, Any]) -> None: + """Update cached SYSTEMS snapshot used by the wire encoder.""" + ... + + @abstractmethod + def set_bridges(self, bridges: dict[str, Any]) -> None: + """Update cached BRIDGES snapshot used by the wire encoder.""" + ... + @abstractmethod def send_config(self, systems: dict[str, Any], *, incremental: bool = False) -> None: """Send CONFIG_SND (pickle systems) or topology / delta JSON.""" diff --git a/src/adn_server/application/report/__init__.py b/src/adn_server/application/report/__init__.py index 908c397..9e03ac5 100644 --- a/src/adn_server/application/report/__init__.py +++ b/src/adn_server/application/report/__init__.py @@ -1,5 +1,11 @@ """Report application layer: payload mapping and protocol mode (no Twisted / wire bytes).""" +from .queue import ( + DEFAULT_MAX_DRAIN_PER_TICK, + DEFAULT_MAX_EVENTS, + BoundedReportQueue, + QueuedReportSender, +) from .payloads import ( REPORT_FEATURES, REPORT_PROTOCOL, @@ -12,6 +18,10 @@ from .payloads import ( ) __all__ = [ + "DEFAULT_MAX_DRAIN_PER_TICK", + "DEFAULT_MAX_EVENTS", + "BoundedReportQueue", + "QueuedReportSender", "REPORT_FEATURES", "REPORT_PROTOCOL", "build_routing_table", diff --git a/src/adn_server/application/report/payloads.py b/src/adn_server/application/report/payloads.py index 31be4c3..dc1afb7 100644 --- a/src/adn_server/application/report/payloads.py +++ b/src/adn_server/application/report/payloads.py @@ -65,11 +65,26 @@ _TOPOLOGY_PEER_FIELDS: tuple[tuple[str, str], ...] = ( ) +def _peer_connected_at(peer: dict[str, Any]) -> int | None: + """Unix time when peer logged in (legacy CONFIG ``CONNECTED``), or None.""" + if not _peer_connected(peer): + return None + raw = peer.get("CONNECTED", 0) + try: + ts = int(float(raw)) + except (TypeError, ValueError): + return None + return ts if ts > 0 else None + + def _topology_peer_row(peer_key: Any, peer: dict[str, Any]) -> dict[str, Any]: row: dict[str, Any] = { "id": _dmr_id(peer_key), "connected": _peer_connected(peer), } + connected_at = _peer_connected_at(peer) + if connected_at is not None: + row["connected_at"] = connected_at if peer.get("IP"): row["ip"] = str(peer["IP"]) if peer.get("PORT") is not None: diff --git a/src/adn_server/application/report/queue.py b/src/adn_server/application/report/queue.py new file mode 100644 index 0000000..cbbba90 --- /dev/null +++ b/src/adn_server/application/report/queue.py @@ -0,0 +1,104 @@ +"""Bounded in-process report queue — decouple hot path from TCP encode/send.""" + +from __future__ import annotations + +import logging +from collections import deque +from dataclasses import dataclass, field +from typing import Any + +from ..ports import ReportSender + +logger = logging.getLogger(__name__) + +DEFAULT_MAX_EVENTS = 2048 +DEFAULT_MAX_DRAIN_PER_TICK = 128 + + +@dataclass +class BoundedReportQueue: + """Coalesce config/bridge snapshots; bound voice-event backlog (drop oldest).""" + + max_events: int = DEFAULT_MAX_EVENTS + max_drain_per_tick: int = DEFAULT_MAX_DRAIN_PER_TICK + _events: deque[str] = field(default_factory=deque, init=False, repr=False) + _pending_config: tuple[dict[str, Any], bool] | None = field(default=None, init=False, repr=False) + _pending_bridge: tuple[dict[str, Any], bool] | None = field(default=None, init=False, repr=False) + dropped_events: int = field(default=0, init=False) + + def enqueue_event(self, event: str) -> None: + if len(self._events) >= self.max_events: + self._events.popleft() + self.dropped_events += 1 + self._events.append(event) + + def enqueue_config(self, systems: dict[str, Any], *, incremental: bool = False) -> None: + self._pending_config = (systems, incremental) + + def enqueue_bridge(self, bridges: dict[str, Any], *, incremental: bool = False) -> None: + self._pending_bridge = (bridges, incremental) + + def pending_count(self) -> int: + n = len(self._events) + if self._pending_config is not None: + n += 1 + if self._pending_bridge is not None: + n += 1 + return n + + def drain(self, sender: ReportSender) -> int: + """Flush pending work to ``sender``; at most ``max_drain_per_tick`` voice events per call.""" + sent = 0 + budget = self.max_drain_per_tick + while self._events and budget > 0: + sender.send_bridge_event(self._events.popleft()) + sent += 1 + budget -= 1 + if self._pending_config is not None: + systems, incremental = self._pending_config + self._pending_config = None + sender.set_systems(systems) + sender.send_config(systems, incremental=incremental) + sent += 1 + if self._pending_bridge is not None: + bridges, incremental = self._pending_bridge + self._pending_bridge = None + sender.set_bridges(bridges) + sender.send_bridge(bridges, incremental=incremental) + sent += 1 + if self.dropped_events and sent: + logger.debug("(REPORT) queue drained %s item(s); dropped_events=%s", sent, self.dropped_events) + return sent + + +class QueuedReportSender(ReportSender): + """``ReportSender`` port: enqueue only; a reactor worker drains to ``inner``.""" + + def __init__(self, queue: BoundedReportQueue, inner: ReportSender) -> None: + self._queue = queue + self._inner = inner + + def set_systems(self, systems: dict[str, Any]) -> None: + self._inner.set_systems(systems) + + def set_bridges(self, bridges: dict[str, Any]) -> None: + self._inner.set_bridges(bridges) + + def send_config(self, systems: dict[str, Any], *, incremental: bool = False) -> None: + self._inner.set_systems(systems) + self._queue.enqueue_config(systems, incremental=incremental) + + def send_bridge(self, bridges: dict[str, Any], *, incremental: bool = False) -> None: + self._inner.set_bridges(bridges) + self._queue.enqueue_bridge(bridges, incremental=incremental) + + def send_bridge_event(self, event: str) -> None: + self._queue.enqueue_event(event) + + @property + def inner(self) -> ReportSender: + return self._inner + + @property + def queue(self) -> BoundedReportQueue: + return self._queue diff --git a/src/adn_server/infrastructure/twisted_adapters/report/__init__.py b/src/adn_server/infrastructure/twisted_adapters/report/__init__.py index c55d5e1..1fe6994 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report/__init__.py +++ b/src/adn_server/infrastructure/twisted_adapters/report/__init__.py @@ -14,11 +14,14 @@ from adn_server.application.ports import ReportWireEncoder from .opcodes import REPORT_OPCODES from .wire import ReportWire +from .worker import DEFAULT_DRAIN_INTERVAL_SEC, start_report_queue_worker __all__ = [ + "DEFAULT_DRAIN_INTERVAL_SEC", "REPORT_OPCODES", "ReportWire", "create_report_wire", + "start_report_queue_worker", ] diff --git a/src/adn_server/infrastructure/twisted_adapters/report/worker.py b/src/adn_server/infrastructure/twisted_adapters/report/worker.py new file mode 100644 index 0000000..a295a17 --- /dev/null +++ b/src/adn_server/infrastructure/twisted_adapters/report/worker.py @@ -0,0 +1,43 @@ +"""Twisted LoopingCall worker that drains the bounded report queue.""" + +from __future__ import annotations + +import logging +from typing import Any, Callable + +from twisted.internet import task + +from adn_server.application.report.queue import BoundedReportQueue +from adn_server.application.ports import ReportSender + +logger = logging.getLogger(__name__) + +DEFAULT_DRAIN_INTERVAL_SEC = 0.05 + + +def start_report_queue_worker( + queue: BoundedReportQueue, + sender: ReportSender, + *, + interval_sec: float = DEFAULT_DRAIN_INTERVAL_SEC, + on_errback: Callable[[Any], None] | None = None, +) -> task.LoopingCall: + """Start periodic drain of ``queue`` into ``sender`` (synchronous TCP send on reactor).""" + + def _drain() -> None: + try: + queue.drain(sender) + except Exception as e: + logger.warning("(REPORT) queue drain failed: %s", e) + + loop = task.LoopingCall(_drain) + deferred = loop.start(interval_sec, now=False) + if on_errback is not None: + deferred.addErrback(on_errback) + logger.info( + "(REPORT) queue worker started interval=%.3fs max_events=%s drain_per_tick=%s", + interval_sec, + queue.max_events, + queue.max_drain_per_tick, + ) + return loop diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 9adb9c3..1dc3775 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -532,7 +532,7 @@ class HBPProtocol(DatagramProtocol): if hasattr(report, "set_systems"): report.set_systems(systems) if hasattr(report, "send_config"): - report.send_config() + report.send_config(systems) logger.debug("(REPORT) Pushed CONFIG_SND after peer state change on %s", self._system) def datagramReceived(self, data: bytes, addr: tuple[str, int]) -> None: diff --git a/src/adn_server/main.py b/src/adn_server/main.py index 2401ca3..b56edc7 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -88,6 +88,8 @@ from .infrastructure.persistence.keys_store import JsonKeysStore from .infrastructure.security.password_download import DefaultSecurityDownloader, StubSecurityDownloader from .infrastructure.security.user_passwords_loader import UserPasswordsLoader from .infrastructure.talker_alias_emblc import default_ta_emblc_encoder +from .application.report.queue import BoundedReportQueue, QueuedReportSender +from .infrastructure.twisted_adapters.report.worker import start_report_queue_worker from .infrastructure.twisted_adapters.report_server import ReportServerFactory from .infrastructure.twisted_adapters.udp_hbp import HBPProtocolFactory from .infrastructure.voice import DefaultVoiceProvider, StubVoiceProvider @@ -100,6 +102,12 @@ class ReportSenderAdapter(ReportSender): def __init__(self, factory: ReportServerFactory) -> None: self._factory = factory + def set_systems(self, systems) -> None: + self._factory.set_systems(systems) + + def set_bridges(self, bridges) -> None: + self._factory.set_bridges(bridges) + def send_config(self, systems, *, incremental: bool = False) -> None: self._factory.set_systems(systems) self._factory.send_config(incremental=incremental) @@ -277,7 +285,9 @@ def main() -> None: if p is not None and hasattr(p, "_obp_send_bcsq"): p._obp_send_bcsq(tgid, stream_id) - report_sender = ReportSenderAdapter(report_factory) + report_inner = ReportSenderAdapter(report_factory) + report_queue = BoundedReportQueue() + report_sender = QueuedReportSender(report_queue, report_inner) report_factory.set_systems(systems_cfg) report_factory.set_bridges(bridge_router.get_bridges()) reporting_use_cases = ReportingUseCases(report_sender, config) @@ -351,6 +361,11 @@ def main() -> None: port = config["REPORTS"].get("REPORT_PORT", 4321) reactor.listenTCP(port, report_factory) logger.info("(REPORT) Report server listening on TCP %s", port) + start_report_queue_worker( + report_queue, + report_inner, + on_errback=lambda e: _looping_errback(e, logger), + ) # Reporting loop (REPORT_INTERVAL) — same logs as legacy after send_config/send_bridge def reporting_loop(): @@ -490,7 +505,7 @@ def main() -> None: return HBPProtocolFactory( system_name, config, - report_factory, + report_sender, router=bridge_router, dmrd_received=bridge_use_cases.dmrd_received, get_user_password_callback=user_passwords_loader.get_user_password, diff --git a/tests/application/test_report_payloads.py b/tests/application/test_report_payloads.py index f88c51e..119136b 100644 --- a/tests/application/test_report_payloads.py +++ b/tests/application/test_report_payloads.py @@ -28,6 +28,25 @@ def validator() -> jsonschema.Draft202012Validator: return jsonschema.Draft202012Validator(schema) +def test_build_topology_exports_peer_connected_at() -> None: + login_ts = 1717555100.0 + systems = { + "MASTER-A": { + "MODE": "MASTER", + "ENABLED": True, + "PEERS": { + bytes_3(3120001): { + "CONNECTION": "YES", + "CONNECTED": login_ts, + } + }, + }, + } + doc = build_topology(systems, seq=1, ts=1717555200.0) + peer = doc["systems"][0]["peers"][0] + assert peer["connected_at"] == int(login_ts) + + def test_build_topology_matches_example(validator: jsonschema.Draft202012Validator) -> None: with (_EXAMPLES_DIR / "topology.json").open(encoding="utf-8") as fh: expected = json.load(fh) diff --git a/tests/application/test_report_queue.py b/tests/application/test_report_queue.py new file mode 100644 index 0000000..5ca2561 --- /dev/null +++ b/tests/application/test_report_queue.py @@ -0,0 +1,91 @@ +"""Unit tests for bounded report queue and queued sender.""" + +from __future__ import annotations + +from typing import Any + +from adn_server.application.report.queue import BoundedReportQueue, QueuedReportSender + + +class RecordingSender: + def __init__(self) -> None: + self.systems: dict[str, Any] | None = None + self.bridges: dict[str, Any] | None = None + self.config_calls: list[tuple[dict, bool]] = [] + self.bridge_calls: list[tuple[dict, bool]] = [] + self.events: list[str] = [] + + def set_systems(self, systems: dict[str, Any]) -> None: + self.systems = systems + + def set_bridges(self, bridges: dict[str, Any]) -> None: + self.bridges = bridges + + def send_config(self, systems: dict[str, Any], *, incremental: bool = False) -> None: + self.config_calls.append((systems, incremental)) + + def send_bridge(self, bridges: dict[str, Any], *, incremental: bool = False) -> None: + self.bridge_calls.append((bridges, incremental)) + + def send_bridge_event(self, event: str) -> None: + self.events.append(event) + + +def test_enqueue_event_is_non_blocking(): + inner = RecordingSender() + queue = BoundedReportQueue(max_events=4, max_drain_per_tick=10) + sender = QueuedReportSender(queue, inner) + sender.send_bridge_event("GROUP VOICE,START,RX,SYS,1,2,3,1,4") + assert inner.events == [] + assert queue.pending_count() == 1 + + +def test_drain_delivers_events_before_snapshots(): + inner = RecordingSender() + queue = BoundedReportQueue(max_drain_per_tick=10) + sender = QueuedReportSender(queue, inner) + sender.send_bridge_event("e1") + sender.send_config({"SYS": {}}, incremental=True) + sender.send_bridge({"99": []}, incremental=False) + queue.drain(inner) + assert inner.events == ["e1"] + assert len(inner.config_calls) == 1 + assert inner.config_calls[0][1] is True + assert len(inner.bridge_calls) == 1 + + +def test_config_and_bridge_coalesce_to_latest(): + inner = RecordingSender() + queue = BoundedReportQueue() + sender = QueuedReportSender(queue, inner) + sender.send_config({"A": 1}) + sender.send_config({"B": 2}, incremental=True) + sender.send_bridge({"old": []}) + sender.send_bridge({"new": []}, incremental=True) + queue.drain(inner) + assert inner.config_calls == [({"B": 2}, True)] + assert inner.bridge_calls == [({"new": []}, True)] + + +def test_bounded_queue_drops_oldest_events(): + inner = RecordingSender() + queue = BoundedReportQueue(max_events=2, max_drain_per_tick=10) + sender = QueuedReportSender(queue, inner) + sender.send_bridge_event("e1") + sender.send_bridge_event("e2") + sender.send_bridge_event("e3") + assert queue.dropped_events == 1 + queue.drain(inner) + assert inner.events == ["e2", "e3"] + + +def test_drain_respects_per_tick_budget(): + inner = RecordingSender() + queue = BoundedReportQueue(max_drain_per_tick=2) + sender = QueuedReportSender(queue, inner) + for i in range(5): + sender.send_bridge_event(f"e{i}") + queue.drain(inner) + assert inner.events == ["e0", "e1"] + queue.drain(inner) + assert inner.events == ["e0", "e1", "e2", "e3"] diff --git a/tests/harness/deterministic.py b/tests/harness/deterministic.py index 2a823c8..2dafce7 100644 --- a/tests/harness/deterministic.py +++ b/tests/harness/deterministic.py @@ -193,6 +193,12 @@ class FakeReportSender: def __init__(self, factory: FakeReportFactory) -> None: self._factory = factory + def set_systems(self, systems: dict[str, Any]) -> None: + pass + + def set_bridges(self, bridges: dict[str, Any]) -> None: + pass + def send_config(self, systems: dict[str, Any], *, incremental: bool = False) -> None: pass