From 1896a2484205f5bcc18deeac14014fd303d6562d Mon Sep 17 00:00:00 2001 From: yo Date: Thu, 24 Sep 2026 00:11:24 +0200 Subject: [PATCH] perf(routing): stop rescanning the subscription store per voice frame Every group voice frame, from HBP or OpenBridge, asked store_has_table() whether its TG has a relay table. It answered by building a tuple of every subscription and calling table_key() on each, so the cost of a frame grew with the size of the network. On a live master (ADN 213, six OpenBridges) this was 26% of the CPU spent on OpenBridge ingress. - has_table() on the store answers from the table index it already keeps; the port gets a default that scans, for stores without one. - The forward plan of a group call (relay tables, resolved legs, the MASTER/PEER dedupe and the target entries) depends only on the store and on the source and target SYSTEMS blocks. It is now kept per (system, slot, TG) and rebuilt when the store's new revision counter moves or a reload replaces one of those blocks. ENABLED, quench, keepalive and contention are still checked per frame in the loop. - One OBP session lookup per target instead of two. - store_has_table is imported once, not on every frame. Per frame, one call bridged to 6 OpenBridges and a MASTER, with N other bridged TGs in the store: N=0 N=200 HBP ingress 81.8us -> 52.1us 447.4us -> 55.3us OBP ingress 95.9us -> 66.8us 469.3us -> 65.7us Co-Authored-By: Claude Opus 5.5 --- src/adn_server/application/ports.py | 4 + .../application/routing/obp_forward.py | 3 +- .../application/routing/voice_subscription.py | 104 ++++++++++++ .../application/routing_use_cases.py | 51 +----- .../subscription/subscription_queries.py | 2 +- .../infrastructure/subscription_store.py | 14 ++ tests/routing/test_forward_plan_cache.py | 157 ++++++++++++++++++ 7 files changed, 288 insertions(+), 47 deletions(-) create mode 100644 tests/routing/test_forward_plan_cache.py diff --git a/src/adn_server/application/ports.py b/src/adn_server/application/ports.py index f4e8d7c..9a97b58 100644 --- a/src/adn_server/application/ports.py +++ b/src/adn_server/application/ports.py @@ -298,6 +298,10 @@ class SubscriptionStore(ABC): def legs_in_table(self, table_key: str) -> tuple["Subscription", ...]: ... + def has_table(self, table_key: str) -> bool: + """True when at least one leg belongs to ``table_key``. Stores with an index override this.""" + return any(sub.table_key() == table_key for sub in self.snapshot()) + @abstractmethod def has_active_target_leg(self, system: str, slot: int, tgid: int) -> bool: ... diff --git a/src/adn_server/application/routing/obp_forward.py b/src/adn_server/application/routing/obp_forward.py index e959594..8b7c17d 100644 --- a/src/adn_server/application/routing/obp_forward.py +++ b/src/adn_server/application/routing/obp_forward.py @@ -52,6 +52,7 @@ from typing import Any from ...domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, int_id from ...domain.dmr import decode from ...domain.dmr.const import LC_OPT +from ..subscription.subscription_queries import store_has_table from .helpers import group_voice_tg_ingress_collision, obp_is_canonical_ingress, unit_data_reportable logger = logging.getLogger(__name__) @@ -456,8 +457,6 @@ class ObpForwardMixin: if self._config.get("GLOBAL", {}).get("GEN_STAT_BRIDGES"): _di = int_id(dst_id) _bk = str(_di) - from ..subscription.subscription_queries import store_has_table - if _di >= 5 and _di != 9 and not store_has_table(self._subscription_store, _bk): logger.debug("(%s) Bridge for STAT TG %s does not exist. Creating", system_name, _di) self.ensure_stat_relay(dst_id) diff --git a/src/adn_server/application/routing/voice_subscription.py b/src/adn_server/application/routing/voice_subscription.py index e762e72..7505b98 100644 --- a/src/adn_server/application/routing/voice_subscription.py +++ b/src/adn_server/application/routing/voice_subscription.py @@ -44,7 +44,9 @@ from __future__ import annotations import logging +from typing import Any, NamedTuple +from ...domain.value_objects import bytes_3 from ...domain.voice_routing import ForwardLeg, VoiceIngress from ..ports import SubscriptionStore from ..subscription.ingress import build_voice_ingress @@ -52,6 +54,27 @@ from ..subscription.router import SubscriptionRouter logger = logging.getLogger(__name__) +# Same set as ``build_voice_ingress``: other call types resolve to no legs. +_BRIDGE_CALL_TYPES = frozenset({"group", "vcsbk"}) +_FORWARD_PLAN_CACHE_MAX = 4096 + + +class ForwardPlan(NamedTuple): + """Where one group voice frame goes, ready for the forwarding loop.""" + + tables: tuple[str, ...] + legs: tuple[ForwardLeg, ...] + # (relay table key, target entry) per leg; shared between frames, read only. + entries: tuple[tuple[str, dict[str, Any]], ...] + + +class _CachedPlan(NamedTuple): + revision: int + # Every SYSTEMS block the plan read, as the object it read: a reload + # replaces the block, which is what makes the plan stale. + blocks: tuple[tuple[str, Any], ...] + plan: ForwardPlan + class VoiceSubscriptionMixin: """Wire ``SubscriptionRouter`` into ``dmrd_received``.""" @@ -90,6 +113,87 @@ class VoiceSubscriptionMixin: stream_id=stream_id, ) + def _group_voice_forward_plan( + self, + *, + system_name: str, + slot: int, + call_type: str, + source_is_obp: bool, + bridge_match_slot: int, + dst_int: int, + ) -> ForwardPlan: + """The forward plan of a group voice frame, rebuilt only when it can differ. + + It depends on the subscription store, the source and target SYSTEMS blocks + and on nothing in the frame beyond the key, so each frame of a call reuses + the plan of the first until the store is written or a block is reloaded. + Per-frame checks (ENABLED, quench, keepalive, contention) stay in the loop. + """ + systems = self._config.get("SYSTEMS", {}) + revision = getattr(self._subscription_store, "revision", None) + src_block = systems.get(system_name) + mode = "OPENBRIDGE" if source_is_obp else (src_block or {}).get("MODE", "") + routable = call_type in _BRIDGE_CALL_TYPES + # VoiceIngress.bridge_match_slot, which is what resolve() matches on. + match_slot = 1 if mode == "OPENBRIDGE" else (1 if int(slot) == 1 else 2) + key = (system_name, bridge_match_slot, match_slot, dst_int, mode == "OPENBRIDGE", routable) + cache: dict[tuple, _CachedPlan] | None = getattr(self, "_forward_plan_cache", None) + if cache is None: + cache = self._forward_plan_cache = {} + hit = cache.get(key) if revision is not None else None + if ( + hit is not None + and hit.revision == revision + and all(systems.get(name) is block for name, block in hit.blocks) + ): + return hit.plan + + tables, legs = self._voice_forward_plan( + system_name=system_name, + peer_id=b"", + rf_src=b"", + dst_id=dst_int.to_bytes(3, "big"), + slot=slot, + call_type=call_type, + stream_id=b"", + source_is_obp=source_is_obp, + bridge_match_slot=bridge_match_slot, + dst_int=dst_int, + ) + # One leg per (target, translated TGID) on MASTER/PEER targets: their + # send_peers() picks each peer's slot itself, so a second leg is the same + # audio twice. OpenBridge targets keep per-slot legs (separate links). + seen_hbp: set[tuple[str, int]] = set() + deduped = [] + for leg in legs: + if systems.get(leg.target_system, {}).get("MODE") != "OPENBRIDGE": + hbp_key = (leg.target_system, int(leg.target_tgid)) + if hbp_key in seen_hbp: + continue + seen_hbp.add(hbp_key) + deduped.append(leg) + table_key = tables[0] if tables else str(dst_int) + entries = tuple( + ( + table_key, + { + "SYSTEM": leg.target_system, + "TS": int(leg.slot), + "TGID": bytes_3(int(leg.target_tgid)), + "ACTIVE": True, + }, + ) + for leg in deduped + ) + plan = ForwardPlan(tables, tuple(deduped), entries) + if revision is not None: + names = {system_name, *(leg.target_system for leg in legs)} + if len(cache) >= _FORWARD_PLAN_CACHE_MAX: + cache.clear() + cache[key] = _CachedPlan(revision, tuple((n, systems.get(n)) for n in names), plan) + return plan + def _voice_relay_tables_with_active_source( self, system_name: str, diff --git a/src/adn_server/application/routing_use_cases.py b/src/adn_server/application/routing_use_cases.py index 79e4482..f5a50c3 100644 --- a/src/adn_server/application/routing_use_cases.py +++ b/src/adn_server/application/routing_use_cases.py @@ -64,7 +64,6 @@ from .routing.helpers import ( obp_publish_flat_bridge_tx, obp_status_plugin_voice, obp_sync_flat_bridge_tx_times, - obp_target_bcsq_quenches_stream, resolve_voice_peer_id, slot_has_active_voice, unit_data_hbp_target_idle, @@ -78,6 +77,7 @@ from .routing.subscription_table import SubscriptionTableMixin from .routing.timers import RoutingTimerMixin from .routing.voice_subscription import VoiceSubscriptionMixin from .server_voice import all_server_voice_ids +from .subscription.subscription_queries import store_has_table from .talker_alias_use_cases import TalkerAliasUseCases logger = logging.getLogger(__name__) @@ -259,8 +259,6 @@ class RoutingUseCases( # Arm ON in-band rules on VHEAD (echo 9990 and UA bridges); VTERM handled in udp_hbp too. if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VHEAD: self.apply_in_band_signalling(system_name, slot, dst_id, pkt_time) - from .subscription.subscription_queries import store_has_table - relay_table_key = str(int_id(dst_id)) dst_int = int_id(dst_id) # Legacy bridge_master to_target: OpenBridge clears TS bit — "all OpenBridge streams are @@ -494,40 +492,16 @@ class RoutingUseCases( source_lc = b"\x00\x00\x20" + dst_id_b + rf_src # Legacy bridge_master routerOBP: _sysIgnore accumulates across each to_target(BRIDGES[_bridge]) # pass; dedupe (SYSTEM, TS) for OpenBridge targets so the same leg is not sent twice per packet. - # SubscriptionRouter.resolve() already applies OBP dedup on OpenBridge targets. - forward_tables, forward_legs = self._voice_forward_plan( + # SubscriptionRouter.resolve() already applies OBP dedup on OpenBridge targets, and the plan + # collapses same-(target, translated TGID) MASTER/PEER legs (see _group_voice_forward_plan). + _plan = self._group_voice_forward_plan( system_name=system_name, - peer_id=peer_id, - rf_src=rf_src, - dst_id=dst_id, slot=slot, call_type=call_type, - stream_id=stream_id, source_is_obp=source_is_obp, bridge_match_slot=bridge_match_slot, dst_int=dst_int, ) - # A BRIDGES scan (legacy parity, kept in SubscriptionRouter.resolve()) can list the - # same MASTER/PEER target on both TS1 and TS2 (e.g. an inject-only proxy whose - # connected hotspots collectively use both slots for one TG). Legacy's dumb - # send_peers() broadcast made that harmless — each hotspot's own radio dropped the - # slot it didn't want. This server's send_peers() instead resolves each peer's - # actual listen slot from its OPTIONS (iter_downlink_voice_slots) regardless of the - # wire slot, so a second identical leg to the same target delivers the same audio - # twice — doubling the downlink rate and making unpaced bridges (e.g. ysf2dmr) sound - # slow. OpenBridge targets keep distinct per-slot legs (real separate links); collapse - # only same-(target, translated TGID) duplicates for MASTER/PEER targets. - if forward_legs: - _seen_hbp_leg: set[tuple[str, int]] = set() - _deduped_legs = [] - for _leg in forward_legs: - if systems_cfg.get(_leg.target_system, {}).get("MODE") != "OPENBRIDGE": - _hbp_key = (_leg.target_system, int(_leg.target_tgid)) - if _hbp_key in _seen_hbp_leg: - continue - _seen_hbp_leg.add(_hbp_key) - _deduped_legs.append(_leg) - forward_legs = tuple(_deduped_legs) _tx_report_peer = int_id(peer_id) if not source_is_obp and not synthetic_announcement: _tx_report_peer = int_id( @@ -543,18 +517,7 @@ class RoutingUseCases( _src_slot_st = src_proto.STATUS.get(slot, {}) if isinstance(_src_slot_st, dict) and _src_slot_st.get("_suppress_uplink"): _suppress_uplink = True - _leg_iter: list[tuple[str, dict[str, Any]]] = [ - ( - forward_tables[0] if forward_tables else str(dst_int), - { - "SYSTEM": leg.target_system, - "TS": int(leg.slot), - "TGID": bytes_3(int(leg.target_tgid)), - "ACTIVE": True, - }, - ) - for leg in forward_legs - ] + _leg_iter = _plan.entries for _relay_table_key, entry in _leg_iter: if _suppress_uplink: @@ -575,10 +538,10 @@ class RoutingUseCases( if isinstance(target_tgid, int): target_tgid = bytes_3(target_tgid) # If target has quenched us, don't send (~1856-1859). - if obp_target_bcsq_quenches_stream(self._config, entry["SYSTEM"], dst_id_b, stream_id): + _target_session = self._obp_session(entry["SYSTEM"]) + if _target_session.quenches(dst_id_b, stream_id): continue # If target has missed keepalives (ENHANCED_OBP), don't send (~1861-1863) - _target_session = self._obp_session(entry["SYSTEM"]) if _target_system.get("ENHANCED_OBP") and not _target_session.keepalive_ok(pkt_time): continue # Talkgroup ACL (global + per-system TG1) (~1865-1873) diff --git a/src/adn_server/application/subscription/subscription_queries.py b/src/adn_server/application/subscription/subscription_queries.py index 3d7c0b1..59e2e07 100644 --- a/src/adn_server/application/subscription/subscription_queries.py +++ b/src/adn_server/application/subscription/subscription_queries.py @@ -31,7 +31,7 @@ def store_has_table(store: SubscriptionStore, table_key: str) -> bool: Indexed: this runs per datagram, and the scan it replaces was building a tuple of every subscription to answer a yes/no question. """ - return bool(store.legs_in_table(table_key)) + return store.has_table(table_key) def system_has_active_leg_in_store( diff --git a/src/adn_server/infrastructure/subscription_store.py b/src/adn_server/infrastructure/subscription_store.py index 116696c..6a7a90c 100644 --- a/src/adn_server/infrastructure/subscription_store.py +++ b/src/adn_server/infrastructure/subscription_store.py @@ -48,6 +48,12 @@ class InMemorySubscriptionStore(SubscriptionStore): # What each leg was indexed under. Callers change a leg in place and then # upsert it, so by the time it is unindexed it may no longer say where it is. self._indexed: dict[SubscriptionId, tuple[str, _IndexKey | None]] = {} + self._revision = 0 + + @property + def revision(self) -> int: + """Bumped on every write, so readers can keep what they derived until it moves.""" + return self._revision def get(self, sub_id: SubscriptionId) -> Subscription | None: return self._items.get(sub_id) @@ -58,12 +64,14 @@ class InMemorySubscriptionStore(SubscriptionStore): self._unindex(old) self._items[subscription.subscription_id] = subscription self._index(subscription) + self._revision += 1 def remove(self, sub_id: SubscriptionId) -> bool: old = self._items.pop(sub_id, None) if old is None: return False self._unindex(old) + self._revision += 1 return True def clear(self) -> None: @@ -72,12 +80,14 @@ class InMemorySubscriptionStore(SubscriptionStore): self._source_tables.clear() self._active_target_counts.clear() self._indexed.clear() + self._revision += 1 def replace_all(self, subscriptions: Sequence[Subscription]) -> None: self.clear() for sub in subscriptions: self._items[sub.subscription_id] = sub self._index(sub) + self._revision += 1 def snapshot(self) -> tuple[Subscription, ...]: return tuple(self._items.values()) @@ -106,6 +116,10 @@ class InMemorySubscriptionStore(SubscriptionStore): return () return tuple(sorted(keys)) + def has_table(self, table_key: str) -> bool: + """O(1): the table index instead of a scan of every subscription.""" + return bool(self._by_table.get(table_key)) + def legs_in_table(self, table_key: str) -> tuple[Subscription, ...]: """All legs for a relay table key (indexed).""" return tuple(self._by_table.get(table_key, ())) diff --git a/tests/routing/test_forward_plan_cache.py b/tests/routing/test_forward_plan_cache.py new file mode 100644 index 0000000..12199ca --- /dev/null +++ b/tests/routing/test_forward_plan_cache.py @@ -0,0 +1,157 @@ +# ADN DMR Peer Server - tests routing forward plan cache +# +# 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 +############################################################################### + +"""A group voice call reuses its forward plan, and the plan follows the store and the config.""" + +from __future__ import annotations + +import copy + +from tests.harness.deterministic import ( + DeterministicScenario, + PacketSpec, + active_routing_table, + add_openbridge_system, + minimal_config, +) + +from adn_server.application.subscription.routing_table_import import subscriptions_from_routing_table +from adn_server.domain.subscription import SubscriptionPhase +from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore + +TG = 213 +OBPS = ("OBP-0", "OBP-1") + + +def _scenario() -> DeterministicScenario: + config = minimal_config(("MASTER-A", "MASTER-B")) + for i, name in enumerate(OBPS): + add_openbridge_system(config, name) + config["SYSTEMS"][name]["NETWORK_ID"] = (100 + i).to_bytes(4, "big") + table = active_routing_table( + TG, (("MASTER-A", 2), ("MASTER-B", 2)) + tuple((n, 1) for n in OBPS), timeout_minutes=10**6 + ) + sc = DeterministicScenario(config=config, routing_table=table) + sc.routing.apply_startup_subscriptions() + return sc + + +def _call(sc: DeterministicScenario, stream_id: int, bursts: int = 3) -> None: + base = PacketSpec(dst_id=TG, stream_id=stream_id, slot=2) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_head_spec(base)) + for seq in range(1, bursts + 1): + sc.clock.advance(0.06) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=seq % 6)) + + +def _plan(sc: DeterministicScenario): + return sc.routing._group_voice_forward_plan( + system_name="MASTER-A", slot=2, call_type="group", source_is_obp=False, + bridge_match_slot=2, dst_int=TG, + ) + + +def _targets(sc: DeterministicScenario) -> set[str]: + return {p.target_system for p in sc.capture.packets} + + +def test_frames_of_a_call_share_one_plan() -> None: + sc = _scenario() + first = _plan(sc) + assert _plan(sc) is first + assert {entry["SYSTEM"] for _, entry in first.entries} == {"MASTER-B", *OBPS} + + +def test_plan_is_rebuilt_after_a_store_write() -> None: + sc = _scenario() + _call(sc, 0x1001) + assert _targets(sc) == {"MASTER-B", *OBPS} + + store = sc.subscription_store + before = _plan(sc) + leg = next(s for s in store.legs_in_table(str(TG)) if s.system.value == "OBP-1") + leg.state.phase = SubscriptionPhase.IDLE + store.upsert(leg) + + after = _plan(sc) + assert after is not before + assert "OBP-1" not in {entry["SYSTEM"] for _, entry in after.entries} + sc.capture.packets.clear() + sc.clock.advance(0.06) + base = PacketSpec(dst_id=TG, stream_id=0x1001, slot=2) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_burst_spec(base, seq=9, dtype_vseq=3)) + assert _targets(sc) == {"MASTER-B", "OBP-0"} + + +def test_plan_is_rebuilt_when_a_reload_replaces_a_system_block() -> None: + sc = _scenario() + before = _plan(sc) + # config_reload assigns a merged copy per system; the old block is gone. + sc.config["SYSTEMS"]["OBP-0"] = copy.deepcopy(sc.config["SYSTEMS"]["OBP-0"]) + assert _plan(sc) is not before + assert _plan(sc) is _plan(sc) + + +def test_cached_plan_matches_a_fresh_one() -> None: + sc = _scenario() + cached = _plan(sc) + sc.routing._forward_plan_cache.clear() + fresh = _plan(sc) + assert fresh is not cached + assert fresh == cached + + +def test_has_table_agrees_with_a_full_scan_through_every_write() -> None: + table = active_routing_table(TG, (("MASTER-A", 2), ("OBP-0", 1))) + table |= active_routing_table(91, (("MASTER-A", 1),)) + subs = subscriptions_from_routing_table(table) + store = InMemorySubscriptionStore() + + def agrees() -> None: + for key in (str(TG), "91", "4000"): + scan = any(s.table_key() == key for s in store.snapshot()) + assert store.has_table(key) is scan, key + + revisions = [store.revision] + store.replace_all(subs) + agrees() + revisions.append(store.revision) + store.remove(subs[0].subscription_id) + agrees() + revisions.append(store.revision) + store.upsert(subs[0]) + agrees() + revisions.append(store.revision) + for sub in store.legs_in_table("91"): + store.remove(sub.subscription_id) + agrees() + revisions.append(store.revision) + store.clear() + agrees() + revisions.append(store.revision) + assert revisions == sorted(set(revisions)), "every write moves the revision" + + +def test_removing_a_missing_subscription_keeps_the_revision() -> None: + store = InMemorySubscriptionStore() + subs = subscriptions_from_routing_table(active_routing_table(TG, (("MASTER-A", 2),))) + before = store.revision + assert store.remove(subs[0].subscription_id) is False + assert store.revision == before