From dc0af385f300c5c7dfa0cf7e0d033d34297ca070 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Tue, 9 Jun 2026 12:01:03 -0400 Subject: [PATCH] feat(report): slim TCP wire with dashboard_state snapshots Send STATE_SND dashboard_state instead of separate topology/routing deltas on the report TCP path, aligned with the monitor FastAPI ingest. --- .../twisted_adapters/report/opcodes.py | 2 + .../twisted_adapters/report/wire.py | 70 +++++++------------ .../test_monitor_report_contract.py | 60 ++++++++++++---- .../infrastructure/test_report_server_wire.py | 39 ++++------- 4 files changed, 88 insertions(+), 83 deletions(-) diff --git a/src/adn_server/infrastructure/twisted_adapters/report/opcodes.py b/src/adn_server/infrastructure/twisted_adapters/report/opcodes.py index 90fc768..421f433 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report/opcodes.py +++ b/src/adn_server/infrastructure/twisted_adapters/report/opcodes.py @@ -15,6 +15,8 @@ REPORT_OPCODES = { "ROUTING_TABLE_SND": b"\x11", "VOICE_EVENT_SND": b"\x12", "DELTA_SND": b"\x13", + "STATE_SND": b"\x14", + "STATE_REQ": b"\x15", "HELLO": b"\xff", } diff --git a/src/adn_server/infrastructure/twisted_adapters/report/wire.py b/src/adn_server/infrastructure/twisted_adapters/report/wire.py index e8f98f3..6d8a4ac 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report/wire.py +++ b/src/adn_server/infrastructure/twisted_adapters/report/wire.py @@ -1,4 +1,4 @@ -"""Report wire: typed JSON topology / routing_table / voice_event / delta.""" +"""Report wire: slim dashboard_state + voice_event to monitor (D-25).""" from __future__ import annotations @@ -11,16 +11,12 @@ from adn_server.application.ports import ReportWireEncoder from adn_server.application.report import ( REPORT_FEATURES, REPORT_PROTOCOL, - build_routing_table, - build_topology, + build_dashboard_state, hello_connected_system_names, parse_bridge_event_csv, - routing_table_delta, - topology_delta, ) from .opcodes import REPORT_OPCODES, SERVER_NAME, server_version -from .pickle_legacy import encode_config_snd_frame logger = logging.getLogger(__name__) @@ -29,14 +25,19 @@ def _json_wire(opcode: bytes, payload: dict[str, Any]) -> bytes: return opcode + json.dumps(payload, separators=(",", ":")).encode("utf-8") +def _state_dedup_key(payload: dict[str, Any]) -> bytes: + return json.dumps( + {"ctable": payload.get("ctable"), "server_id": payload.get("server_id")}, + separators=(",", ":"), + sort_keys=True, + ).encode("utf-8") + + class ReportWire(ReportWireEncoder): - """Stateful report encoder — seq counters and delta patches.""" + """Slim monitor encoder — ``dashboard_state`` snapshots + ``voice_event`` only.""" def __init__(self) -> None: - self._topology_seq = 0 - self._routing_seq = 0 - self._last_topology: dict[str, Any] | None = None - self._last_routing: dict[str, Any] | None = None + self._last_state_key: bytes | None = None def hello_frames(self, systems: dict[str, Any]) -> tuple[bytes, ...]: names = hello_connected_system_names(systems) @@ -53,48 +54,27 @@ class ReportWire(ReportWireEncoder): return (REPORT_OPCODES["HELLO"] + payload,) def config_frames(self, systems: dict[str, Any], *, full_snapshot: bool) -> tuple[bytes, ...]: + return self.state_frames(systems, force=full_snapshot) + + def state_frames(self, systems: dict[str, Any], *, force: bool = False) -> tuple[bytes, ...]: ts = time.time() - if full_snapshot or self._last_topology is None: - self._topology_seq += 1 - topology = build_topology(systems, seq=self._topology_seq, ts=ts) - self._last_topology = topology - logger.debug("(REPORT) CONFIG_SND pickle + TOPOLOGY_SND seq=%s", self._topology_seq) - return ( - encode_config_snd_frame(systems), - _json_wire(REPORT_OPCODES["TOPOLOGY_SND"], topology), - ) - self._topology_seq += 1 - current = build_topology(systems, seq=self._topology_seq, ts=ts) - delta = topology_delta(self._last_topology, current, seq=self._topology_seq, ts=ts) - if delta is None: - self._topology_seq -= 1 + payload = build_dashboard_state(systems, ts=ts) + key = _state_dedup_key(payload) + if not force and self._last_state_key == key: + logger.debug("(REPORT) STATE_SND unchanged, skip") return () - self._last_topology = current - logger.debug("(REPORT) DELTA_SND topology seq=%s", delta["seq"]) - return (_json_wire(REPORT_OPCODES["DELTA_SND"], delta),) + self._last_state_key = key + logger.debug("(REPORT) STATE_SND ts=%s", payload.get("ts")) + return (_json_wire(REPORT_OPCODES["STATE_SND"], payload),) def bridge_frames(self, bridges: dict[str, Any], *, full_snapshot: bool) -> tuple[bytes, ...]: - ts = time.time() - if full_snapshot or self._last_routing is None: - self._routing_seq += 1 - routing = build_routing_table(bridges, seq=self._routing_seq, ts=ts) - self._last_routing = routing - logger.debug("(REPORT) ROUTING_TABLE_SND seq=%s", self._routing_seq) - return (_json_wire(REPORT_OPCODES["ROUTING_TABLE_SND"], routing),) - self._routing_seq += 1 - current = build_routing_table(bridges, seq=self._routing_seq, ts=ts) - delta = routing_table_delta(self._last_routing, current, seq=self._routing_seq, ts=ts) - if delta is None: - self._routing_seq -= 1 - return () - self._last_routing = current - logger.debug("(REPORT) DELTA_SND routing seq=%s", delta["seq"]) - return (_json_wire(REPORT_OPCODES["DELTA_SND"], delta),) + """Routing is not exported to the monitor (D-25).""" + return () def bridge_event_frames(self, event: str) -> tuple[bytes, ...]: voice = parse_bridge_event_csv(event) if voice is None: - logger.warning("(REPORT) BRDG_EVENT not mapped to voice_event: %s", event[:120]) + logger.warning("(REPORT) voice_event not emitted (unmapped CSV): %s", event[:120]) return () logger.debug( "(REPORT) VOICE_EVENT_SND %s %s %s", diff --git a/tests/application/test_monitor_report_contract.py b/tests/application/test_monitor_report_contract.py index f2722e1..25f611a 100644 --- a/tests/application/test_monitor_report_contract.py +++ b/tests/application/test_monitor_report_contract.py @@ -8,7 +8,6 @@ These tests model that behaviour and assert the server report snapshot is safe. from __future__ import annotations import json -import pickle from typing import Any import pytest @@ -164,19 +163,51 @@ def test_truncated_monitor_ctable_cannot_show_system_n_peers() -> None: assert full["SYSTEM-6"]["PEERS"] -def test_report_wire_emits_legacy_pickle_before_topology() -> None: +def test_report_wire_emits_dashboard_state_only() -> None: peer_specs = [(730039101, 4)] snapshot = _production_snapshot(peer_specs) - frames = ReportWire().config_frames(snapshot, full_snapshot=True) + frames = ReportWire().state_frames(snapshot, force=True) - assert len(frames) == 2 - assert frames[0][:1] == REPORT_OPCODES["CONFIG_SND"] - assert frames[1][:1] == REPORT_OPCODES["TOPOLOGY_SND"] - pickle_cfg = pickle.loads(frames[0][1:]) - topo = json.loads(frames[1][1:].decode()) - assert len(_system_n_names(pickle_cfg)) == MAX_PEERS - topo_names = {row["name"] for row in topo["systems"]} - assert topo_names == set(pickle_cfg) + assert len(frames) == 1 + assert frames[0][:1] == REPORT_OPCODES["STATE_SND"] + doc = json.loads(frames[0][1:].decode()) + assert doc["type"] == "dashboard_state" + assert "SYSTEM-4" in doc["ctable"]["MASTERS"] + assert len(_system_n_names(snapshot)) == MAX_PEERS + + +def test_report_wire_skips_unchanged_dashboard_state() -> None: + peer_specs = [(730039101, 4)] + snapshot = _production_snapshot(peer_specs) + wire = ReportWire() + assert len(wire.state_frames(snapshot, force=True)) == 1 + assert wire.state_frames(snapshot, force=False) == () + + +def test_report_wire_bridge_frames_empty() -> None: + bridges = { + "52090": [ + {"SYSTEM": "SYSTEM-0", "ACTIVE": True, "TS": 2, "TGID": 52090}, + ] + } + assert ReportWire().bridge_frames(bridges, full_snapshot=True) == () + + +def test_report_wire_emits_voice_event_only() -> None: + csv = "GROUP VOICE,START,RX,SYSTEM,3262598598,730039101,730039101,2,730444" + frames = ReportWire().bridge_event_frames(csv) + + assert len(frames) == 1 + assert frames[0][:1] == REPORT_OPCODES["VOICE_EVENT_SND"] + voice = json.loads(frames[0][1:].decode()) + assert voice["type"] == "voice_event" + + +def test_report_wire_skips_unmapped_voice_event() -> None: + csv = "GROUP VOICE,START,RX,SYSTEM-0,1001,3120001,2,52090" + frames = ReportWire().bridge_event_frames(csv) + + assert frames == () def test_report_factory_push_matches_monitor_contract() -> None: @@ -187,12 +218,13 @@ def test_report_factory_push_matches_monitor_contract() -> None: factory.set_peer_slot_map(lambda: _peer_slots(peer_specs)) snapshot = factory._systems_for_report() - frames = ReportWire().config_frames(snapshot, full_snapshot=True) - config = pickle.loads(frames[0][1:]) + frames = ReportWire().state_frames(snapshot, force=True) + assert len(frames) == 1 + assert frames[0][:1] == REPORT_OPCODES["STATE_SND"] ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS) before = count_masters(ctable) - update_ctable_from_config(config, ctable) + update_ctable_from_config(snapshot, ctable) assert len(_system_n_names(snapshot)) == MAX_PEERS assert count_masters(ctable) == before diff --git a/tests/infrastructure/test_report_server_wire.py b/tests/infrastructure/test_report_server_wire.py index ea50b0b..1d72e84 100644 --- a/tests/infrastructure/test_report_server_wire.py +++ b/tests/infrastructure/test_report_server_wire.py @@ -40,18 +40,14 @@ def test_connect_sends_hello_topology_and_routing() -> None: client = _CapturingClient() factory.clients.append(client) factory._send_hello_to(client) - factory._send_config_to(client, full_snapshot=True) - factory._send_bridge_to(client, full_snapshot=True) + factory._send_state_to(client, force=True) - assert len(client.messages) == 4 + assert len(client.messages) == 2 hello = json.loads(client.messages[0][1:].decode()) assert hello["report_protocol"] == 2 assert "REPORT_V2" in hello["features"] - assert client.messages[1][:1] == REPORT_OPCODES["CONFIG_SND"] - assert client.messages[2][:1] == REPORT_OPCODES["TOPOLOGY_SND"] - assert json.loads(client.messages[2][1:])["type"] == "topology" - assert client.messages[3][:1] == REPORT_OPCODES["ROUTING_TABLE_SND"] - assert json.loads(client.messages[3][1:])["type"] == "routing_table" + assert client.messages[1][:1] == REPORT_OPCODES["STATE_SND"] + assert json.loads(client.messages[1][1:])["type"] == "dashboard_state" def test_bridge_event_emits_voice_event_json() -> None: @@ -78,7 +74,7 @@ def test_incremental_bridge_update_sends_delta() -> None: ) client = _CapturingClient() factory.clients.append(client) - factory._send_bridge_to(client, full_snapshot=True) + factory._send_state_to(client, force=True) factory.set_bridges( { "52090": [ @@ -86,11 +82,8 @@ def test_incremental_bridge_update_sends_delta() -> None: ], } ) - factory._send_bridge_to(client, full_snapshot=False) - assert client.messages[1][:1] == REPORT_OPCODES["DELTA_SND"] - delta = json.loads(client.messages[1][1:].decode()) - assert delta["type"] == "delta" - assert delta["patch"]["type"] == "routing_table" + factory.send_bridge(incremental=True) + assert len(client.messages) == 1 def test_inject_proxy_topology_expanded_for_monitor() -> None: @@ -122,19 +115,17 @@ def test_inject_proxy_topology_expanded_for_monitor() -> None: factory.set_peer_slot_map(lambda: {peer: 2}) client = _CapturingClient() factory.clients.append(client) - factory._send_config_to(client, full_snapshot=True) - assert client.messages[0][:1] == REPORT_OPCODES["CONFIG_SND"] - topology = json.loads(client.messages[1][1:].decode()) - names = {system["name"] for system in topology["systems"]} + factory._send_state_to(client, force=True) + assert client.messages[0][:1] == REPORT_OPCODES["STATE_SND"] + state_doc = json.loads(client.messages[0][1:].decode()) + names = set((state_doc.get("ctable") or {}).get("MASTERS", {})) + names |= set((state_doc.get("ctable") or {}).get("OPENBRIDGES", {})) assert "SYSTEM" not in names assert "SYSTEM-2" in names - assert "SYSTEM-0" in names - assert len([name for name in names if name.startswith("SYSTEM-")]) == 102 - system = next(item for item in topology["systems"] if item["name"] == "SYSTEM-2") + system = state_doc["ctable"]["MASTERS"]["SYSTEM-2"] assert system["port"] == 56402 - assert system["peers"][0]["id"] == 730039101 - empty = next(item for item in topology["systems"] if item["name"] == "SYSTEM-0") - assert empty["peers"] == [] + peer_keys = {int(k) for k in system["peers"]} + assert 730039101 in peer_keys def test_send_bridge_event_remaps_inject_proxy_system_name() -> None: