perf(obp): stop repeating per-frame work in OBP ingress

- Build BridgePolicy/AdmissionContext and PeerMeshConfig once, not per frame.
  Dropped on reload, and on the _SERVER_IDS swap an alias refresh does.
- Carry the DMRE timestamp on MeshIngress: the trailer was parsed twice.
- Let the engine be the only one vetting the source; it already answers None.
- Resolve the session and its peer once per datagram instead of 3-5 times.
- Dispatch effects most-frequent-first.

47.1us -> 28.6us per frame on the v5 path. Recorded effects corpus unchanged.
pull/93/head
Rodrigo Pérez 7 days ago
parent 27a0a1031c
commit d0fc7020f7

@ -190,10 +190,11 @@ def accepts_source(
bound, which NAT may rewrite and which needs not be the port we send to. With
no name to go by, RELAX_CHECKS decides as it always has.
"""
if addr == session.peer:
peer = session.peer
if addr == peer:
return True
if session.dns_anchored:
return bool(addr) and addr[0] == session.peer[0]
return bool(addr) and addr[0] == peer[0]
return bool(policy.relax_checks)

@ -47,6 +47,7 @@ class MeshIngress:
source_server: bytes
source_rptr: bytes
embedded_ver: int | None = None
timestamp: bytes = b"\x00" * 8
@dataclass(frozen=True, slots=True)

@ -103,8 +103,8 @@ class ObpBridgeSession:
"""
if not self.dns_host:
return None
host, port = self.resolved_peer or self.configured_peer
return (host, port) if host else None
anchor = self.resolved_peer or self.configured_peer
return anchor if anchor[0] else None
@property
def dns_anchored(self) -> bool:

@ -80,6 +80,7 @@ class DmreV5PeerTransport:
source_server=trailer.source_server,
source_rptr=trailer.source_rptr,
embedded_ver=trailer.embedded_version,
timestamp=trailer.timestamp,
)
def encode(self, egress: MeshEgress, config: PeerMeshConfig) -> bytes | None:

@ -155,6 +155,8 @@ logger = logging.getLogger(__name__)
_DEFAULT_MESH_REGISTRY = MeshCodecRegistry()
_NO_SERVER_IDS: frozenset[str] = frozenset()
OBP_DNS_REFRESH_S = 300.0
# Floor for on-demand lookups: unknown traffic must not drive one per packet.
OBP_DNS_MIN_INTERVAL_S = 15.0
@ -262,6 +264,8 @@ class HBPProtocol(DatagramProtocol):
self._mesh_registry = mesh_registry if mesh_registry is not None else _DEFAULT_MESH_REGISTRY
self._dynamic_tg_uc = dynamic_tg_uc
self._config = config.get("SYSTEMS", {}).get(system_name, {})
self._obp_policy_cache: tuple[Any, BridgePolicy] | None = None
self._peer_mesh_config_cache: PeerMeshConfig | None = None
if self._config.get("MODE") == "OPENBRIDGE":
self._laststrid = deque([], 20)
# Legacy parity: routerOBP.__init__ uses a flat dict keyed by stream_id
@ -349,6 +353,8 @@ class HBPProtocol(DatagramProtocol):
self._CONFIG = config
sys_cfg = config.get("SYSTEMS", {}).get(self._system, {})
self._config = sys_cfg
self._obp_policy_cache = None
self._peer_mesh_config_cache = None
if sys_cfg.get("MODE") == "MASTER":
self._peers = sys_cfg.setdefault("PEERS", {})
self._refresh_connected_peer_count()
@ -518,6 +524,14 @@ class HBPProtocol(DatagramProtocol):
del sub_map[rf_src]
def _peer_mesh_config(self) -> PeerMeshConfig:
cached = self._peer_mesh_config_cache
if cached is not None:
return cached
cached = self._build_peer_mesh_config()
self._peer_mesh_config_cache = cached
return cached
def _build_peer_mesh_config(self) -> PeerMeshConfig:
_global = self._CONFIG.get("GLOBAL", {})
_sid = _global.get("SERVER_ID", b"\x00\x00\x00\x00")
_server_id = (
@ -2144,7 +2158,12 @@ class HBPProtocol(DatagramProtocol):
else:
logger.debug("(%s) *BridgeControl* not sending BCVE, TARGET not currently known", self._system)
def _obp_sync_target_sock_from_peer(self, _sockaddr: tuple[str, int], _stream_id: bytes | None = None) -> None:
def _obp_sync_target_sock_from_peer(
self,
_sockaddr: tuple[str, int],
_stream_id: bytes | None = None,
_session: ObpBridgeSession | None = None,
) -> None:
"""Learn the address RELAX_CHECKS just accepted, so replies go back to it.
The configured peer stays where the operator put it; only the session
@ -2154,8 +2173,11 @@ class HBPProtocol(DatagramProtocol):
"""
if self._config.get("MODE") != "OPENBRIDGE" or not self._config.get("RELAX_CHECKS"):
return
_session = self._session
if _session is None:
_session = self._session
cur = _session.peer
if cur == _sockaddr:
return # already where we send: learn_peer would refuse it anyway
if not _session.learn_peer(_sockaddr, at=time.time()):
return
_once = getattr(self, "_obp_target_sync_log_once", None)
@ -2266,12 +2288,26 @@ class HBPProtocol(DatagramProtocol):
),
server_id=server_prefix(_global.get("SERVER_ID", 0)),
validate_server_ids=bool(_global.get("VALIDATE_SERVER_IDS")),
known_server_prefixes=self._CONFIG.get("_SERVER_IDS", set()),
known_server_prefixes=self._CONFIG.get("_SERVER_IDS", _NO_SERVER_IDS),
resolve_server_id=self.validate_obp_source_server_id,
)
def _obp_policy(self) -> BridgePolicy:
"""This bridge's rules, read once per frame and handed to the engine."""
"""This bridge's rules, handed to the engine.
Everything here comes from configuration, so it is built once and kept until
a reload swaps the dicts (``apply_system_config``) or an alias refresh swaps
``_SERVER_IDS`` under us, which is the one input that moves on its own.
"""
server_ids = self._CONFIG.get("_SERVER_IDS", _NO_SERVER_IDS)
cached = self._obp_policy_cache
if cached is not None and cached[0] is server_ids:
return cached[1]
policy = self._build_obp_policy()
self._obp_policy_cache = (server_ids, policy)
return policy
def _build_obp_policy(self) -> BridgePolicy:
return BridgePolicy(
system=self._system,
network_id=self._config.get("NETWORK_ID", b""),
@ -2282,26 +2318,14 @@ class HBPProtocol(DatagramProtocol):
)
def _obp_apply(self, effects: list[Effect] | None) -> None:
"""Carry out what the engine decided, in order."""
"""Carry out what the engine decided, in order.
Branches are mutually exclusive, so they run most-frequent-first.
"""
if not effects:
return
for effect in effects:
if isinstance(effect, Reject):
self._obp_reject(effect.rejection, effect.dst_id, effect.stream_id)
self._session.count_drop(effect.reason)
elif isinstance(effect, Log):
logger.log(effect.level, effect.message, *effect.args)
elif isinstance(effect, NoteStream):
self.note_dmrd_stream(effect.peer_id, effect.rf_src, effect.stream_id)
elif isinstance(effect, StoreTalkerAlias):
self.store_ta_from_voice_burst(
effect.peer_id,
effect.rf_src,
effect.stream_id,
effect.dtype_vseq,
effect.burst,
)
elif isinstance(effect, Deliver):
if isinstance(effect, Deliver):
if self._dmrd_received:
self._dmrd_received(
self._system,
@ -2322,6 +2346,21 @@ class HBPProtocol(DatagramProtocol):
obp_rssi=effect.rssi,
obp_source_rptr=effect.source_rptr,
)
elif isinstance(effect, NoteStream):
self.note_dmrd_stream(effect.peer_id, effect.rf_src, effect.stream_id)
elif isinstance(effect, StoreTalkerAlias):
self.store_ta_from_voice_burst(
effect.peer_id,
effect.rf_src,
effect.stream_id,
effect.dtype_vseq,
effect.burst,
)
elif isinstance(effect, Log):
logger.log(effect.level, effect.message, *effect.args)
elif isinstance(effect, Reject):
self._obp_reject(effect.rejection, effect.dst_id, effect.stream_id)
self._session.count_drop(effect.reason)
elif isinstance(effect, RequestVersion):
self._obp_send_bcve()
@ -2359,19 +2398,20 @@ class HBPProtocol(DatagramProtocol):
self._obp_apply(reject_v1_protocol(_stream_id, policy=_policy))
return
_ingress = self._try_decode_mesh_ingress(_packet)
_session = self._session
if _ingress is not None and _ingress.codec == "obp_v1" and not accepts_source(
_sockaddr, policy=_policy, session=self._session
_sockaddr, policy=_policy, session=_session
):
self._obp_reject_source(_packet[:4], _sockaddr)
return
if _ingress is not None and _ingress.codec == "obp_v1":
self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id)
self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id, _session)
self._obp_apply(
ingest_dmrd_v1(
_ingress,
_sockaddr,
policy=_policy,
session=self._session,
session=_session,
now=time.time(),
)
)
@ -2382,22 +2422,21 @@ class HBPProtocol(DatagramProtocol):
if _ingress is None or _ingress.codec != "dmre_v5":
return
_policy = self._obp_policy()
if not accepts_source(_sockaddr, policy=_policy, session=self._session):
_session = self._session
# None means the engine refused the source; it checks, so we do not.
_effects = ingest_dmre_v5(
_ingress,
_sockaddr,
policy=_policy,
session=_session,
timestamp_ns=int.from_bytes(_ingress.timestamp, "big"),
now=time.time(),
)
if _effects is None:
self._obp_reject_source(_packet[:4], _sockaddr)
return
_trailer = parse_dmre_trailer(_packet)
_timestamp = _trailer.timestamp if _trailer is not None else b"\x00" * 8
self._obp_sync_target_sock_from_peer(_sockaddr, _ingress.voice_frame[16:20])
self._obp_apply(
ingest_dmre_v5(
_ingress,
_sockaddr,
policy=_policy,
session=self._session,
timestamp_ns=int.from_bytes(_timestamp, "big"),
now=time.time(),
)
)
self._obp_sync_target_sock_from_peer(_sockaddr, _ingress.voice_frame[16:20], _session)
self._obp_apply(_effects)
elif _packet[:4] == EOBP:
logger.warning("(%s) *ProtoControl* KF7EEL EOBP protocol not supported", self._system)
elif self._config.get("ENHANCED_OBP") and _packet[:2] == BC:

@ -0,0 +1,90 @@
# ADN DMR Peer Server - tests infrastructure obp policy 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 bridge's policy is built from configuration and kept until that changes.
Holding it across datagrams is only safe while every input that can move is
known to invalidate it, so each of those is asserted here: a config reload
swaps the dicts, and an alias refresh swaps ``_SERVER_IDS`` on its own.
"""
from __future__ import annotations
import copy
from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol
from tests.harness.obp_ingress import PEER, build_config
def _protocol() -> HBPProtocol:
case = {"kind": "v5", "name": "policy", "desc": "policy", "validate_server_ids": True}
return HBPProtocol("OBP-FR", build_config(case), router=None, dmrd_received=lambda *a, **k: None)
def test_the_policy_is_built_once_while_nothing_moves() -> None:
protocol = _protocol()
assert protocol._obp_policy() is protocol._obp_policy()
def test_an_alias_refresh_that_swaps_server_ids_rebuilds_the_policy() -> None:
"""alias_loader replaces _SERVER_IDS in place on the live config, without a
reload. A policy still holding the old set would admit or refuse by stale data.
"""
protocol = _protocol()
before = protocol._obp_policy()
protocol._CONFIG["_SERVER_IDS"] = {"2084", "7302"}
after = protocol._obp_policy()
assert after is not before
assert after.admission.known_server_prefixes == {"2084", "7302"}
def test_a_config_reload_rebuilds_the_policy() -> None:
"""config_reload carries _SERVER_IDS across untouched, so the identity check
cannot notice a reload: apply_system_config has to drop the policy itself.
"""
protocol = _protocol()
before = protocol._obp_policy()
reloaded = copy.deepcopy(protocol._CONFIG)
reloaded["_SERVER_IDS"] = protocol._CONFIG.get("_SERVER_IDS")
reloaded["SYSTEMS"]["OBP-FR"]["RELAX_CHECKS"] = not before.relax_checks
protocol.apply_system_config(reloaded)
after = protocol._obp_policy()
assert after is not before
assert after.relax_checks is not before.relax_checks
def test_a_config_reload_rebuilds_the_peer_mesh_config() -> None:
protocol = _protocol()
before = protocol._peer_mesh_config()
assert protocol._peer_mesh_config() is before
reloaded = copy.deepcopy(protocol._CONFIG)
reloaded["SYSTEMS"]["OBP-FR"]["PASSPHRASE"] = b"rotated"
protocol.apply_system_config(reloaded)
assert protocol._peer_mesh_config().passphrase == b"rotated"
def test_egress_still_reaches_the_peer_after_a_reload() -> None:
"""The cached passphrase signs what leaves, so a stale one would sign wrongly."""
protocol = _protocol()
protocol._peer_mesh_config()
reloaded = copy.deepcopy(protocol._CONFIG)
reloaded["SYSTEMS"]["OBP-FR"]["TARGET_SOCK"] = PEER
protocol.apply_system_config(reloaded)
assert protocol._peer_mesh_config() is not None
Loading…
Cancel
Save

Powered by TurnKey Linux.