From 2235203e2248a3f2ea15dfa57d03020a9efb18de Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 12 Jun 2026 01:58:25 -0400 Subject: [PATCH] perf: index inject-only downlink peers by slot and TGID Narrow send_peers and REPEAT fan-out on inject-only proxies using a precomputed OPTIONS index, cache connected peer count, and debounce CONFIG_SND pushes to the monitor. --- .../routing/peer_downlink_index.py | 151 ++++++++++++++++++ .../twisted_adapters/udp_hbp.py | 88 ++++++++-- .../test_peer_downlink_fanout.py | 54 +++++++ tests/routing/test_peer_downlink_index.py | 45 ++++++ tests/support/hbp_repeat_stack.py | 2 + 5 files changed, 324 insertions(+), 16 deletions(-) create mode 100644 src/adn_server/application/routing/peer_downlink_index.py create mode 100644 tests/infrastructure/test_peer_downlink_fanout.py create mode 100644 tests/routing/test_peer_downlink_index.py diff --git a/src/adn_server/application/routing/peer_downlink_index.py b/src/adn_server/application/routing/peer_downlink_index.py new file mode 100644 index 0000000..587b00f --- /dev/null +++ b/src/adn_server/application/routing/peer_downlink_index.py @@ -0,0 +1,151 @@ +"""Inject-only MASTER downlink: narrow peer fan-out by (slot, TGID). + +Legacy ``send_peers`` scans every registered peer per packet. On inject-only +proxies with hundreds of hotspots, that is O(peers × pkt/s). This index +builds a candidate set from static OPTIONS and UA session state; each +candidate is still checked with :func:`peer_should_receive_group_voice`. +""" + +from __future__ import annotations + +import time +from dataclasses import dataclass, field +from typing import Any + +from ...domain import bytes_4, int_id +from .helpers import peer_single_exclusive_tgid + + +def invalidate_peer_options_cache(peer: dict[str, Any]) -> None: + """Drop cached OPTIONS parse after RPTO.""" + peer.pop("_CACHED_OPTIONS_STATIC", None) + + +def cached_peer_static_tgs(peer: dict[str, Any]) -> tuple[tuple[str, ...], tuple[str, ...]]: + """Memoize ``parse_peer_options_static`` per peer OPTIONS blob.""" + opts = peer.get("OPTIONS") + key = opts if isinstance(opts, bytes) else b"" + cached = peer.get("_CACHED_OPTIONS_STATIC") + if cached and cached[0] == key: + return cached[1], cached[2] + from adn_server.application.report.payloads import parse_peer_options_static + + ts1, ts2 = parse_peer_options_static(opts) + t1, t2 = tuple(ts1), tuple(ts2) + peer["_CACHED_OPTIONS_STATIC"] = (key, t1, t2) + return t1, t2 + + +def count_connected_peers(peers: dict[bytes, dict[str, Any]]) -> int: + return sum(1 for p in peers.values() if p.get("CONNECTION") == "YES") + + +@dataclass +class PeerDownlinkIndex: + """Precomputed (slot, TGID) → peer candidates for connected hotspots.""" + + static_by_slot_tgid: dict[tuple[int, int], frozenset[bytes]] = field(default_factory=dict) + ua_by_slot_tgid: dict[tuple[int, int], frozenset[bytes]] = field(default_factory=dict) + connected: frozenset[bytes] = frozenset() + + def candidates(self, slot: int, tgid: int, *, connected_count: int) -> frozenset[bytes]: + if connected_count == 1: + return self.connected + out: set[bytes] = set() + key = (int(slot), int(tgid)) + out.update(self.static_by_slot_tgid.get(key, ())) + out.update(self.ua_by_slot_tgid.get(key, ())) + return frozenset(out) + + +def _add_index_entry( + index: dict[tuple[int, int], set[bytes]], + slot: int, + tgid: int, + peer_id: bytes, +) -> None: + try: + tgid_i = int(tgid) + except (TypeError, ValueError): + return + if tgid_i <= 0: + return + index.setdefault((int(slot), tgid_i), set()).add(peer_id) + + +def build_peer_downlink_index( + peers: dict[bytes, dict[str, Any]], + sys_cfg: dict[str, Any], + *, + now: float | None = None, +) -> PeerDownlinkIndex: + """Rebuild candidate map from all connected peers (call when index is dirty).""" + pkt_time = time.time() if now is None else now + static_map: dict[tuple[int, int], set[bytes]] = {} + ua_map: dict[tuple[int, int], set[bytes]] = {} + connected: set[bytes] = set() + + for peer_id, peer in peers.items(): + if peer.get("CONNECTION") != "YES": + continue + connected.add(peer_id) + ts1, ts2 = cached_peer_static_tgs(peer) + for tg in ts1: + _add_index_entry(static_map, 1, tg, peer_id) + for tg in ts2: + _add_index_entry(static_map, 2, tg, peer_id) + + for slot in (1, 2): + locked = peer_single_exclusive_tgid( + peer, slot, sys_cfg, peer_id=peer_id, now=pkt_time, + ) + if locked is not None: + _add_index_entry(ua_map, slot, locked, peer_id) + + sessions = peer.get("_UA_SESSION") + if isinstance(sessions, dict): + for slot, entry in sessions.items(): + if not isinstance(entry, dict): + continue + if pkt_time >= float(entry.get("expires", 0)): + continue + locked = entry.get("tgid") + if locked is not None: + _add_index_entry(ua_map, int(slot), locked, peer_id) + + multi_store = sys_cfg.get("_PEER_UA_MULTI_TGS") + if isinstance(multi_store, dict): + for pk, per_slot in multi_store.items(): + if not isinstance(per_slot, dict): + continue + peer_id = pk if isinstance(pk, bytes) else bytes_4(int_id(pk)) + if peer_id not in connected: + continue + for slot, tg_set in per_slot.items(): + if not isinstance(tg_set, set): + continue + for tgid in tg_set: + _add_index_entry(ua_map, int(slot), tgid, peer_id) + + ua_sessions = sys_cfg.get("_PEER_UA_SESSIONS") + if isinstance(ua_sessions, dict): + for pk, per_slot in ua_sessions.items(): + if not isinstance(per_slot, dict): + continue + peer_id = pk if isinstance(pk, bytes) else bytes_4(int_id(pk)) + if peer_id not in connected: + continue + for slot, entry in per_slot.items(): + if not isinstance(entry, dict): + continue + if pkt_time >= float(entry.get("expires", 0)): + continue + locked = entry.get("tgid") + if locked is not None: + _add_index_entry(ua_map, int(slot), locked, peer_id) + + return PeerDownlinkIndex( + static_by_slot_tgid={k: frozenset(v) for k, v in static_map.items()}, + ua_by_slot_tgid={k: frozenset(v) for k, v in ua_map.items()}, + connected=frozenset(connected), + ) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 0f5b295..9cd8892 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -53,6 +53,11 @@ from ...application.routing.helpers import ( seed_peer_ua_session_from_status, tg4000_reset_on_vhead, ) +from ...application.routing.peer_downlink_index import ( + build_peer_downlink_index, + count_connected_peers, + invalidate_peer_options_cache, +) from ...application.proxy.deployment import is_proxy_inject_only from ...domain import bytes_4, int_id from ...domain.talker_alias import ( @@ -219,6 +224,11 @@ class HBPProtocol(DatagramProtocol): self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} if self._config.get("MODE") == "MASTER": self._peers = self._config.setdefault("PEERS", {}) + self._downlink_index_dirty = True + self._downlink_index = None + self._connected_peer_count = 0 + self._config_push_delayed = None + self._refresh_connected_peer_count() else: self._peers = {} if self._config.get("MODE") in ("MASTER", "OPENBRIDGE"): @@ -274,19 +284,60 @@ class HBPProtocol(DatagramProtocol): self._config = sys_cfg if sys_cfg.get("MODE") == "MASTER": self._peers = sys_cfg.setdefault("PEERS", {}) + self._refresh_connected_peer_count() + self._mark_downlink_index_dirty() # ── Exact port of hblink.py send_peers / send_peer / send_master / send_system ── def send_peers(self, _packet: bytes, _hops: bytes = b"", _ber: bytes = b"\x00", _rssi: bytes = b"\x00", _source_server: bytes = b"\x00\x00\x00\x00", _source_rptr: bytes = b"\x00\x00\x00\x00") -> None: - for _peer in self._peers: - if len(_packet) < 54: - _packet = b"".join([_packet, _ber, _rssi]) + if len(_packet) < 54: + _packet = b"".join([_packet, _ber, _rssi]) + for _peer in self._iter_downlink_peers(_packet): self.send_peer(_peer, _packet) def _inject_multi_peer_options_filter(self) -> bool: """Inject-only proxy: always filter downlink by each peer's own OPTIONS.""" return is_proxy_inject_only(self._CONFIG, self._system) + def _mark_downlink_index_dirty(self) -> None: + self._downlink_index_dirty = True + + def _refresh_connected_peer_count(self) -> None: + self._connected_peer_count = count_connected_peers(self._peers) + + def _cached_connected_peer_count(self) -> int: + """Return cached count; refresh once if cache is zero but peers exist.""" + n = self._connected_peer_count + if n <= 0 and self._peers: + self._refresh_connected_peer_count() + n = self._connected_peer_count + return n + + def _ensure_downlink_index(self): + if not self._downlink_index_dirty and self._downlink_index is not None: + return self._downlink_index + self._downlink_index = build_peer_downlink_index(self._peers, self._config) + self._downlink_index_dirty = False + return self._downlink_index + + def _iter_downlink_peers(self, packet: bytes): + """Peer ids to consider for MASTER downlink / REPEAT (indexed when inject-only).""" + if not self._inject_multi_peer_options_filter(): + return self._peers.keys() + parsed = parse_dmrd_route_fields(packet) + if parsed is None: + return self._peers.keys() + slot, tgid, call_type = parsed + if call_type not in ("group", "vcsbk"): + return self._peers.keys() + if is_special_tg(str(tgid)): + return self._peers.keys() + connected = self._cached_connected_peer_count() + if connected <= 0: + return () + index = self._ensure_downlink_index() + return index.candidates(slot, tgid, connected_count=connected) + def _peer_mesh_config(self) -> PeerMeshConfig: _global = self._CONFIG.get("GLOBAL", {}) _sid = _global.get("SERVER_ID", b"\x00\x00\x00\x00") @@ -386,10 +437,7 @@ class HBPProtocol(DatagramProtocol): return False parsed = parse_dmrd_route_fields(packet) if parsed is None: - connected = sum( - 1 for peer in self._peers.values() if peer.get("CONNECTION") == "YES" - ) - return connected <= 1 + return self._cached_connected_peer_count() <= 1 slot, tgid, call_type = parsed if call_type not in ("group", "vcsbk"): return True @@ -402,13 +450,8 @@ class HBPProtocol(DatagramProtocol): return bytes_4(int_id(peer_id)) == bytes_4(int_id(rx_peer)) if len(packet) >= 8 and peer_matches_rf_source(peer_id, packet[5:8], self._peers): return True - connected = sum( - 1 for peer in self._peers.values() if peer.get("CONNECTION") == "YES" - ) - return connected == 1 - connected = sum( - 1 for peer in self._peers.values() if peer.get("CONNECTION") == "YES" - ) + return self._cached_connected_peer_count() == 1 + connected = self._cached_connected_peer_count() store = self._get_subscription_store() if self._get_subscription_store else None return peer_should_receive_group_voice( self._peers[peer_id], @@ -690,7 +733,13 @@ class HBPProtocol(DatagramProtocol): self.proxy_IPBlackList(_pi, self._peers[_pi]["SOCKADDR"]) def _push_config_to_monitor(self) -> None: - """Push CONFIG_SND when MASTER peer list or connection state changes.""" + """Schedule debounced CONFIG_SND when MASTER peer list or OPTIONS change.""" + if self._config_push_delayed is not None: + return + self._config_push_delayed = reactor.callLater(0.3, self._flush_config_to_monitor) + + def _flush_config_to_monitor(self) -> None: + self._config_push_delayed = None report = self._report if report is None: return @@ -759,6 +808,8 @@ class HBPProtocol(DatagramProtocol): def _remove_peer(self, peer_id: bytes) -> None: self._on_peer_disconnected(peer_id) self._peers.pop(peer_id, None) + self._refresh_connected_peer_count() + self._mark_downlink_index_dirty() def _master_datagram_received(self, _data: bytes, _sockaddr: tuple[str, int]) -> None: """Direct port of hblink.py master_datagramReceived (lines 888-1146).""" @@ -861,6 +912,7 @@ class HBPProtocol(DatagramProtocol): self._config, now=pkt_time, ) + self._mark_downlink_index_dirty() if _prev_single_tg != _int_dst_id: self._push_config_to_monitor() if ( @@ -898,7 +950,7 @@ class HBPProtocol(DatagramProtocol): ) _repeat_tail = b"".join([_data[15:20], _dmrpkt_out, _data[53:]]) _repeat_pkt = b"".join([_data[:11], _peer_id, _repeat_tail]) - for _peer in self._peers: + for _peer in self._iter_downlink_peers(_repeat_pkt): if _peer != _peer_id: self.send_peer(_peer, _repeat_pkt) # TG 4000: reset after REPEAT so peers see the packet (legacy order) @@ -1117,6 +1169,8 @@ class HBPProtocol(DatagramProtocol): else: self.send_peer(_peer_id, b"".join([RPTACK, _peer_id])) logger.info("(%s) Peer %s (%s) has sent repeater configuration, Package ID: %s, Software ID: %s, Desc: %s", self._system, _this_peer["CALLSIGN"], _this_peer["RADIO_ID"], self._peers[_peer_id]["PACKAGE_ID"].decode("utf8", errors="replace").rstrip(), self._peers[_peer_id]["SOFTWARE_ID"].decode("utf8", errors="replace").rstrip(), self._peers[_peer_id]["DESCRIPTION"].decode("utf8", errors="replace").rstrip()) + self._refresh_connected_peer_count() + self._mark_downlink_index_dirty() self._push_config_to_monitor() else: self.transport.write(b"".join([MSTNAK, _peer_id]), _sockaddr) @@ -1127,6 +1181,8 @@ class HBPProtocol(DatagramProtocol): if _peer_id in self._peers and self._peers[_peer_id]["SOCKADDR"] == _sockaddr: _this_peer = self._peers[_peer_id] _this_peer["OPTIONS"] = _data[8:] + invalidate_peer_options_cache(_this_peer) + self._mark_downlink_index_dirty() self.send_peer(_peer_id, b"".join([RPTACK, _peer_id])) logger.info("(%s) Peer %s has sent options %s", self._system, _this_peer["CALLSIGN"], _this_peer["OPTIONS"]) if is_proxy_inject_only(self._CONFIG, self._system): diff --git a/tests/infrastructure/test_peer_downlink_fanout.py b/tests/infrastructure/test_peer_downlink_fanout.py new file mode 100644 index 0000000..bf9400d --- /dev/null +++ b/tests/infrastructure/test_peer_downlink_fanout.py @@ -0,0 +1,54 @@ +"""Inject-only MASTER: downlink index limits send_peer fan-out.""" + +from __future__ import annotations + +import pytest + +from adn_server.domain import bytes_4 +from adn_server.infrastructure.hbp_constants import DMRD +from tests.harness.deterministic import DeterministicScenario, PacketSpec +from tests.support.hbp_repeat_stack import build_hbp_repeat_stack + +pytestmark = pytest.mark.integration + +_TG = 52090 + + +def _voice_burst(peer_id: bytes, rf_src: int = 7300444) -> bytes: + spec = PacketSpec( + peer_id=int.from_bytes(peer_id, "big"), + rf_src=rf_src, + dst_id=_TG, + slot=2, + stream_id=0x11223344, + payload=b"\x00" * 33, + ) + return DeterministicScenario.voice_burst_spec(spec, seq=1, dtype_vseq=1).data() + + +def test_send_peers_scales_with_matching_options_not_total_peers() -> None: + stack = build_hbp_repeat_stack(talker_alias=False, system_name="MASTER-A") + stack.config["PROXY"] = {"TARGET_SYSTEM": "MASTER-A"} + stack.hbp._CONFIG = stack.config + + tx = bytes_4(730044401) + match = bytes_4(730044402) + stack.register_peer(tx, ("10.0.0.10", 62010), options=f"TS2={_TG};") + stack.register_peer(match, ("10.0.0.11", 62011), options=f"TS2={_TG};") + + extras = [] + for i in range(50): + pid = bytes_4(730100000 + i) + extras.append(pid) + stack.register_peer(pid, (f"10.1.0.{i}", 62100 + i), options="TS2=91;") + + stack.hbp._refresh_connected_peer_count() + stack.hbp._mark_downlink_index_dirty() + + burst = _voice_burst(tx) + stack.transport.clear() + stack.hbp.send_peers(burst) + + dmrd_sends = [pkt for pkt, _ in stack.transport.sent if pkt[:4] == DMRD] + # TX peer + one matching peer (not 50 unrelated hotspots). + assert len(dmrd_sends) == 2 diff --git a/tests/routing/test_peer_downlink_index.py b/tests/routing/test_peer_downlink_index.py new file mode 100644 index 0000000..07a5e27 --- /dev/null +++ b/tests/routing/test_peer_downlink_index.py @@ -0,0 +1,45 @@ +"""Peer downlink index: inject-only fan-out narrowing.""" + +from __future__ import annotations + +from adn_server.application.routing.peer_downlink_index import ( + build_peer_downlink_index, + cached_peer_static_tgs, + invalidate_peer_options_cache, +) +from adn_server.domain import bytes_4 + + +def _peer(options: str) -> dict: + return {"CONNECTION": "YES", "OPTIONS": options.encode()} + + +def test_build_index_maps_static_tgs_per_slot() -> None: + p1 = bytes_4(730044401) + p2 = bytes_4(730044402) + peers = { + p1: _peer("TS2=52090,314569;"), + p2: _peer("TS2=730170;"), + } + idx = build_peer_downlink_index(peers, {}) + assert p1 in idx.candidates(2, 52090, connected_count=5) + assert p2 not in idx.candidates(2, 52090, connected_count=5) + assert p2 in idx.candidates(2, 730170, connected_count=5) + assert len(idx.candidates(2, 52090, connected_count=5)) == 1 + + +def test_single_connected_peer_returns_all() -> None: + p1 = bytes_4(730044401) + peers = {p1: _peer("TS2=91;")} + idx = build_peer_downlink_index(peers, {}) + assert idx.candidates(2, 52090, connected_count=1) == frozenset({p1}) + + +def test_options_cache_invalidates_on_rpto() -> None: + peer = _peer("TS2=52090;") + cached_peer_static_tgs(peer) + assert "_CACHED_OPTIONS_STATIC" in peer + peer["OPTIONS"] = b"TS2=730170;" + invalidate_peer_options_cache(peer) + ts1, ts2 = cached_peer_static_tgs(peer) + assert "730170" in ts2 diff --git a/tests/support/hbp_repeat_stack.py b/tests/support/hbp_repeat_stack.py index 605bebc..ddb919b 100644 --- a/tests/support/hbp_repeat_stack.py +++ b/tests/support/hbp_repeat_stack.py @@ -73,6 +73,8 @@ class HbpRepeatStack: self.config["SYSTEMS"][self.system_name].setdefault("PEERS", {})[peer_id] = ( self.hbp._peers[peer_id] ) + self.hbp._refresh_connected_peer_count() + self.hbp._mark_downlink_index_dirty() def inject(self, packet: bytes, sockaddr: tuple[str, int]) -> None: self.hbp.datagramReceived(packet, sockaddr)