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

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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.