From 783f5d196d635d753a78b558ffbcfa73483432ff Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Fri, 18 Sep 2026 18:21:08 -0300 Subject: [PATCH] fix: log OBP RELAX_CHECKS target sync at debug, once per stream --- .../twisted_adapters/udp_hbp.py | 34 +++++++---- .../test_obp_relax_checks_sync_log.py | 57 +++++++++++++++++++ 2 files changed, 79 insertions(+), 12 deletions(-) create mode 100644 tests/infrastructure/test_obp_relax_checks_sync_log.py diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 1675633..b0f31cc 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -247,6 +247,7 @@ class HBPProtocol(DatagramProtocol): # cleaned uniformly. No pre-seed of slot keys here. self.STATUS: dict[Any, Any] = {} self._bcsq_log_once: deque = deque(maxlen=1024) + self._obp_target_sync_log_once: deque = deque(maxlen=1024) else: self._laststrid = {1: b"", 2: b""} self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} @@ -2109,9 +2110,14 @@ 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]) -> None: + def _obp_sync_target_sock_from_peer(self, _sockaddr: tuple[str, int], _stream_id: bytes | None = None) -> None: """If RELAX_CHECKS accepted traffic from a different IP:port than TARGET_SOCK, sync (same idea as BCKA). - Ensures BCSQ and outbound DMR go to the peer address we actually receive from.""" + Ensures BCSQ and outbound DMR go to the peer address we actually receive from. + + A peer behind per-packet load-balanced NAT can flip source address on every frame of the + same call, so the sync itself still runs every packet but the log is debug and capped to + once per stream_id (same _log_once deque idiom as _bcsq_log_once above). + """ if self._config.get("MODE") != "OPENBRIDGE" or not self._config.get("RELAX_CHECKS"): return if not _sockaddr or not _sockaddr[0]: @@ -2120,14 +2126,18 @@ class HBPProtocol(DatagramProtocol): if cur == _sockaddr: return h, p = _sockaddr[0], int(_sockaddr[1]) - logger.info( - "(%s) *BridgeControl* OBP peer address sync to %s:%s (RELAX_CHECKS; was %s:%s)", - self._system, - h, - p, - (cur[0] if cur and cur[0] else "?"), - (cur[1] if cur and len(cur) > 1 else "?"), - ) + _once = getattr(self, "_obp_target_sync_log_once", None) + if _stream_id is None or not isinstance(_once, deque) or _stream_id not in _once: + logger.debug( + "(%s) *BridgeControl* OBP peer address sync to %s:%s (RELAX_CHECKS; was %s:%s)", + self._system, + h, + p, + (cur[0] if cur and cur[0] else "?"), + (cur[1] if cur and len(cur) > 1 else "?"), + ) + if _stream_id is not None and isinstance(_once, deque): + _once.append(_stream_id) self._config["TARGET_IP"] = h self._config["TARGET_PORT"] = p self._config["TARGET_SOCK"] = (h, p) @@ -2166,7 +2176,7 @@ class HBPProtocol(DatagramProtocol): _ingress = self._try_decode_mesh_ingress(_packet) if _ingress is not None and _ingress.codec == "obp_v1" and (_sockaddr == self._config.get("TARGET_SOCK") or self._config.get("RELAX_CHECKS")): _data = _ingress.voice_frame - self._obp_sync_target_sock_from_peer(_sockaddr) + self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id) _peer_id = _data[11:15] if self._config.get("NETWORK_ID") != _peer_id: if _stream_id not in self._laststrid: @@ -2294,7 +2304,7 @@ class HBPProtocol(DatagramProtocol): _trailer = parse_dmre_trailer(_packet) _timestamp = _trailer.timestamp if _trailer is not None else b"\x00" * 8 _stream_id = _data[16:20] - self._obp_sync_target_sock_from_peer(_sockaddr) + self._obp_sync_target_sock_from_peer(_sockaddr, _stream_id) _peer_id = _data[11:15] if self._config.get("NETWORK_ID") != _peer_id: if _stream_id not in self._laststrid: diff --git a/tests/infrastructure/test_obp_relax_checks_sync_log.py b/tests/infrastructure/test_obp_relax_checks_sync_log.py new file mode 100644 index 0000000..fc36a0e --- /dev/null +++ b/tests/infrastructure/test_obp_relax_checks_sync_log.py @@ -0,0 +1,57 @@ +# ADN DMR Peer Server - OBP RELAX_CHECKS target-address sync log is debug, once per stream + +from __future__ import annotations + +import logging +from collections import deque +from types import SimpleNamespace + +from adn_server.infrastructure.twisted_adapters.udp_hbp import HBPProtocol + +_SYNC = HBPProtocol._obp_sync_target_sock_from_peer + + +def _fake_obp(target_sock=None) -> SimpleNamespace: + return SimpleNamespace( + _system="OBP-USA", + _config={"MODE": "OPENBRIDGE", "RELAX_CHECKS": True, "TARGET_SOCK": target_sock}, + _obp_target_sync_log_once=deque(maxlen=1024), + ) + + +def test_sync_updates_target_sock_every_packet() -> None: + fake = _fake_obp(target_sock=("1.1.1.1", 62044)) + _SYNC(fake, ("2.2.2.2", 62044), b"strm") + assert fake._config["TARGET_SOCK"] == ("2.2.2.2", 62044) + _SYNC(fake, ("1.1.1.1", 62044), b"strm") + assert fake._config["TARGET_SOCK"] == ("1.1.1.1", 62044) + + +def test_sync_logs_debug_once_per_stream_even_when_flapping(caplog) -> None: + fake = _fake_obp(target_sock=("1.1.1.1", 62044)) + with caplog.at_level(logging.DEBUG, logger="adn_server.infrastructure.twisted_adapters.udp_hbp"): + for _ in range(20): + _SYNC(fake, ("2.2.2.2", 62044), b"strm-1") + _SYNC(fake, ("1.1.1.1", 62044), b"strm-1") + sync_records = [r for r in caplog.records if "OBP peer address sync" in r.message] + assert len(sync_records) == 1 + assert sync_records[0].levelno == logging.DEBUG + # still tracked the address on every flap despite logging once + assert fake._config["TARGET_SOCK"] == ("1.1.1.1", 62044) + + +def test_sync_logs_again_for_a_new_stream(caplog) -> None: + fake = _fake_obp(target_sock=("1.1.1.1", 62044)) + with caplog.at_level(logging.DEBUG, logger="adn_server.infrastructure.twisted_adapters.udp_hbp"): + _SYNC(fake, ("2.2.2.2", 62044), b"strm-1") + _SYNC(fake, ("1.1.1.1", 62044), b"strm-1") + _SYNC(fake, ("2.2.2.2", 62044), b"strm-2") + sync_records = [r for r in caplog.records if "OBP peer address sync" in r.message] + assert len(sync_records) == 2 + + +def test_no_relax_checks_no_sync() -> None: + fake = _fake_obp(target_sock=("1.1.1.1", 62044)) + fake._config["RELAX_CHECKS"] = False + _SYNC(fake, ("2.2.2.2", 62044), b"strm") + assert fake._config["TARGET_SOCK"] == ("1.1.1.1", 62044)