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.
pull/1/head
Rodrigo Pérez 4 months ago
parent cdf486feb4
commit dc0af385f3

@ -15,6 +15,8 @@ REPORT_OPCODES = {
"ROUTING_TABLE_SND": b"\x11", "ROUTING_TABLE_SND": b"\x11",
"VOICE_EVENT_SND": b"\x12", "VOICE_EVENT_SND": b"\x12",
"DELTA_SND": b"\x13", "DELTA_SND": b"\x13",
"STATE_SND": b"\x14",
"STATE_REQ": b"\x15",
"HELLO": b"\xff", "HELLO": b"\xff",
} }

@ -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 from __future__ import annotations
@ -11,16 +11,12 @@ from adn_server.application.ports import ReportWireEncoder
from adn_server.application.report import ( from adn_server.application.report import (
REPORT_FEATURES, REPORT_FEATURES,
REPORT_PROTOCOL, REPORT_PROTOCOL,
build_routing_table, build_dashboard_state,
build_topology,
hello_connected_system_names, hello_connected_system_names,
parse_bridge_event_csv, parse_bridge_event_csv,
routing_table_delta,
topology_delta,
) )
from .opcodes import REPORT_OPCODES, SERVER_NAME, server_version from .opcodes import REPORT_OPCODES, SERVER_NAME, server_version
from .pickle_legacy import encode_config_snd_frame
logger = logging.getLogger(__name__) 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") 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): class ReportWire(ReportWireEncoder):
"""Stateful report encoder — seq counters and delta patches.""" """Slim monitor encoder — ``dashboard_state`` snapshots + ``voice_event`` only."""
def __init__(self) -> None: def __init__(self) -> None:
self._topology_seq = 0 self._last_state_key: bytes | None = None
self._routing_seq = 0
self._last_topology: dict[str, Any] | None = None
self._last_routing: dict[str, Any] | None = None
def hello_frames(self, systems: dict[str, Any]) -> tuple[bytes, ...]: def hello_frames(self, systems: dict[str, Any]) -> tuple[bytes, ...]:
names = hello_connected_system_names(systems) names = hello_connected_system_names(systems)
@ -53,48 +54,27 @@ class ReportWire(ReportWireEncoder):
return (REPORT_OPCODES["HELLO"] + payload,) return (REPORT_OPCODES["HELLO"] + payload,)
def config_frames(self, systems: dict[str, Any], *, full_snapshot: bool) -> tuple[bytes, ...]: 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() ts = time.time()
if full_snapshot or self._last_topology is None: payload = build_dashboard_state(systems, ts=ts)
self._topology_seq += 1 key = _state_dedup_key(payload)
topology = build_topology(systems, seq=self._topology_seq, ts=ts) if not force and self._last_state_key == key:
self._last_topology = topology logger.debug("(REPORT) STATE_SND unchanged, skip")
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
return () return ()
self._last_topology = current self._last_state_key = key
logger.debug("(REPORT) DELTA_SND topology seq=%s", delta["seq"]) logger.debug("(REPORT) STATE_SND ts=%s", payload.get("ts"))
return (_json_wire(REPORT_OPCODES["DELTA_SND"], delta),) return (_json_wire(REPORT_OPCODES["STATE_SND"], payload),)
def bridge_frames(self, bridges: dict[str, Any], *, full_snapshot: bool) -> tuple[bytes, ...]: def bridge_frames(self, bridges: dict[str, Any], *, full_snapshot: bool) -> tuple[bytes, ...]:
ts = time.time() """Routing is not exported to the monitor (D-25)."""
if full_snapshot or self._last_routing is None: return ()
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),)
def bridge_event_frames(self, event: str) -> tuple[bytes, ...]: def bridge_event_frames(self, event: str) -> tuple[bytes, ...]:
voice = parse_bridge_event_csv(event) voice = parse_bridge_event_csv(event)
if voice is None: 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 () return ()
logger.debug( logger.debug(
"(REPORT) VOICE_EVENT_SND %s %s %s", "(REPORT) VOICE_EVENT_SND %s %s %s",

@ -8,7 +8,6 @@ These tests model that behaviour and assert the server report snapshot is safe.
from __future__ import annotations from __future__ import annotations
import json import json
import pickle
from typing import Any from typing import Any
import pytest import pytest
@ -164,19 +163,51 @@ def test_truncated_monitor_ctable_cannot_show_system_n_peers() -> None:
assert full["SYSTEM-6"]["PEERS"] 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)] peer_specs = [(730039101, 4)]
snapshot = _production_snapshot(peer_specs) 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 len(frames) == 1
assert frames[0][:1] == REPORT_OPCODES["CONFIG_SND"] assert frames[0][:1] == REPORT_OPCODES["STATE_SND"]
assert frames[1][:1] == REPORT_OPCODES["TOPOLOGY_SND"] doc = json.loads(frames[0][1:].decode())
pickle_cfg = pickle.loads(frames[0][1:]) assert doc["type"] == "dashboard_state"
topo = json.loads(frames[1][1:].decode()) assert "SYSTEM-4" in doc["ctable"]["MASTERS"]
assert len(_system_n_names(pickle_cfg)) == MAX_PEERS assert len(_system_n_names(snapshot)) == MAX_PEERS
topo_names = {row["name"] for row in topo["systems"]}
assert topo_names == set(pickle_cfg)
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: 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)) factory.set_peer_slot_map(lambda: _peer_slots(peer_specs))
snapshot = factory._systems_for_report() snapshot = factory._systems_for_report()
frames = ReportWire().config_frames(snapshot, full_snapshot=True) frames = ReportWire().state_frames(snapshot, force=True)
config = pickle.loads(frames[0][1:]) assert len(frames) == 1
assert frames[0][:1] == REPORT_OPCODES["STATE_SND"]
ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS) ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS)
before = count_masters(ctable) 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 len(_system_n_names(snapshot)) == MAX_PEERS
assert count_masters(ctable) == before assert count_masters(ctable) == before

@ -40,18 +40,14 @@ def test_connect_sends_hello_topology_and_routing() -> None:
client = _CapturingClient() client = _CapturingClient()
factory.clients.append(client) factory.clients.append(client)
factory._send_hello_to(client) factory._send_hello_to(client)
factory._send_config_to(client, full_snapshot=True) factory._send_state_to(client, force=True)
factory._send_bridge_to(client, full_snapshot=True)
assert len(client.messages) == 4 assert len(client.messages) == 2
hello = json.loads(client.messages[0][1:].decode()) hello = json.loads(client.messages[0][1:].decode())
assert hello["report_protocol"] == 2 assert hello["report_protocol"] == 2
assert "REPORT_V2" in hello["features"] assert "REPORT_V2" in hello["features"]
assert client.messages[1][:1] == REPORT_OPCODES["CONFIG_SND"] assert client.messages[1][:1] == REPORT_OPCODES["STATE_SND"]
assert client.messages[2][:1] == REPORT_OPCODES["TOPOLOGY_SND"] assert json.loads(client.messages[1][1:])["type"] == "dashboard_state"
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"
def test_bridge_event_emits_voice_event_json() -> None: def test_bridge_event_emits_voice_event_json() -> None:
@ -78,7 +74,7 @@ def test_incremental_bridge_update_sends_delta() -> None:
) )
client = _CapturingClient() client = _CapturingClient()
factory.clients.append(client) factory.clients.append(client)
factory._send_bridge_to(client, full_snapshot=True) factory._send_state_to(client, force=True)
factory.set_bridges( factory.set_bridges(
{ {
"52090": [ "52090": [
@ -86,11 +82,8 @@ def test_incremental_bridge_update_sends_delta() -> None:
], ],
} }
) )
factory._send_bridge_to(client, full_snapshot=False) factory.send_bridge(incremental=True)
assert client.messages[1][:1] == REPORT_OPCODES["DELTA_SND"] assert len(client.messages) == 1
delta = json.loads(client.messages[1][1:].decode())
assert delta["type"] == "delta"
assert delta["patch"]["type"] == "routing_table"
def test_inject_proxy_topology_expanded_for_monitor() -> None: 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}) factory.set_peer_slot_map(lambda: {peer: 2})
client = _CapturingClient() client = _CapturingClient()
factory.clients.append(client) factory.clients.append(client)
factory._send_config_to(client, full_snapshot=True) factory._send_state_to(client, force=True)
assert client.messages[0][:1] == REPORT_OPCODES["CONFIG_SND"] assert client.messages[0][:1] == REPORT_OPCODES["STATE_SND"]
topology = json.loads(client.messages[1][1:].decode()) state_doc = json.loads(client.messages[0][1:].decode())
names = {system["name"] for system in topology["systems"]} 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" not in names
assert "SYSTEM-2" in names assert "SYSTEM-2" in names
assert "SYSTEM-0" in names system = state_doc["ctable"]["MASTERS"]["SYSTEM-2"]
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")
assert system["port"] == 56402 assert system["port"] == 56402
assert system["peers"][0]["id"] == 730039101 peer_keys = {int(k) for k in system["peers"]}
empty = next(item for item in topology["systems"] if item["name"] == "SYSTEM-0") assert 730039101 in peer_keys
assert empty["peers"] == []
def test_send_bridge_event_remaps_inject_proxy_system_name() -> None: def test_send_bridge_event_remaps_inject_proxy_system_name() -> None:

Loading…
Cancel
Save

Powered by TurnKey Linux.