Merge pull request #88 from Amateur-Digital-Network/fix/obp-count-refused-sources

fix(obp): make refused sources visible — count them and name each frame type
pull/89/head
ce5rpy 1 week ago committed by GitHub
commit 427ae67889
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194

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

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

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

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

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

Loading…
Cancel
Save

Powered by TurnKey Linux.