diff --git a/src/adn_server/domain/mesh_engine.py b/src/adn_server/domain/mesh_engine.py index ee798ad..c3a0d79 100644 --- a/src/adn_server/domain/mesh_engine.py +++ b/src/adn_server/domain/mesh_engine.py @@ -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) diff --git a/src/adn_server/domain/mesh_routing.py b/src/adn_server/domain/mesh_routing.py index b6b3b0a..72eb06a 100644 --- a/src/adn_server/domain/mesh_routing.py +++ b/src/adn_server/domain/mesh_routing.py @@ -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) diff --git a/src/adn_server/domain/mesh_session.py b/src/adn_server/domain/mesh_session.py index b6cca17..47562b1 100644 --- a/src/adn_server/domain/mesh_session.py +++ b/src/adn_server/domain/mesh_session.py @@ -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: diff --git a/src/adn_server/infrastructure/mesh/transports.py b/src/adn_server/infrastructure/mesh/transports.py index 5729b9d..9d98883 100644 --- a/src/adn_server/infrastructure/mesh/transports.py +++ b/src/adn_server/infrastructure/mesh/transports.py @@ -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: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index bc1380b..0e65b0a 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -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: diff --git a/tests/infrastructure/test_obp_policy_cache.py b/tests/infrastructure/test_obp_policy_cache.py new file mode 100644 index 0000000..1180569 --- /dev/null +++ b/tests/infrastructure/test_obp_policy_cache.py @@ -0,0 +1,90 @@ +# ADN DMR Peer Server - tests infrastructure obp policy 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 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