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

@ -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] == "#":

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

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

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

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

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

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

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

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

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

@ -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"}
Loading…
Cancel
Save

Powered by TurnKey Linux.