feat(report): bounded queue worker and topology connected_at

Decouple bridge events from TCP send via a reactor-drained queue (50ms, 2048
events). Export peer connected_at in topology so monitor timers survive refresh.
pull/1/head
Rodrigo Pérez 4 months ago
parent 048c27ab23
commit 99169e77ca

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

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

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

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

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

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

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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.