From cb70a1cae6547769418b4a505fbcbfbdf63c8b58 Mon Sep 17 00:00:00 2001 From: yo Date: Wed, 23 Sep 2026 22:01:22 +0200 Subject: [PATCH] perf(hbp): parse a peer's OPTIONS and RF mode when they change, not per frame The MASTER ingress path asks the same two questions of every voice frame. peer_single_mode parses the peer's OPTIONS blob through parse_peer_options_fields, and peer_rf_mode falls back to derive_peer_rf_mode because nothing ever wrote the RF_MODE key its docstring anticipates. Around a static-TG hotspot the downlink helpers add five more parse_peer_options_static calls on the same blob, and each of those parses it twice (peer_options_static_valid, then _parse_options_kv). A hotspot with TS2_1=214;TS1_1=91;SINGLE=1;TIMER=15; had that string re-parsed more than ten times per frame. None of it moves at frame rate: OPTIONS only changes on RPTO, which already drops the memo added for cached_peer_static_tgs, and SLOTS with the two frequencies only arrive in RPTC. So peer_options_fields memoizes against the blob the same way, the five remaining parse_peer_options_static call sites go through cached_peer_static_tgs, and RPTC classifies the RF mode once at login. Measured on the MASTER ingress path (tests/support/hbp_repeat_stack, 20k group voice frames from a static-TG hotspot, ACLs on): 138.8us -> 39.1us per frame, or 42.0us for a peer that has not sent RPTC yet and still derives its mode. Emitted packets are byte-identical across 72 combinations of OPTIONS blobs (including invalid ones and PASS=) and simplex/duplex peers. Co-Authored-By: Claude Opus 5 --- .../application/routing/downlink.py | 4 +- src/adn_server/application/routing/helpers.py | 37 +++-- .../routing/peer_downlink_index.py | 3 +- .../twisted_adapters/udp_hbp.py | 6 + tests/routing/test_peer_options_cache.py | 134 ++++++++++++++++++ 5 files changed, 169 insertions(+), 15 deletions(-) create mode 100644 tests/routing/test_peer_options_cache.py diff --git a/src/adn_server/application/routing/downlink.py b/src/adn_server/application/routing/downlink.py index 1bc6f5e..30126dd 100644 --- a/src/adn_server/application/routing/downlink.py +++ b/src/adn_server/application/routing/downlink.py @@ -100,9 +100,9 @@ def peer_listen_slots(peer: dict[str, Any], tgid: int) -> list[int]: RF timeslots and is expected to key up on both, same as a real repeater configured that way. Simplex peers/bridges always collapse to one slot. """ - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) if peer_is_simplex(peer): tg = str(tgid) if tg in ts1 or tg in ts2: diff --git a/src/adn_server/application/routing/helpers.py b/src/adn_server/application/routing/helpers.py index c4ae2f9..587cbcf 100644 --- a/src/adn_server/application/routing/helpers.py +++ b/src/adn_server/application/routing/helpers.py @@ -1318,10 +1318,23 @@ def _system_has_active_bridge_leg( def peer_options_fields(peer: dict[str, Any]) -> dict[str, Any]: - """Parse hotspot OPTIONS into fields used by SINGLE/TIMER resolution.""" + """Parse hotspot OPTIONS into fields used by SINGLE/TIMER resolution. + + Memoized against the OPTIONS blob, the way ``cached_peer_static_tgs`` already + memoizes the static lists: OPTIONS only changes on RPTO, which drops the cache + (``invalidate_peer_options_cache``), while ingress asks for these fields + several times per voice frame through ``peer_single_mode``. + """ + opts = peer.get("OPTIONS") + key = opts if isinstance(opts, bytes) else b"" + cached = peer.get("_CACHED_OPTIONS_FIELDS") + if cached is not None and cached[0] == key: + return cached[1] from adn_server.application.report.payloads import parse_peer_options_fields - return parse_peer_options_fields(peer.get("OPTIONS")) + fields = parse_peer_options_fields(opts) + peer["_CACHED_OPTIONS_FIELDS"] = (key, fields) + return fields def _peer_ua_session_entry( @@ -1381,9 +1394,9 @@ def _peer_static_tg_blocks_slot(peer: dict[str, Any], slot: int, tgid: int) -> b match on one slot must not block genuinely independent dynamic activity on the *other* slot (e.g. TG static on TS2, this same peer separately keying up the same TG on TS1).""" - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) tg = str(tgid) if peer_is_simplex(peer): return tg in ts1 or tg in ts2 @@ -1868,9 +1881,9 @@ def peer_single_blocks_foreign_same_tg_downlink( def peer_static_options_tg_count(peer: dict[str, Any]) -> int: """Count distinct static group TGs listed in peer OPTIONS (TS1 ∪ TS2).""" - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) return len(set(ts1) | set(ts2)) @@ -1890,18 +1903,18 @@ def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool slot while self-service lists the TG on the other. """ del slot - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) tg = str(tgid) return tg in ts1 or tg in ts2 def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None: """Timeslot (1 or 2) where peer OPTIONS list ``tgid``, when unambiguous.""" - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) tg = str(tgid) in_ts1 = tg in ts1 in_ts2 = tg in ts2 @@ -1958,9 +1971,9 @@ def peer_downlink_voice_slot( static = peer_options_static_tg_slot(peer, tgid) if static is not None: return static - from adn_server.application.report.payloads import parse_peer_options_static + from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs - ts1, ts2 = parse_peer_options_static(peer.get("OPTIONS")) + ts1, ts2 = cached_peer_static_tgs(peer) tg = str(tgid) if tg in ts1 and tg in ts2: # Static on both slots (peer_options_static_tg_slot returns None diff --git a/src/adn_server/application/routing/peer_downlink_index.py b/src/adn_server/application/routing/peer_downlink_index.py index 69bccb1..b6f5585 100644 --- a/src/adn_server/application/routing/peer_downlink_index.py +++ b/src/adn_server/application/routing/peer_downlink_index.py @@ -37,8 +37,9 @@ from .helpers import peer_single_exclusive_tgid def invalidate_peer_options_cache(peer: dict[str, Any]) -> None: - """Drop cached OPTIONS parse after RPTO.""" + """Drop cached OPTIONS parses after RPTO.""" peer.pop("_CACHED_OPTIONS_STATIC", None) + peer.pop("_CACHED_OPTIONS_FIELDS", None) def cached_peer_static_tgs(peer: dict[str, Any]) -> tuple[tuple[str, ...], tuple[str, ...]]: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 0e65b0a..ef34174 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -52,6 +52,7 @@ from ...application.routing.downlink import ( from ...application.routing.helpers import ( clear_peer_rx_status_slots, clear_peer_ua_sessions, + derive_peer_rf_mode, hbp_master_ingress_repeat_allowed, is_on_demand_service_dst, is_server_originated_voice, @@ -1732,6 +1733,11 @@ class HBPProtocol(DatagramProtocol): _this_peer["URL"] = _data[98:222] _this_peer["SOFTWARE_ID"] = _data[222:262] _this_peer["PACKAGE_ID"] = _data[262:302] + # RPTC carries the only inputs simplex/duplex depends on (SLOTS and + # the two frequencies), so classify here instead of on every frame: + # peer_rf_mode() reads this and the downlink path asks it three + # times per voice frame. + _this_peer["RF_MODE"] = derive_peer_rf_mode(_this_peer) _sent_call = _rptc_field_str(_this_peer["CALLSIGN"]) if ("ALLOW_UNREG_ID" in self._config and not self._config["ALLOW_UNREG_ID"]) and _sent_call != self.validate_id(_peer_id): self._remove_peer(_peer_id) diff --git a/tests/routing/test_peer_options_cache.py b/tests/routing/test_peer_options_cache.py new file mode 100644 index 0000000..67d50ef --- /dev/null +++ b/tests/routing/test_peer_options_cache.py @@ -0,0 +1,134 @@ +# ADN DMR Peer Server - tests peer OPTIONS / RF mode caching +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Per-peer OPTIONS fields and RF mode are parsed on change, not per frame.""" + +from __future__ import annotations + +from typing import Any + +from tests.support.hbp_repeat_stack import build_hbp_repeat_stack + +from adn_server.application.routing.helpers import ( + RF_MODE_DUPLEX, + RF_MODE_SIMPLEX, + peer_options_fields, + peer_rf_mode, + peer_single_mode, +) +from adn_server.application.routing.peer_downlink_index import cached_peer_static_tgs + +_PEER = (1234567).to_bytes(4, "big") +_ADDR = ("10.0.0.9", 54321) + + +def _rptc(peer_id: bytes, *, slots: bytes, rx: bytes, tx: bytes) -> bytes: + """MMDVM RPTC frame (field widths per hblink.py / MMDVMHost RPTC layout).""" + return b"".join( + [ + b"RPTC", + peer_id, + b"CE5RPY ", + rx.ljust(9, b"0"), + tx.ljust(9, b"0"), + b"01", + b"01", + b"00.00000", + b"000.00000", + b"000", + b"Andorra".ljust(20), + b"hotspot".ljust(19), + slots, + b"".ljust(124), + b"MMDVM".ljust(40), + b"".ljust(40), + ] + ) + + +def _waiting_peer(stack: Any, peer_id: bytes, addr: tuple[str, int]) -> dict[str, Any]: + peer: dict[str, Any] = { + "CONNECTION": "WAITING_CONFIG", + "SOCKADDR": addr, + "RADIO_ID": str(int.from_bytes(peer_id, "big")), + "CALLSIGN": b"CE5RPY ", + } + stack.hbp._peers[peer_id] = peer + return peer + + +def test_options_fields_reused_while_the_blob_is_unchanged() -> None: + """Same OPTIONS must not be re-parsed: ingress asks several times per frame.""" + stack = build_hbp_repeat_stack() + stack.register_peer(_PEER, _ADDR) + peer = stack.hbp._peers[_PEER] + peer["OPTIONS"] = b"TS2_1=214;SINGLE=1;TIMER=15;" + + first = peer_options_fields(peer) + assert peer_options_fields(peer) is first + + +def test_rpto_makes_the_new_options_take_effect() -> None: + """RPTO drops the cached parse, so SINGLE/TIMER and static TGs follow the new blob.""" + stack = build_hbp_repeat_stack() + sys_cfg = stack.config["SYSTEMS"][stack.system_name] + stack.register_peer(_PEER, _ADDR) + peer = stack.hbp._peers[_PEER] + + stack.hbp._master_datagram_received(b"RPTO" + _PEER + b"TS2_1=214;SINGLE=1;", _ADDR) + assert peer_single_mode(peer, sys_cfg) is True + assert cached_peer_static_tgs(peer) == ((), ("214",)) + + stack.hbp._master_datagram_received(b"RPTO" + _PEER + b"TS1_1=91;SINGLE=0;", _ADDR) + assert peer_single_mode(peer, sys_cfg) is False + assert cached_peer_static_tgs(peer) == (("91",), ()) + assert peer_options_fields(peer).get("SINGLE") == "0" + + +def test_rptc_classifies_rf_mode_at_login() -> None: + """Simplex/duplex is decided from the RPTC fields, not on every frame.""" + stack = build_hbp_repeat_stack() + stack.config["SYSTEMS"][stack.system_name]["ALLOW_UNREG_ID"] = True + peer = _waiting_peer(stack, _PEER, _ADDR) + + stack.hbp._master_datagram_received( + _rptc(_PEER, slots=b"3", rx=b"438500000", tx=b"431100000"), _ADDR, + ) + assert peer["CONNECTION"] == "YES" + assert peer["RF_MODE"] == RF_MODE_DUPLEX + assert peer_rf_mode(peer) == RF_MODE_DUPLEX + + +def test_relogin_with_new_slots_reclassifies_rf_mode() -> None: + """A hotspot that comes back as simplex must not keep the duplex verdict.""" + stack = build_hbp_repeat_stack() + stack.config["SYSTEMS"][stack.system_name]["ALLOW_UNREG_ID"] = True + peer = _waiting_peer(stack, _PEER, _ADDR) + stack.hbp._master_datagram_received( + _rptc(_PEER, slots=b"3", rx=b"438500000", tx=b"431100000"), _ADDR, + ) + assert peer["RF_MODE"] == RF_MODE_DUPLEX + + peer["CONNECTION"] = "WAITING_CONFIG" + stack.hbp._master_datagram_received( + _rptc(_PEER, slots=b"4", rx=b"145500000", tx=b"145500000"), _ADDR, + ) + assert peer["RF_MODE"] == RF_MODE_SIMPLEX + assert peer_rf_mode(peer) == RF_MODE_SIMPLEX