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 <noreply@anthropic.com>
pull/99/head
yo 5 days ago
parent 595f4764d0
commit 1896a24842

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

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

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

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

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

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

@ -0,0 +1,157 @@
# ADN DMR Peer Server - tests routing forward plan cache
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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
Loading…
Cancel
Save

Powered by TurnKey Linux.