From ca59b886fb39f0c1a4950a7384e76603e37f0091 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 8 Jun 2026 23:22:57 -0400 Subject: [PATCH] feat(report): monitor topology expansion and proxy peer resolution Map proxy upstream slots into VOICE_EVENT_SND, resolve RX peer IDs for dashboard chips, and keep ECHO/parrot TX-RX colors aligned with legacy. --- src/adn_server/application/bridge/helpers.py | 66 ++++- .../application/bridge_use_cases.py | 49 +++- .../application/report/monitor_topology.py | 230 ++++++++++++++++ .../twisted_adapters/report/pickle_legacy.py | 9 + .../twisted_adapters/report/wire.py | 8 +- .../twisted_adapters/report_server.py | 31 ++- .../test_monitor_report_contract.py | 235 ++++++++++++++++ tests/application/test_monitor_topology.py | 252 ++++++++++++++++++ tests/bridge/test_report_peer_id.py | 49 ++++ .../infrastructure/test_report_server_wire.py | 93 ++++++- tests/support/monitor_ctable_sim.py | 144 ++++++++++ 11 files changed, 1149 insertions(+), 17 deletions(-) create mode 100644 src/adn_server/application/report/monitor_topology.py create mode 100644 tests/application/test_monitor_report_contract.py create mode 100644 tests/application/test_monitor_topology.py create mode 100644 tests/bridge/test_report_peer_id.py create mode 100644 tests/support/monitor_ctable_sim.py diff --git a/src/adn_server/application/bridge/helpers.py b/src/adn_server/application/bridge/helpers.py index 6e4211a..b033344 100644 --- a/src/adn_server/application/bridge/helpers.py +++ b/src/adn_server/application/bridge/helpers.py @@ -7,7 +7,7 @@ from __future__ import annotations from typing import Any -from ...domain import bytes_3, int_id +from ...domain import bytes_3, bytes_4, int_id # Embedded LC codeword sits at bits 116:148 inside the 48-bit EMB field (108:156). # Legacy bridge_master.py replaces dmrbits[116:148] on bursts B–E (dtype_vseq 1–4). @@ -35,6 +35,70 @@ def obp_target_bcsq_quenches_stream( return False +def _peer_key_from_int(peer_key: Any) -> bytes: + if isinstance(peer_key, bytes): + return peer_key + if isinstance(peer_key, int): + return bytes_4(peer_key) + if isinstance(peer_key, str) and peer_key.isdigit(): + return bytes_4(int(peer_key)) + return bytes_4(int_id(peer_key)) + + +def _fuzzy_peer_matches( + val: int, + peers: dict[Any, Any], +) -> list[bytes]: + val_str = str(val) + matches: list[bytes] = [] + for pk in peers: + try: + pk_b = _peer_key_from_int(pk) + except (TypeError, ValueError): + continue + pk_int = int_id(pk_b) + pk_str = str(pk_int) + if pk_int == val or pk_int // 100 == val: + matches.append(pk_b) + continue + if len(val_str) >= 5 and len(pk_str) >= 7 and pk_str.startswith(val_str): + matches.append(pk_b) + return matches + + +def resolve_voice_peer_id( + peer_id: bytes, + rf_src: bytes, + system_name: str, + systems_cfg: dict[str, Any], +) -> bytes: + """Resolve BRDG_EVENT field 5 for RX legs from a MASTER (hotspot transmitting). + + Legacy bridge uses ``_peer_id`` from DMRD for TX legs unchanged. Only RX source + events need the full hotspot radio id so monitor ``rts_update`` marks that peer RX. + """ + peers = systems_cfg.get(system_name, {}).get("PEERS", {}) + if not isinstance(peers, dict) or not peers: + return peer_id + peer_b = peer_id if isinstance(peer_id, bytes) else bytes_4(int_id(peer_id)) + if peer_b in peers: + return peer_b + rf_b = rf_src if isinstance(rf_src, bytes) else bytes_3(int_id(rf_src)) + if rf_b in peers: + return rf_b + peer_matches = _fuzzy_peer_matches(int_id(peer_id), peers) + if len(peer_matches) == 1: + return peer_matches[0] + rf_matches = _fuzzy_peer_matches(int_id(rf_src), peers) + if len(rf_matches) == 1: + return rf_matches[0] + return peer_id + + +# Back-compat alias for tests and imports. +report_peer_id_for_hbp_target = resolve_voice_peer_id + + def is_special_tg(bridge_key: str) -> bool: """True if bridge is special TGID 9990-9999 (excluded from infinite timer).""" if bridge_key and bridge_key[0:1] == "#": diff --git a/src/adn_server/application/bridge_use_cases.py b/src/adn_server/application/bridge_use_cases.py index eb1d096..8363587 100644 --- a/src/adn_server/application/bridge_use_cases.py +++ b/src/adn_server/application/bridge_use_cases.py @@ -42,7 +42,7 @@ from ..domain.dmr import bptc from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4, int_id from .ports import BridgeRouter, DmrEmbeddedLcEncoder, TalkerAliasEmblcEncoder from .talker_alias_use_cases import TalkerAliasUseCases -from .bridge.helpers import obp_target_bcsq_quenches_stream +from .bridge.helpers import obp_target_bcsq_quenches_stream, resolve_voice_peer_id from .bridge.timers import BridgeTimerMixin from .bridge.obp_forward import BridgeObpForwardMixin from .bridge.hbp_forward import BridgeHbpForwardMixin @@ -247,9 +247,22 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi _obp_grp = source_is_obp and call_type in ("group", "vcsbk") if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: if not _obp_grp: + _rx_report_peer = peer_id + if not source_is_obp: + _rx_report_peer = resolve_voice_peer_id( + peer_id, + rf_src, + system_name, + systems_cfg, + ) self._send_bridge_event( "GROUP VOICE,START,RX,{},{},{},{},{},{}".format( - system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id) + system_name, + int_id(stream_id), + int_id(_rx_report_peer), + int_id(rf_src), + slot, + int_id(dst_id), ) ) elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: @@ -265,9 +278,23 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi start = st.get(slot, {}).get("RX_START") if start is not None: duration = pkt_time - start + _rx_report_peer = peer_id + if not source_is_obp: + _rx_report_peer = resolve_voice_peer_id( + peer_id, + rf_src, + system_name, + systems_cfg, + ) self._send_bridge_event( "GROUP VOICE,END,RX,{},{},{},{},{},{},{:.2f}".format( - system_name, int_id(stream_id), int_id(peer_id), int_id(rf_src), slot, int_id(dst_id), duration + system_name, + int_id(stream_id), + int_id(_rx_report_peer), + int_id(rf_src), + slot, + int_id(dst_id), + duration, ) ) # ── Exact port of legacy bridge.py routerOBP/routerHBP forwarding to targets ── @@ -539,7 +566,12 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi ) self._send_bridge_event( "GROUP VOICE,START,TX,{},{},{},{},{},{}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b) + entry["SYSTEM"], + int_id(stream_id), + int_id(peer_id), + int_id(rf_src), + entry_ts, + int_id(entry_tgid_b), ) ) # First successful forward to HBP (may not be VHEAD if earlier frames were hangtime-blocked). @@ -563,9 +595,16 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi elif frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: dmrbits = _ts_st["TX_T_LC"][0:98] + dmrbits[98:166] + _ts_st["TX_T_LC"][98:197] call_duration = pkt_time - _ts_st.get("TX_START", pkt_time) + _end_peer = _ts_st.get("TX_PEER", peer_id) self._send_bridge_event( "GROUP VOICE,END,TX,{},{},{},{},{},{},{:.2f}".format( - entry["SYSTEM"], int_id(stream_id), int_id(peer_id), int_id(rf_src), entry_ts, int_id(entry_tgid_b), call_duration + entry["SYSTEM"], + int_id(stream_id), + int_id(_end_peer), + int_id(rf_src), + entry_ts, + int_id(entry_tgid_b), + call_duration, ) ) elif dtype_vseq in (1, 2, 3, 4): diff --git a/src/adn_server/application/report/monitor_topology.py b/src/adn_server/application/report/monitor_topology.py new file mode 100644 index 0000000..7ce93fa --- /dev/null +++ b/src/adn_server/application/report/monitor_topology.py @@ -0,0 +1,230 @@ +"""Monitor topology parity for inject-only proxy (legacy SYSTEM-N / 56400+N). + +Hotspot radio IDs often share a user/subscriber prefix, e.g. user ``7300391`` with +HS1 ``730039101``, HS2 ``730039102``, … HS99 ``730039199`` (``user * 100 + n``). +Voice-event peer resolution must not guess when several connected peers match the +same user or 6-digit legacy prefix. +""" + +from __future__ import annotations + +import copy +from typing import Any + +from adn_server.application.proxy.deployment import is_proxy_inject_only, proxy_target_system +from adn_server.domain.value_objects import bytes_4, int_id + +DEFAULT_REPORT_BASE_PORT = 56400 + + +def _connected_peers(peers: dict[Any, dict[str, Any]]) -> list[tuple[Any, dict[str, Any]]]: + return [ + (peer_key, peer) + for peer_key, peer in peers.items() + if isinstance(peer, dict) and peer.get("CONNECTION") == "YES" + ] + + +def _resolve_slot_map( + connected: list[tuple[Any, dict[str, Any]]], + peer_slots: dict[bytes, int] | None, + *, + max_slots: int, +) -> dict[Any, int]: + """Map peer keys to upstream slot indices for monitor ``SYSTEM-N`` rows.""" + slot_map: dict[Any, int] = {} + if peer_slots: + for peer_key, _peer in connected: + if isinstance(peer_key, bytes) and peer_key in peer_slots: + slot_map[peer_key] = peer_slots[peer_key] + used = set(slot_map.values()) + for peer_key, _peer in sorted(connected, key=lambda item: int_id(item[0])): + if peer_key in slot_map: + continue + for index in range(max_slots): + if index not in used: + slot_map[peer_key] = index + used.add(index) + break + return slot_map + + +def expand_inject_proxy_systems( + config: dict[str, Any], + systems: dict[str, Any], + peer_slots: dict[bytes, int] | None = None, +) -> dict[str, Any]: + """Fan inject-only ``SYSTEM`` peers into ``SYSTEM-N`` masters for monitor/report. + + Runtime HBP stays on a single inject target; only the topology snapshot sent to + adn-monitor matches legacy ``expand_generator`` + adn-proxy upstream ports. + """ + target = proxy_target_system(config) + if not target or not is_proxy_inject_only(config, target): + return systems + sys_cfg = systems.get(target) + if not isinstance(sys_cfg, dict) or sys_cfg.get("MODE") != "MASTER": + return systems + peers = sys_cfg.get("PEERS", {}) + if not isinstance(peers, dict): + return systems + + max_slots = int(sys_cfg.get("MAX_PEERS", 1)) + base_port = int(sys_cfg.get("_REPORT_BASE_PORT", DEFAULT_REPORT_BASE_PORT)) + connected = _connected_peers(peers) + slot_map = _resolve_slot_map(connected, peer_slots, max_slots=max_slots) + + out = {name: cfg for name, cfg in systems.items() if name != target} + # Legacy ``expand_generator``: emit every virtual master (SYSTEM-0..N-1) so the + # monitor does not delete unused upstream slots on topology/config update. + for slot in range(max_slots): + virtual_name = f"{target}-{slot}" + virtual = copy.deepcopy(sys_cfg) + virtual["PORT"] = base_port + slot + virtual["PEERS"] = { + peer_key: peers[peer_key] + for peer_key, mapped in slot_map.items() + if mapped == slot and peer_key in peers + } + out[virtual_name] = virtual + return out + + +def _slot_for_voice_peer( + peer_key: bytes, + *, + peers: dict[Any, dict[str, Any]], + peer_slots: dict[bytes, int] | None, + max_slots: int, +) -> int | None: + connected = _connected_peers(peers) + slot_map = _resolve_slot_map(connected, peer_slots, max_slots=max_slots) + return slot_map.get(peer_key) + + +def _peer_key_from_int(peer_key: Any) -> bytes: + if isinstance(peer_key, bytes): + return peer_key + return bytes_4(int_id(peer_key)) + + +def _connected_peer_keys(peers: dict[Any, Any]) -> list[bytes]: + keys: list[bytes] = [] + for peer_key, peer in peers.items(): + if not isinstance(peer, dict) or peer.get("CONNECTION") != "YES": + continue + keys.append(_peer_key_from_int(peer_key)) + return keys + + +def _unique_peer_match(matches: list[bytes]) -> bytes | None: + unique = list(dict.fromkeys(matches)) + return unique[0] if len(unique) == 1 else None + + +def _peers_for_voice_candidate(val: int, connected: list[bytes]) -> list[bytes]: + """Match a BRDG_EVENT peer/subscriber field to connected hotspot radio ids.""" + exact = bytes_4(val) + if exact in connected: + return [exact] + val_str = str(val) + matches: list[bytes] = [] + for peer_key in connected: + peer_int = int_id(peer_key) + peer_str = str(peer_int) + if peer_str == val_str: + matches.append(peer_key) + continue + # user 7300391 → radios 730039101..730039199 (user * 100 + hs) + if peer_int // 100 == val: + matches.append(peer_key) + continue + if len(val_str) >= 5 and len(peer_str) >= 7 and peer_str.startswith(val_str): + matches.append(peer_key) + continue + if len(val_str) >= 7 and peer_str.startswith(val_str) and len(peer_str) == len(val_str) + 2: + matches.append(peer_key) + continue + # legacy 6-digit dst (bridge.py hotspot match) — ambiguous when user has >1 HS + if len(val_str) >= 6 and len(peer_str) >= 6 and peer_str[:6] == val_str[:6]: + matches.append(peer_key) + return matches + + +def _peer_key_from_voice_csv(parts: list[str], peers: dict[Any, Any]) -> bytes | None: + """Resolve hotspot radio id from legacy BRDG_EVENT CSV (peer_id, then rf_src). + + Prefer exact radio ids. Fuzzy user/6-digit matching applies only when a single + connected peer matches (e.g. one HS online for user 7300391). + """ + connected = _connected_peer_keys(peers) + if not connected: + return None + field_values: list[tuple[int, int]] = [] + for idx in (5, 6): + if len(parts) <= idx: + continue + raw = parts[idx].strip() + if not raw: + continue + try: + field_values.append((idx, int(raw))) + except ValueError: + continue + for _idx, val in field_values: + key = bytes_4(val) + if key in connected: + return key + for _idx, val in field_values: + matched = _peers_for_voice_candidate(val, connected) + resolved = _unique_peer_match(matched) + if resolved is not None: + return resolved + return None + + +def remap_inject_proxy_voice_event( + event: str, + config: dict[str, Any], + systems: dict[str, Any], + peer_slots: dict[bytes, int] | None = None, +) -> str: + """Map inject-only ``SYSTEM`` voice events to ``SYSTEM-N`` for monitor ``rts_update``. + + Topology expansion removes the runtime inject target from CONFIG/TOPOLOGY; voice + events must use the same virtual master name or dashboard hotspot chips stay idle. + """ + target = proxy_target_system(config) + if not target or not is_proxy_inject_only(config, target): + return event + parts = event.split(",") + if len(parts) < 6 or parts[3].strip() != target: + return event + sys_cfg = systems.get(target, {}) + if not isinstance(sys_cfg, dict): + return event + peers = sys_cfg.get("PEERS", {}) + if not isinstance(peers, dict): + return event + peer_key = _peer_key_from_voice_csv(parts, peers) + if peer_key is None: + return event + max_slots = int(sys_cfg.get("MAX_PEERS", 1)) + slot = _slot_for_voice_peer( + peer_key, peers=peers, peer_slots=peer_slots, max_slots=max_slots + ) + if slot is None: + return event + parts[3] = f"{target}-{slot}" + # RX legs: field 5 is the RF source peer — normalize to full hotspot radio id. + # TX legs: keep legacy field 5 (echo 9990, OBP server id) so the hotspot chip + # shows TX/green while receiving; rewriting to the hotspot id would mark RX/red. + if len(parts) > 5 and len(parts) > 2 and parts[2].strip() == "RX": + resolved_peer = int_id(peer_key) + try: + reported_peer = int(parts[5].strip()) + except ValueError: + reported_peer = None + if reported_peer != resolved_peer: + parts[5] = str(resolved_peer) + return ",".join(parts) diff --git a/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py b/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py index 9007740..438e353 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py +++ b/src/adn_server/infrastructure/twisted_adapters/report/pickle_legacy.py @@ -10,6 +10,15 @@ from .opcodes import REPORT_OPCODES _PICKLE_PROTOCOL = 2 +def encode_config_snd_frame( + systems: dict[str, Any], + *, + protocol: int = _PICKLE_PROTOCOL, +) -> bytes: + """``CONFIG_SND`` opcode + ``pickle.dumps(SYSTEMS, protocol=2)`` (legacy ``hblink.send_config``).""" + return REPORT_OPCODES["CONFIG_SND"] + pickle.dumps(systems, protocol=protocol) + + def encode_bridge_snd_frame( bridges: dict[str, Any], *, diff --git a/src/adn_server/infrastructure/twisted_adapters/report/wire.py b/src/adn_server/infrastructure/twisted_adapters/report/wire.py index 647f661..e8f98f3 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report/wire.py +++ b/src/adn_server/infrastructure/twisted_adapters/report/wire.py @@ -20,6 +20,7 @@ from adn_server.application.report import ( ) from .opcodes import REPORT_OPCODES, SERVER_NAME, server_version +from .pickle_legacy import encode_config_snd_frame logger = logging.getLogger(__name__) @@ -57,8 +58,11 @@ class ReportWire(ReportWireEncoder): self._topology_seq += 1 topology = build_topology(systems, seq=self._topology_seq, ts=ts) self._last_topology = topology - logger.debug("(REPORT) TOPOLOGY_SND seq=%s", self._topology_seq) - return (_json_wire(REPORT_OPCODES["TOPOLOGY_SND"], 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) diff --git a/src/adn_server/infrastructure/twisted_adapters/report_server.py b/src/adn_server/infrastructure/twisted_adapters/report_server.py index e0c0df8..f5c9ffe 100644 --- a/src/adn_server/infrastructure/twisted_adapters/report_server.py +++ b/src/adn_server/infrastructure/twisted_adapters/report_server.py @@ -26,12 +26,17 @@ from __future__ import annotations import logging +from collections.abc import Callable from typing import Any from twisted.internet.protocol import Factory from twisted.protocols.basic import NetstringReceiver from adn_server.application.ports import ReportMqttPublisher, ReportWireEncoder +from adn_server.application.report.monitor_topology import ( + expand_inject_proxy_systems, + remap_inject_proxy_voice_event, +) from .report import REPORT_OPCODES, create_report_wire @@ -82,6 +87,7 @@ class ReportServerFactory(Factory): self.clients: list[ReportProtocol] = [] self._systems: dict[str, Any] = {} self._bridges: dict[str, Any] = {} + self._peer_slot_map: Callable[[], dict[bytes, int]] | None = None def buildProtocol(self, addr: Any) -> ReportProtocol | None: allowed = self._config.get("REPORTS", {}).get("REPORT_CLIENTS", ["127.0.0.1"]) @@ -103,6 +109,15 @@ class ReportServerFactory(Factory): def set_bridges(self, bridges: dict[str, Any]) -> None: self._bridges = bridges + def set_peer_slot_map(self, provider: Callable[[], dict[bytes, int]] | None) -> None: + """Provide proxy upstream slot indices for monitor topology expansion.""" + self._peer_slot_map = provider + + def _systems_for_report(self) -> dict[str, Any]: + systems = self._systems + peer_slots = self._peer_slot_map() if self._peer_slot_map is not None else None + return expand_inject_proxy_systems(self._config, systems, peer_slots) + def _send_frames(self, client: ReportProtocol, frames: tuple[bytes, ...]) -> None: for frame in frames: client.sendString(frame) @@ -119,27 +134,35 @@ class ReportServerFactory(Factory): def _send_hello_to(self, client: ReportProtocol) -> None: try: - self._send_frames(client, self._wire.hello_frames(self._systems)) + self._send_frames(client, self._wire.hello_frames(self._systems_for_report())) except Exception as e: logger.warning("(REPORT) Failed to send HELLO: %s", e) def _send_config_to(self, client: ReportProtocol, *, full_snapshot: bool) -> None: - self._send_frames(client, self._wire.config_frames(self._systems, full_snapshot=full_snapshot)) + self._send_frames( + client, + self._wire.config_frames(self._systems_for_report(), full_snapshot=full_snapshot), + ) def _send_bridge_to(self, client: ReportProtocol, *, full_snapshot: bool) -> None: self._send_frames(client, self._wire.bridge_frames(self._bridges, full_snapshot=full_snapshot)) def send_config(self, *, incremental: bool = False) -> None: - frames = self._wire.config_frames(self._systems, full_snapshot=not incremental) + systems = self._systems_for_report() + frames = self._wire.config_frames(systems, full_snapshot=not incremental) self._broadcast_frames(frames) if self._mqtt is not None: - self._mqtt.publish_dashboard(self._systems) + self._mqtt.publish_dashboard(systems) def send_bridge(self, *, incremental: bool = False) -> None: frames = self._wire.bridge_frames(self._bridges, full_snapshot=not incremental) self._broadcast_frames(frames) def send_bridge_event(self, event: str) -> None: + peer_slots = self._peer_slot_map() if self._peer_slot_map is not None else None + event = remap_inject_proxy_voice_event( + event, self._config, self._systems, peer_slots + ) frames = self._wire.bridge_event_frames(event) self._broadcast_frames(frames) if self._mqtt is not None: diff --git a/tests/application/test_monitor_report_contract.py b/tests/application/test_monitor_report_contract.py new file mode 100644 index 0000000..f2722e1 --- /dev/null +++ b/tests/application/test_monitor_report_contract.py @@ -0,0 +1,235 @@ +"""Report ↔ monitor contract tests (would have caught the 2026-06 production outage). + +Unit tests on expansion/topology alone are insufficient: adn-monitor deletes masters +not present in CONFIG and only adds peers under masters that already exist in CTABLE. +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 + +from adn_server.application.report.monitor_topology import ( + expand_inject_proxy_systems, + remap_inject_proxy_voice_event, +) +from adn_server.application.report.payloads import build_topology +from adn_server.domain.value_objects import bytes_4 +from adn_server.infrastructure.twisted_adapters.report.opcodes import REPORT_OPCODES +from adn_server.infrastructure.twisted_adapters.report.wire import ReportWire +from adn_server.infrastructure.twisted_adapters.report_server import ReportServerFactory +from tests.support.monitor_ctable_sim import ( + apply_config_to_ctable, + count_master_peers, + count_masters, + ctable_with_virtual_masters, + empty_ctable, + sparse_expand_buggy, + update_ctable_from_config, +) + +MAX_PEERS = 102 +BASE_PORT = 56400 + + +def _peer(peer_id: int, *, slot: int | None = None) -> dict[str, Any]: + return { + "CONNECTION": "YES", + "CONNECTED": 1_700_000_000, + "IP": "203.0.113.10", + "PORT": 62031, + "CALLSIGN": b"CE5RPY ", + } + + +def _inject_runtime_systems(peer_specs: list[tuple[int, int]]) -> dict[str, Any]: + """Runtime SYSTEM dict (single inject target) with peers keyed by radio ID.""" + peers = {bytes_4(pid): _peer(pid) for pid, _slot in peer_specs} + return { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": MAX_PEERS, + "_REPORT_BASE_PORT": BASE_PORT, + "PEERS": peers, + }, + "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, + "D-APRS-0": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, + "OBP-CL": { + "MODE": "OPENBRIDGE", + "ENABLED": True, + "NETWORK_ID": bytes_4(73010), + "PEERS": {}, + }, + } + + +def _proxy_config() -> dict[str, Any]: + return { + "REPORTS": {"REPORT": True, "REPORT_CLIENTS": ["*"]}, + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + } + + +def _peer_slots(peer_specs: list[tuple[int, int]]) -> dict[bytes, int]: + return {bytes_4(pid): slot for pid, slot in peer_specs} + + +def _production_snapshot( + peer_specs: list[tuple[int, int]], +) -> dict[str, Any]: + systems = _inject_runtime_systems(peer_specs) + return expand_inject_proxy_systems( + _proxy_config(), + systems, + _peer_slots(peer_specs), + ) + + +def _system_n_names(config: dict[str, Any], *, target: str = "SYSTEM") -> set[str]: + return {name for name in config if name.startswith(f"{target}-")} + + +def test_sparse_expansion_fails_monitor_master_preservation_contract() -> None: + """Regression: sparse SYSTEM-N-only snapshots delete 99+ masters on monitor update.""" + peer_specs = [(730039101, 4), (7301896, 6), (7301795, 1)] + runtime = _inject_runtime_systems(peer_specs) + slots = _peer_slots(peer_specs) + sparse = sparse_expand_buggy(_proxy_config(), runtime, slots) + + ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS) + before = count_masters(ctable) + update_ctable_from_config(sparse, ctable) + + assert before == MAX_PEERS + 2 # SYSTEM-0..101 + ECHO + D-APRS-0 + assert count_masters(ctable) == 2 + len(_system_n_names(sparse)) + assert count_masters(ctable) < 10 + + +def test_full_expansion_passes_monitor_master_preservation_contract() -> None: + """Healthy monitor CTABLE must survive every server report push.""" + peer_specs = [(730039101, 4), (7301896, 6), (7301795, 1)] + full = _production_snapshot(peer_specs) + + ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS) + before = count_masters(ctable) + update_ctable_from_config(full, ctable) + + assert count_masters(ctable) == before + assert len(_system_n_names(full)) == MAX_PEERS + + +def test_full_expansion_surfaces_all_hotspot_peers_on_fresh_monitor() -> None: + """First connect (empty CTABLE): all connected peers must appear under SYSTEM-N.""" + peer_specs = [ + (730050, 0), + (7301795, 1), + (7303246, 2), + (7300444, 3), + (730039101, 4), + (730266501, 5), + (7301896, 6), + ] + full = _production_snapshot(peer_specs) + ctable = empty_ctable() + apply_config_to_ctable(full, ctable) + + assert count_masters(ctable) == MAX_PEERS + 2 + assert count_master_peers(ctable) == len(peer_specs) + + +def test_truncated_monitor_ctable_cannot_show_system_n_peers() -> None: + """Documents post-outage state: update path never creates SYSTEM-N rows or their peers.""" + peer_specs = [(730039101, 4), (7301896, 6)] + full = _production_snapshot(peer_specs) + damaged = empty_ctable() + damaged["MASTERS"] = { + "ECHO": {"PEERS": {}}, + "D-APRS-0": {"PEERS": {}}, + "D-APRS-1": {"PEERS": {}}, + } + damaged["OPENBRIDGES"] = {"OBP-CL": {"STREAMS": {}}} + + update_ctable_from_config(full, damaged) + + assert count_masters(damaged) == 2 + assert count_master_peers(damaged) == 0 + assert "SYSTEM-4" not in damaged["MASTERS"] + assert "SYSTEM-6" not in damaged["MASTERS"] + assert full["SYSTEM-4"]["PEERS"] + assert full["SYSTEM-6"]["PEERS"] + + +def test_report_wire_emits_legacy_pickle_before_topology() -> None: + peer_specs = [(730039101, 4)] + snapshot = _production_snapshot(peer_specs) + frames = ReportWire().config_frames(snapshot, full_snapshot=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) + + +def test_report_factory_push_matches_monitor_contract() -> None: + """End-to-end: factory._systems_for_report + wire frames safe for monitor.""" + peer_specs = [(730039101, 4), (7301896, 6)] + factory = ReportServerFactory(_proxy_config()) + factory.set_systems(_inject_runtime_systems(peer_specs)) + 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:]) + + ctable = ctable_with_virtual_masters(max_slots=MAX_PEERS) + before = count_masters(ctable) + update_ctable_from_config(config, ctable) + + assert len(_system_n_names(snapshot)) == MAX_PEERS + assert count_masters(ctable) == before + assert count_master_peers(ctable) == len(peer_specs) + + +def test_voice_event_remap_updates_monitor_hotspot_timeslot() -> None: + """Dashboard chips need SYSTEM-N in voice events (not runtime inject SYSTEM).""" + peer_specs = [(730039101, 4)] + runtime = _inject_runtime_systems(peer_specs) + full = _production_snapshot(peer_specs) + ctable = empty_ctable() + apply_config_to_ctable(full, ctable) + + raw = "GROUP VOICE,START,RX,SYSTEM,3262598598,730039101,730039101,2,730444" + remapped = remap_inject_proxy_voice_event( + raw, _proxy_config(), runtime, _peer_slots(peer_specs) + ) + parts = remapped.split(",") + assert parts[3] == "SYSTEM-4" + assert 730039101 in ctable["MASTERS"]["SYSTEM-4"]["PEERS"] + assert "SYSTEM" not in ctable["MASTERS"] + + +@pytest.mark.parametrize("max_peers", [4, 16, 102]) +def test_expansion_always_emits_full_virtual_master_range(max_peers: int) -> None: + peer = bytes_4(730039101) + systems = { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": max_peers, + "_REPORT_BASE_PORT": BASE_PORT, + "PEERS": {peer: _peer(730039101)}, + } + } + expanded = expand_inject_proxy_systems(_proxy_config(), systems, {peer: 1}) + assert len(_system_n_names(expanded)) == max_peers + topo = build_topology(expanded, seq=1) + assert len([s for s in topo["systems"] if s["name"].startswith("SYSTEM-")]) == max_peers diff --git a/tests/application/test_monitor_topology.py b/tests/application/test_monitor_topology.py new file mode 100644 index 0000000..3687b93 --- /dev/null +++ b/tests/application/test_monitor_topology.py @@ -0,0 +1,252 @@ +"""Monitor topology parity for inject-only proxy (legacy SYSTEM-N report shape).""" + +from __future__ import annotations + +from adn_server.application.report.monitor_topology import ( + expand_inject_proxy_systems, + remap_inject_proxy_voice_event, +) +from adn_server.application.report.payloads import build_topology +from adn_server.domain.value_objects import bytes_4 + + +def _peer(*, connected: bool = True) -> dict: + return { + "CONNECTION": "YES" if connected else "NO", + "CONNECTED": 1_700_000_000 if connected else 0, + "IP": "203.0.113.10", + "PORT": 62031, + "CALLSIGN": b"CE5RPY ", + "RX_FREQ": b"145625000", + "TX_FREQ": b"145625000", + } + + +def test_expand_inject_proxy_fans_peers_into_system_n() -> None: + peer_a = bytes_4(730039101) + peer_b = bytes_4(7301896) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "_REPORT_BASE_PORT": 56400, + "PEERS": { + peer_a: _peer(), + peer_b: _peer(), + }, + }, + "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, + }, + } + expanded = expand_inject_proxy_systems( + config, + config["SYSTEMS"], + {peer_a: 2, peer_b: 66}, + ) + assert "SYSTEM" not in expanded + assert expanded["SYSTEM-2"]["PORT"] == 56402 + assert expanded["SYSTEM-66"]["PORT"] == 56466 + assert list(expanded["SYSTEM-2"]["PEERS"]) == [peer_a] + assert list(expanded["SYSTEM-66"]["PEERS"]) == [peer_b] + assert "ECHO" in expanded + + +def test_build_topology_after_expand_matches_monitor_shape() -> None: + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "_REPORT_BASE_PORT": 56400, + "PEERS": {peer: _peer()}, + } + }, + } + expanded = expand_inject_proxy_systems(config, config["SYSTEMS"], {peer: 2}) + doc = build_topology(expanded, seq=1) + names = {system["name"] for system in doc["systems"]} + assert "SYSTEM" not in names + assert "SYSTEM-2" in names + system = next(item for item in doc["systems"] if item["name"] == "SYSTEM-2") + assert system["port"] == 56402 + assert len(system["peers"]) == 1 + row = system["peers"][0] + assert row["id"] == 730039101 + assert row["connected"] is True + assert row["ip"] == "203.0.113.10" + assert row["callsign"] == "CE5RPY" + + +def test_expand_inject_proxy_emits_all_virtual_masters() -> None: + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 4, + "_REPORT_BASE_PORT": 56400, + "PEERS": {peer: _peer()}, + } + }, + } + expanded = expand_inject_proxy_systems(config, config["SYSTEMS"], {peer: 2}) + for slot in range(4): + assert f"SYSTEM-{slot}" in expanded + assert expanded[f"SYSTEM-{slot}"]["PORT"] == 56400 + slot + assert list(expanded["SYSTEM-2"]["PEERS"]) == [peer] + assert expanded["SYSTEM-0"]["PEERS"] == {} + assert expanded["SYSTEM-1"]["PEERS"] == {} + + +def test_remap_voice_event_system_to_virtual_master() -> None: + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {peer: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,RX,SYSTEM,3262598598,730039101,730039101,2,730444" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {peer: 4} + ) + assert remapped.startswith("GROUP VOICE,START,RX,SYSTEM-4,") + + +def test_remap_voice_event_tx_keeps_echo_peer_id_for_hotspot_rx_display() -> None: + """TX echo→hotspot: field 5 stays 9990 so hotspot chip is TX/green while receiving.""" + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {peer: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039101,2,730444" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {peer: 4} + ) + parts = remapped.split(",") + assert parts[3] == "SYSTEM-4" + assert parts[5] == "9990" + + +def test_remap_voice_event_rx_normalizes_field5_to_hotspot_radio_id() -> None: + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {peer: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,RX,SYSTEM,4100887026,73003,7300392,2,9990" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {peer: 4} + ) + parts = remapped.split(",") + assert parts[3] == "SYSTEM-4" + assert parts[5] == "730039101" + + +def test_remap_voice_event_matches_user_prefix_when_single_hotspot_online() -> None: + """User 7300391 with only HS1 online: legacy short subscriber id can resolve.""" + peer = bytes_4(730039101) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {peer: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300392,2,9990" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {peer: 4} + ) + parts = remapped.split(",") + assert parts[3] == "SYSTEM-4" + assert parts[5] == "9990" + + +def test_remap_voice_event_skips_ambiguous_user_with_multiple_hotspots() -> None: + """User 7300391 + HS1/HS2 online: do not pick the wrong radio from rf_src alone.""" + hs1 = bytes_4(730039101) + hs2 = bytes_4(730039102) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {hs1: _peer(), hs2: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300391,2,9990" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {hs1: 4, hs2: 5} + ) + assert remapped == raw + + +def test_remap_voice_event_resolves_full_radio_id_with_sibling_hotspots_online() -> None: + """Full radio id in rf_src still maps correctly when sibling HS are connected.""" + hs1 = bytes_4(730039101) + hs2 = bytes_4(730039102) + config = { + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + "SYSTEMS": { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": {hs1: _peer(), hs2: _peer()}, + } + }, + } + raw = "GROUP VOICE,START,TX,SYSTEM,4100887026,9990,730039102,2,730444" + remapped = remap_inject_proxy_voice_event( + raw, config, config["SYSTEMS"], {hs1: 4, hs2: 5} + ) + parts = remapped.split(",") + assert parts[3] == "SYSTEM-5" + assert parts[5] == "9990" + + +def test_remap_voice_event_passes_through_non_proxy_systems() -> None: + raw = "GROUP VOICE,START,RX,ECHO,1,9990,730039101,2,9990" + assert remap_inject_proxy_voice_event(raw, {}, {}) == raw + + +def test_non_proxy_systems_pass_through_unchanged() -> None: + systems = { + "ECHO": {"MODE": "MASTER", "ENABLED": True, "PEERS": {}}, + } + assert expand_inject_proxy_systems({"PROXY": {"TARGET_SYSTEM": "SYSTEM"}}, systems) is systems diff --git a/tests/bridge/test_report_peer_id.py b/tests/bridge/test_report_peer_id.py new file mode 100644 index 0000000..3d5950a --- /dev/null +++ b/tests/bridge/test_report_peer_id.py @@ -0,0 +1,49 @@ +"""BRDG_EVENT peer_id resolution for RX legs (hotspot transmitting).""" + +from __future__ import annotations + +from adn_server.application.bridge.helpers import resolve_voice_peer_id +from adn_server.domain.value_objects import bytes_3, bytes_4 + + +def test_resolve_voice_peer_uses_rf_src_when_field5_is_network_id() -> None: + hotspot = bytes_4(730039101) + systems = { + "SYSTEM": { + "PEERS": { + hotspot: {"CONNECTION": "YES"}, + } + } + } + network_peer = bytes_4(73003) + resolved = resolve_voice_peer_id( + network_peer, + bytes_3(730039101), + "SYSTEM", + systems, + ) + assert resolved == hotspot + + +def test_resolve_voice_peer_keeps_known_hotspot_id() -> None: + hotspot = bytes_4(730039102) + systems = {"SYSTEM": {"PEERS": {hotspot: {"CONNECTION": "YES"}}}} + resolved = resolve_voice_peer_id( + hotspot, + bytes_3(730039102), + "SYSTEM", + systems, + ) + assert resolved == hotspot + + +def test_resolve_voice_peer_resolves_network_prefix_for_single_hotspot() -> None: + hotspot = bytes_4(730039101) + systems = {"SYSTEM": {"PEERS": {hotspot: {"CONNECTION": "YES"}}}} + resolved = resolve_voice_peer_id( + bytes_4(73003), + bytes_3(7300392), + "SYSTEM", + systems, + ) + assert resolved == hotspot diff --git a/tests/infrastructure/test_report_server_wire.py b/tests/infrastructure/test_report_server_wire.py index 3a0d007..ea50b0b 100644 --- a/tests/infrastructure/test_report_server_wire.py +++ b/tests/infrastructure/test_report_server_wire.py @@ -4,6 +4,7 @@ from __future__ import annotations import json +from adn_server.domain.value_objects import bytes_4 from adn_server.infrastructure.twisted_adapters.report_server import ( REPORT_OPCODES, ReportServerFactory, @@ -42,14 +43,15 @@ def test_connect_sends_hello_topology_and_routing() -> None: factory._send_config_to(client, full_snapshot=True) factory._send_bridge_to(client, full_snapshot=True) - assert len(client.messages) == 3 + assert len(client.messages) == 4 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["TOPOLOGY_SND"] - assert json.loads(client.messages[1][1:])["type"] == "topology" - assert client.messages[2][:1] == REPORT_OPCODES["ROUTING_TABLE_SND"] - assert json.loads(client.messages[2][1:])["type"] == "routing_table" + 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" def test_bridge_event_emits_voice_event_json() -> None: @@ -90,3 +92,84 @@ def test_incremental_bridge_update_sends_delta() -> None: assert delta["type"] == "delta" assert delta["patch"]["type"] == "routing_table" + +def test_inject_proxy_topology_expanded_for_monitor() -> None: + peer = bytes_4(730039101) + config = { + "REPORTS": {"REPORT": True, "REPORT_CLIENTS": ["*"]}, + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + } + factory = ReportServerFactory(config) + factory.set_systems( + { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "_REPORT_BASE_PORT": 56400, + "PEERS": { + peer: { + "CONNECTION": "YES", + "CONNECTED": 1_700_000_000, + "IP": "203.0.113.10", + "PORT": 62031, + "CALLSIGN": b"CE5RPY ", + } + }, + } + } + ) + 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"]} + 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") + 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"] == [] + + +def test_send_bridge_event_remaps_inject_proxy_system_name() -> None: + peer = bytes_4(730039101) + config = { + "REPORTS": {"REPORT": True, "REPORT_CLIENTS": ["*"]}, + "PROXY": {"TARGET_SYSTEM": "SYSTEM"}, + } + factory = ReportServerFactory(config) + factory.set_systems( + { + "SYSTEM": { + "MODE": "MASTER", + "ENABLED": True, + "MAX_PEERS": 102, + "PEERS": { + peer: { + "CONNECTION": "YES", + "CONNECTED": 1_700_000_000, + "IP": "203.0.113.10", + "PORT": 62031, + } + }, + } + } + ) + factory.set_peer_slot_map(lambda: {peer: 4}) + client = _CapturingClient() + factory.clients.append(client) + factory.send_bridge_event( + "GROUP VOICE,START,TX,SYSTEM,2693411696,9990,7300392,2,9990" + ) + assert len(client.messages) == 1 + payload = json.loads(client.messages[0][1:].decode()) + assert payload["type"] == "voice_event" + assert payload["system"] == "SYSTEM-4" + assert payload["peer_id"] == 9990 + diff --git a/tests/support/monitor_ctable_sim.py b/tests/support/monitor_ctable_sim.py new file mode 100644 index 0000000..caa15f6 --- /dev/null +++ b/tests/support/monitor_ctable_sim.py @@ -0,0 +1,144 @@ +"""Minimal adn-monitor CTABLE semantics for report contract tests. + +Mirrors the behaviour that broke production (``update_hblink_table_impl`` in +adn-monitor): masters missing from CONFIG are deleted; peers on ``SYSTEM-N`` are +only visible when that master already exists in CTABLE (update does not create +new master rows). +""" + +from __future__ import annotations + +import copy +from typing import Any + +from adn_server.domain.value_objects import bytes_4, int_id + + +def empty_ctable() -> dict[str, Any]: + return {"MASTERS": {}, "PEERS": {}, "OPENBRIDGES": {}} + + +def build_ctable_from_config(config: dict[str, Any]) -> dict[str, Any]: + """Monitor ``build_hblink_table`` path (empty CTABLE on first connect).""" + ctable = empty_ctable() + for name, data in config.items(): + if not data.get("ENABLED", True): + continue + mode = data.get("MODE") + if mode == "MASTER": + ctable["MASTERS"][name] = {"PEERS": {}} + for peer_key, peer_conf in data.get("PEERS", {}).items(): + if peer_conf.get("CONNECTION") == "YES": + _add_peer(ctable["MASTERS"][name]["PEERS"], peer_key) + elif mode == "OPENBRIDGE": + ctable["OPENBRIDGES"][name] = {"STREAMS": {}} + return ctable + + +def update_ctable_from_config(config: dict[str, Any], ctable: dict[str, Any]) -> None: + """Monitor ``update_hblink_table`` path (CTABLE already populated).""" + for key in ("MASTERS", "PEERS", "OPENBRIDGES"): + for name in list(ctable.get(key, {})): + if name not in config: + del ctable[key][name] + for name, data in config.items(): + if data.get("MODE") != "MASTER": + continue + # Peers attach only to masters that already exist — no master creation here. + masters_peers = ctable["MASTERS"].get(name, {}).get("PEERS", {}) + for peer_key, peer_conf in data.get("PEERS", {}).items(): + if peer_conf.get("CONNECTION") != "YES": + continue + pid = int_id(peer_key) + if pid not in masters_peers: + _add_peer(masters_peers, peer_key) + for peer_id in list(masters_peers): + if bytes_4(peer_id) not in data.get("PEERS", {}): + del masters_peers[peer_id] + + +def apply_config_to_ctable(config: dict[str, Any], ctable: dict[str, Any]) -> None: + """``_apply_config_to_state``: build when empty, else update.""" + if ctable["MASTERS"]: + update_ctable_from_config(config, ctable) + else: + built = build_ctable_from_config(config) + ctable["MASTERS"] = built["MASTERS"] + ctable["PEERS"] = built["PEERS"] + ctable["OPENBRIDGES"] = built["OPENBRIDGES"] + + +def count_masters(ctable: dict[str, Any]) -> int: + return len(ctable.get("MASTERS", {})) + + +def count_master_peers(ctable: dict[str, Any]) -> int: + total = 0 + for master in ctable.get("MASTERS", {}).values(): + total += len(master.get("PEERS", {})) + return total + + +def ctable_with_virtual_masters( + *, + target: str = "SYSTEM", + max_slots: int, + extra: tuple[str, ...] = ("ECHO", "D-APRS-0", "OBP-CL"), +) -> dict[str, Any]: + """CTABLE like a healthy monitor before the sparse-topology regression.""" + ctable = empty_ctable() + for slot in range(max_slots): + ctable["MASTERS"][f"{target}-{slot}"] = {"PEERS": {}} + for name in extra: + if name == "OBP-CL": + ctable["OPENBRIDGES"][name] = {"STREAMS": {}} + else: + ctable["MASTERS"][name] = {"PEERS": {}} + return ctable + + +def sparse_expand_buggy( + config: dict[str, Any], + systems: dict[str, Any], + peer_slots: dict[bytes, int] | None, +) -> dict[str, Any]: + """Old broken behaviour: only ``SYSTEM-N`` slots that carry peers.""" + import copy as _copy + + from adn_server.application.proxy.deployment import is_proxy_inject_only, proxy_target_system + from adn_server.application.report.monitor_topology import ( + _connected_peers, + _resolve_slot_map, + ) + + target = proxy_target_system(config) + if not target or not is_proxy_inject_only(config, target): + return systems + sys_cfg = systems.get(target) + if not isinstance(sys_cfg, dict) or sys_cfg.get("MODE") != "MASTER": + return systems + peers = sys_cfg.get("PEERS", {}) + max_slots = int(sys_cfg.get("MAX_PEERS", 1)) + base_port = int(sys_cfg.get("_REPORT_BASE_PORT", 56400)) + connected = _connected_peers(peers) + slot_map = _resolve_slot_map(connected, peer_slots, max_slots=max_slots) + out = {name: cfg for name, cfg in systems.items() if name != target} + for slot in sorted({s for s in slot_map.values()}): + virtual_name = f"{target}-{slot}" + virtual = _copy.deepcopy(sys_cfg) + virtual["PORT"] = base_port + slot + virtual["PEERS"] = { + peer_key: peers[peer_key] + for peer_key, mapped in slot_map.items() + if mapped == slot and peer_key in peers + } + out[virtual_name] = virtual + return out + + +def deep_copy_ctable(ctable: dict[str, Any]) -> dict[str, Any]: + return copy.deepcopy(ctable) + + +def _add_peer(peers: dict[int, dict[str, str]], peer_key: Any) -> None: + peers[int_id(peer_key)] = {"CONNECTION": "YES"}