From 55faf797dae1102bb2e92fbca1b11cc6b7743fe3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 21 Sep 2026 01:35:08 -0300 Subject: [PATCH] fix(obp): count frames refused for coming from the wrong address MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The refusal log is capped to once per source address so a clone pinging every 10s cannot flood the log. That left an operator with no way to tell whether it was still happening: the first line scrolls away and nothing replaces it, while the fan-in still prints "RX b'BCKA' from -> OBP-USA", which reads as if the frame had been delivered. Refusing skips the engine, and the engine is what calls count_drop() for every other reason, so these frames were tallied nowhere either — not in the periodic "(ROUTER) system X refused frames: ..." line, not in the report. Count them as source-not-peer, so the tally answers "is the clone still knocking?" once a minute without repeating the explanation. --- src/adn_server/domain/mesh_engine.py | 6 +-- src/adn_server/domain/mesh_session.py | 13 ++--- .../infrastructure/config_normalizer.py | 3 +- .../twisted_adapters/udp_hbp.py | 52 +++++++++---------- tests/infrastructure/test_obp_dns_anchor.py | 35 +++++++++++-- 5 files changed, 64 insertions(+), 45 deletions(-) diff --git a/src/adn_server/domain/mesh_engine.py b/src/adn_server/domain/mesh_engine.py index a94ade2..76460cd 100644 --- a/src/adn_server/domain/mesh_engine.py +++ b/src/adn_server/domain/mesh_engine.py @@ -186,10 +186,8 @@ def accepts_source( ) -> bool: """A frame counts as ours when it comes from the peer. - RELAX_CHECKS widens that to any address, which is how a peer on a dynamic IP - keeps working. It does not widen it when DNS owns the peer: there the name is - the identity and only a re-resolution may move it, so a second host holding - the same passphrase is not mistaken for the peer. + RELAX_CHECKS widens that to any address, for a peer on a dynamic IP — but not + when DNS owns the peer, or a second host with the passphrase would pass as it. """ if addr == session.peer: return True diff --git a/src/adn_server/domain/mesh_session.py b/src/adn_server/domain/mesh_session.py index 425180f..5697606 100644 --- a/src/adn_server/domain/mesh_session.py +++ b/src/adn_server/domain/mesh_session.py @@ -61,12 +61,8 @@ def _peer_from_config(sys_cfg: dict[str, Any] | None) -> tuple[str | None, int]: def dns_host_from_config(sys_cfg: dict[str, Any] | None) -> str | None: - """``TARGET_IP`` as written, when it was a hostname rather than a literal address. - - ``normalize_obp_config`` keeps the original under ``_TARGET_IP`` before it - overwrites ``TARGET_IP`` with what the name resolved to, the same way PEER - systems keep ``_MASTER_IP``. - """ + """``TARGET_IP`` as written, when a hostname: ``normalize_obp_config`` keeps it + under ``_TARGET_IP`` before overwriting ``TARGET_IP`` with the resolved address.""" original = (sys_cfg or {}).get("_TARGET_IP") if not original: return None @@ -90,8 +86,7 @@ class ObpBridgeSession: quenched: dict[bytes, bytes] = field(default_factory=dict) stunned: bool = False drops: dict[str, int] = field(default_factory=dict) - # TARGET_IP as the operator wrote it, when that was a hostname. Set means DNS - # owns this peer's address: nothing the wire says can move it. + # Set means DNS owns this address: nothing the wire says can move it. dns_host: str | None = None resolved_peer: tuple[str, int] | None = None dns_checked_at: float = 0.0 @@ -117,7 +112,7 @@ class ObpBridgeSession: """Remember the address a datagram really came from. True when it moved.""" if not addr or not addr[0]: return False - if self.dns_anchored: # only a re-resolution may move a DNS-anchored peer + if self.dns_anchored: return False host, port = str(addr[0]), int(addr[1]) if self.peer == (host, port): diff --git a/src/adn_server/infrastructure/config_normalizer.py b/src/adn_server/infrastructure/config_normalizer.py index 4d93c02..aae97ed 100644 --- a/src/adn_server/infrastructure/config_normalizer.py +++ b/src/adn_server/infrastructure/config_normalizer.py @@ -151,8 +151,7 @@ def normalize_obp_config(config: dict) -> None: sys_cfg["NETWORK_ID"] = (net_id & 0xFFFFFFFF).to_bytes(4, "big") target_ip = str(sys_cfg.get("TARGET_IP", "")) target_port = int(sys_cfg.get("TARGET_PORT", 62044)) - # Keep the name the operator wrote, like PEER keeps _MASTER_IP: resolving - # here overwrites TARGET_IP, and a hostname has to stay re-resolvable. + # Resolving overwrites TARGET_IP, so keep the name, as PEER keeps _MASTER_IP. sys_cfg["_TARGET_IP"] = target_ip if target_ip: try: diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index b9149ab..a9e6f8e 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -155,11 +155,11 @@ logger = logging.getLogger(__name__) _DEFAULT_MESH_REGISTRY = MeshCodecRegistry() -# A DNS-anchored OBP peer is re-resolved on this cadence, and on demand when a -# frame arrives from somewhere else — never more often than the shorter one, so -# unknown traffic cannot drive a lookup per packet. 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 +# How often a refused source is named again, per address and frame type. +OBP_FOREIGN_LOG_INTERVAL_S = 60.0 def get_user_password(radio_id: int): @@ -272,7 +272,7 @@ class HBPProtocol(DatagramProtocol): self._bcsq_log_once: deque = deque(maxlen=1024) self._obp_target_sync_log_once: deque = deque(maxlen=1024) self._obp_foreign_bcka_log_once: deque = deque(maxlen=1024) - self._obp_foreign_source_log_once: deque = deque(maxlen=1024) + self._obp_foreign_source_seen: dict[tuple[Any, bytes], float] = {} else: self._laststrid = {1: b"", 2: b""} self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} @@ -2172,18 +2172,23 @@ class HBPProtocol(DatagramProtocol): _once.append(_stream_id) def _obp_reject_source(self, _opcode: bytes, _sockaddr: tuple[str, int]) -> None: - """Refuse a frame from an address DNS does not give for this peer, and ask again. + """Refuse a frame from an address this peer's name does not resolve to. - A peer really moving is exactly what this looks like, so the refusal - schedules a re-resolution: if the name now answers with this address, the - next frame is accepted. + A peer that really moved looks exactly like this, so the refusal asks DNS + again. Counting is explicit because refusing here never reaches the engine, + which is what tallies every other reason. """ + self._session.count_drop("source-not-peer") self._obp_resolve_target() - _once = getattr(self, "_obp_foreign_source_log_once", None) - if isinstance(_once, deque) and _sockaddr in _once: - return - if isinstance(_once, deque): - _once.append(_sockaddr) + _seen = getattr(self, "_obp_foreign_source_seen", None) + if isinstance(_seen, dict): + _key = (_sockaddr, _opcode) + _now = time.time() + if _now - _seen.get(_key, 0.0) < OBP_FOREIGN_LOG_INTERVAL_S: + return + if len(_seen) > 256: + _seen.clear() + _seen[_key] = _now _peer = self._session.peer _why = ( f"{self._session.dns_host} resolves to {_peer[0]}:{_peer[1]}" @@ -2200,11 +2205,7 @@ class HBPProtocol(DatagramProtocol): ) def _obp_resolve_target(self) -> None: - """Ask DNS where this peer is now: non-blocking, and rate limited. - - Never resolve on the datagram thread — an unknown source must not be able - to drive a lookup per packet. - """ + """Ask DNS where this peer is now, off the datagram path and rate limited.""" _session = self._session if not _session.dns_anchored: return @@ -2228,9 +2229,9 @@ class HBPProtocol(DatagramProtocol): _was[0], _was[1], ) - _once = getattr(self, "_obp_foreign_source_log_once", None) - if isinstance(_once, deque): - _once.clear() + _seen = getattr(self, "_obp_foreign_source_seen", None) + if isinstance(_seen, dict): + _seen.clear() def _obp_target_resolve_failed(self, _failure: Any) -> None: """Keep the address we have: a name server hiccup must not drop the link.""" @@ -2400,8 +2401,7 @@ class HBPProtocol(DatagramProtocol): 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: - # Control frames carry no NETWORK_ID, so the source address is the only - # thing telling one sender from another on a shared-passphrase mesh. + # No NETWORK_ID here, so the source address is all that tells senders apart. if not accepts_source(_sockaddr, policy=self._obp_policy(), session=self._session): self._obp_reject_source(_packet[:4], _sockaddr) return @@ -2411,10 +2411,8 @@ class HBPProtocol(DatagramProtocol): _session = self._session _now = time.time() _session.note_keepalive(_now) - # BCKA carries no NETWORK_ID, so with a shared passphrase anyone's - # keepalive verifies here: it may bootstrap a peer we have no address - # for, never move one we already have. DMRD/DMRE identify themselves - # and do that instead (_obp_sync_target_sock_from_peer). + # Anyone with the passphrase can send one, so it may bootstrap a peer + # we have no address for, never move one we have. DMRD/DMRE do that. if not _session.peer_known: if _session.learn_peer(_sockaddr, at=_now): logger.info( diff --git a/tests/infrastructure/test_obp_dns_anchor.py b/tests/infrastructure/test_obp_dns_anchor.py index 3ca8fb5..4a33e99 100644 --- a/tests/infrastructure/test_obp_dns_anchor.py +++ b/tests/infrastructure/test_obp_dns_anchor.py @@ -3,7 +3,6 @@ from __future__ import annotations import logging -from collections import deque from types import SimpleNamespace from twisted.internet import defer @@ -30,7 +29,7 @@ def _fake_obp(*, dns_host: str | None = _HOST) -> SimpleNamespace: _session=ObpBridgeSession( system_name="OBP-USA", configured_peer=_CONFIGURED, dns_host=dns_host ), - _obp_foreign_source_log_once=deque(maxlen=1024), + _obp_foreign_source_seen={}, ) fake._obp_resolve_target = lambda: _RESOLVE(fake) fake._obp_target_resolved = lambda host: _RESOLVED(fake, host) @@ -91,7 +90,9 @@ def test_a_frame_from_elsewhere_is_refused_and_asks_dns_again(monkeypatch, caplo assert "discarded" in caplog.text and _HOST in caplog.text -def test_a_refused_source_is_logged_once(monkeypatch, caplog) -> None: +def test_a_refused_source_is_rate_limited_but_counted_every_time(monkeypatch, caplog) -> None: + """A clone pinging every 10s must not flood the log, but every frame it sends + still has to reach the tally.""" resolver = _FakeResolver(_CONFIGURED[0]) monkeypatch.setattr(udp_hbp, "reactor", resolver) fake = _fake_obp() @@ -99,6 +100,34 @@ def test_a_refused_source_is_logged_once(monkeypatch, caplog) -> None: for _ in range(30): _REJECT(fake, b"BCKA", _ZOMBIE) assert sum("discarded" in line for line in caplog.text.splitlines()) == 1 + assert fake._session.drops == {"source-not-peer": 30} + + +def test_each_frame_type_from_a_refused_source_is_named(monkeypatch, caplog) -> None: + """Keying the cap on the address alone hid every frame type after the first: + a BCVE arriving a millisecond behind a BCKA was silently dropped.""" + resolver = _FakeResolver(_CONFIGURED[0]) + monkeypatch.setattr(udp_hbp, "reactor", resolver) + fake = _fake_obp() + with caplog.at_level(logging.INFO, logger="adn_server.infrastructure.twisted_adapters.udp_hbp"): + _REJECT(fake, b"BCKA", _ZOMBIE) + _REJECT(fake, b"BCVE", _ZOMBIE) + _REJECT(fake, b"DMRD", _ZOMBIE) + assert "BCKA from" in caplog.text + assert "BCVE from" in caplog.text + assert "DMRD from" in caplog.text + + +def test_a_refused_source_is_named_again_after_the_interval(monkeypatch, caplog) -> None: + resolver = _FakeResolver(_CONFIGURED[0]) + monkeypatch.setattr(udp_hbp, "reactor", resolver) + fake = _fake_obp() + with caplog.at_level(logging.INFO, logger="adn_server.infrastructure.twisted_adapters.udp_hbp"): + _REJECT(fake, b"BCKA", _ZOMBIE) + for key in fake._obp_foreign_source_seen: # as if the interval had elapsed + fake._obp_foreign_source_seen[key] -= udp_hbp.OBP_FOREIGN_LOG_INTERVAL_S + 1 + _REJECT(fake, b"BCKA", _ZOMBIE) + assert sum("discarded" in line for line in caplog.text.splitlines()) == 2 def test_a_peer_that_really_moved_is_adopted_on_the_next_resolution(monkeypatch) -> None: