fix: log OBP RELAX_CHECKS target sync at debug, once per stream

pull/75/head
Rodrigo Pérez 1 week ago
parent 834034ea16
commit 783f5d196d

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

@ -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)
Loading…
Cancel
Save

Powered by TurnKey Linux.