From 35f5c01b3f9f516c4bbb4b2db2fa2e2cf41fee37 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Rodrigo=20P=C3=A9rez?= Date: Mon, 21 Sep 2026 01:11:26 -0300 Subject: [PATCH] fix(obp): keep keepalives from relocating a bridge's peer MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit BCKA carries no NETWORK_ID, so on a shared-passphrase mesh anyone's keepalive verifies against any bridge. It was moving session.peer with no gate at all — not even RELAX_CHECKS, which the recorded corpus shows: "bcka from 9.9.9.9:62201 relax=False" moved egress to 9.9.9.9. Since session.peer is where voice and control are sent, a second instance of a peer (seen in production on OBP-USA: two hosts, two 10s keepalive timers) took the traffic over every few seconds. A keepalive now only confirms liveness. It may still bootstrap a bridge that has no address yet (inbound-only, no TARGET_IP), since there is nothing to steal there and it is the only way such a bridge learns where to answer. Relocation is left to DMRD/DMRE, which identify themselves. Also: the fan-in demux ranked bridges by sys_cfg's TARGET_SOCK, frozen since #81 moved runtime state into the session, so the "live" ranks were dead code and a peer that really moved was no longer recognised. It now reads learned_peer from the session store, closing the integration #81 left pending. --- .../infrastructure/proxy/obp_fanin.py | 43 ++++++++++--------- .../infrastructure/proxy/obp_runtime.py | 4 +- .../twisted_adapters/udp_hbp.py | 30 +++++++++++-- tests/fixtures/obp_ingress_effects.jsonl | 9 ++-- tests/harness/obp_ingress.py | 4 ++ .../test_obp_ingress_effects.py | 36 +++++++++++++++- tests/infrastructure/test_obp_proxy.py | 36 ++++++++++++---- 7 files changed, 124 insertions(+), 38 deletions(-) diff --git a/src/adn_server/infrastructure/proxy/obp_fanin.py b/src/adn_server/infrastructure/proxy/obp_fanin.py index ca71e95..b2ab53a 100644 --- a/src/adn_server/infrastructure/proxy/obp_fanin.py +++ b/src/adn_server/infrastructure/proxy/obp_fanin.py @@ -28,6 +28,7 @@ from typing import Any, Protocol from twisted.internet.protocol import DatagramProtocol +from adn_server.domain.mesh_session import MeshSessionStore from adn_server.infrastructure.hbp_constants import BCKA, BCSQ, BCST, BCVE, DMRD, DMRE from adn_server.infrastructure.mesh.obp_v1 import verify_bcka, verify_bcsq, verify_bcst, verify_bcve from adn_server.infrastructure.udp_rcvbuf import apply_udp_rcvbuf, udp_rcvbuf_bytes @@ -106,10 +107,7 @@ class ObpBridgeEntry: sink: InProcessObpSink reply_transport: ObpIngressReplyTransport legacy_port: int | None = None - # Live SYSTEMS. dict; used to tell bridges apart by peer address. - sys_cfg: dict[str, Any] | None = None - # Peer as configured, snapshotted when the registry is built. RELAX_CHECKS - # rewrites TARGET_SOCK in sys_cfg at runtime; this one never moves. + # Peer as configured, snapshotted when the registry is built. peer_hint: tuple[str, int] | None = None @@ -120,6 +118,8 @@ class ObpBridgeRegistry: by_network_id: dict[bytes, str] = field(default_factory=dict) by_legacy_port: dict[int, str] = field(default_factory=dict) bridges: dict[str, ObpBridgeEntry] = field(default_factory=dict) + # Where each bridge's peer actually is, as DMRD/DMRE taught it. + sessions: MeshSessionStore | None = None def register(self, entry: ObpBridgeEntry) -> None: self.bridges[entry.system_name] = entry @@ -127,35 +127,38 @@ class ObpBridgeRegistry: if entry.legacy_port is not None: self.by_legacy_port[entry.legacy_port] = entry.system_name - @staticmethod - def _peer_rank(entry: ObpBridgeEntry, addr: tuple[str, int]) -> int: + def _learned_peer(self, system_name: str) -> tuple[str, int] | None: + if self.sessions is None: + return None + session = self.sessions.get(system_name) + return session.learned_peer if session is not None else None + + def _peer_rank(self, entry: ObpBridgeEntry, addr: tuple[str, int]) -> int: host, port = addr hint = entry.peer_hint - live = peer_sock_from_config(entry.sys_cfg) + learned = self._learned_peer(entry.system_name) if hint is not None and hint == (host, port): return 0 - if live is not None and live == (host, port): + if learned is not None and learned == (host, port): return 1 if hint is not None and hint[0] == host: return 2 - if live is not None and live[0] == host: + if learned is not None and learned[0] == host: return 3 return 4 def bridges_by_peer(self, addr: tuple[str, int] | None) -> list[tuple[str, ObpBridgeEntry]]: """Bridges ordered by how well their peer matches ``addr``. - The configured peer (``peer_hint``, taken from the YAML when the registry is - built) is tried before the live ``TARGET_SOCK``, which RELAX_CHECKS rewrites - in place whenever a bridge accepts traffic from an unexpected source: ranking - the live value first would let a single misattributed control frame move a - bridge's target and then keep matching that same wrong address. Matching the - live value after it still lets a peer that legitimately moved (dynamic IP, - learned from a DMRD, which carries NETWORK_ID) be recognised. - - Order: configured IP:port, live IP:port, configured IP, live IP (a peer that - answers from another local port, e.g. its own fan-in), then registration - order. + The configured peer (``peer_hint``, from the YAML) is tried before the learned + one so a wrong match can never entrench itself; the learned one is what lets a + peer that really moved still be recognised. It comes from the session store, + where only DMRD/DMRE put it — they carry NETWORK_ID, so they cannot be + misattributed the way a control frame can. + + Order: configured IP:port, learned IP:port, configured IP, learned IP (a peer + answering from another local port, e.g. its own fan-in), then registration + order — where a shared passphrase leaves nothing else to go on. """ items = list(self.bridges.items()) if addr is None: diff --git a/src/adn_server/infrastructure/proxy/obp_runtime.py b/src/adn_server/infrastructure/proxy/obp_runtime.py index 66c849e..b784160 100644 --- a/src/adn_server/infrastructure/proxy/obp_runtime.py +++ b/src/adn_server/infrastructure/proxy/obp_runtime.py @@ -29,6 +29,7 @@ from typing import Any from twisted.internet import reactor from adn_server.application.proxy.deployment import obp_bridge_legacy_listen_port, obp_proxy_enabled +from adn_server.domain.mesh_session import mesh_sessions from adn_server.infrastructure.proxy.obp_config import obp_proxy_settings from adn_server.infrastructure.proxy.obp_fanin import ( InProcessObpSink, @@ -50,7 +51,7 @@ def build_obp_bridge_registry( primary_transport: Any, ) -> ObpBridgeRegistry: """Register enabled OPENBRIDGE systems for fan-in demux.""" - registry = ObpBridgeRegistry() + registry = ObpBridgeRegistry(sessions=mesh_sessions(config)) systems = config.get("SYSTEMS", {}) if not isinstance(systems, dict): return registry @@ -83,7 +84,6 @@ def build_obp_bridge_registry( sink=InProcessObpSink(proto), reply_transport=reply, legacy_port=legacy_port if legacy_port and legacy_port > 0 else None, - sys_cfg=sys_cfg, peer_hint=peer_sock_from_config(sys_cfg), ) registry.register(entry) diff --git a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py index 6605eb2..749a7d6 100644 --- a/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py +++ b/src/adn_server/infrastructure/twisted_adapters/udp_hbp.py @@ -265,6 +265,7 @@ class HBPProtocol(DatagramProtocol): self.STATUS: dict[Any, Any] = {} 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) else: self._laststrid = {1: b"", 2: b""} self.STATUS = {1: _make_slot_status(), 2: _make_slot_status()} @@ -2320,11 +2321,34 @@ class HBPProtocol(DatagramProtocol): if _packet[:4] == BCKA and len(_packet) >= 24: if verify_bcka(_packet, _passphrase): _session = self._session - _was = _session.peer _now = time.time() _session.note_keepalive(_now) - if _session.learn_peer(_sockaddr, at=_now): - logger.info("(%s) *BridgeControl* Source IP and Port has changed for OBP from %s:%s to %s:%s, updating", self._system, _was[0], _was[1], _sockaddr[0], _sockaddr[1]) + # 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). + if not _session.peer_known: + if _session.learn_peer(_sockaddr, at=_now): + logger.info( + "(%s) *BridgeControl* OBP peer address learned from keepalive: %s:%s", + self._system, + _sockaddr[0], + _sockaddr[1], + ) + elif _sockaddr != _session.peer: + _once = getattr(self, "_obp_foreign_bcka_log_once", None) + if not isinstance(_once, deque) or _sockaddr not in _once: + if isinstance(_once, deque): + _once.append(_sockaddr) + logger.debug( + "(%s) *BridgeControl* BCKA from %s:%s is not this bridge's peer %s:%s " + "(keepalive only; a second instance of the peer, or another bridge sharing the passphrase)", + self._system, + _sockaddr[0], + _sockaddr[1], + _session.peer[0], + _session.peer[1], + ) else: logger.info("(%s) *BridgeControl* BCKA invalid KeepAlive, packet discarded", self._system) # Source quench — legacy hblink.py OPENBRIDGE ~629-639 (sets CONFIG['_bcsq'][tgid]=stream_id) diff --git a/tests/fixtures/obp_ingress_effects.jsonl b/tests/fixtures/obp_ingress_effects.jsonl index 700aa99..7ee6524 100644 --- a/tests/fixtures/obp_ingress_effects.jsonl +++ b/tests/fixtures/obp_ingress_effects.jsonl @@ -78,10 +78,10 @@ {"case":{"desc":"from an unexpected address","from":["9.9.9.9",40000],"kind":"v5","name":"v5 from an unexpected address","stream":1358954574},"effects":{"delivered":[[2130003,214]],"egress":[[24,["9.9.9.9",40000]],[89,["9.9.9.9",40000]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* OBP peer address sync to 9.9.9.9:40000 (RELAX_CHECKS; was 82.65.127.86:62201)"]],"quenched":[]}} {"case":{"desc":"from 82.65.127.86:62201 relax=False","from":["82.65.127.86",62201],"kind":"bcka","name":"bcka from 82.65.127.86:62201 relax=False","relax":false,"stream":1358954575},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"]],"quenched":[]}} {"case":{"desc":"from 82.65.127.86:62201 relax=True","from":["82.65.127.86",62201],"kind":"bcka","name":"bcka from 82.65.127.86:62201 relax=True","relax":true,"stream":1358954576},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"]],"quenched":[]}} -{"case":{"desc":"from 82.65.127.86:40000 relax=False","from":["82.65.127.86",40000],"kind":"bcka","name":"bcka from 82.65.127.86:40000 relax=False","relax":false,"stream":1358954577},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",40000]],[73,["82.65.127.86",40000]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* Source IP and Port has changed for OBP from 82.65.127.86:62201 to 82.65.127.86:40000, updating"]],"quenched":[]}} -{"case":{"desc":"from 82.65.127.86:40000 relax=True","from":["82.65.127.86",40000],"kind":"bcka","name":"bcka from 82.65.127.86:40000 relax=True","relax":true,"stream":1358954578},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",40000]],[73,["82.65.127.86",40000]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* Source IP and Port has changed for OBP from 82.65.127.86:62201 to 82.65.127.86:40000, updating"]],"quenched":[]}} -{"case":{"desc":"from 9.9.9.9:62201 relax=False","from":["9.9.9.9",62201],"kind":"bcka","name":"bcka from 9.9.9.9:62201 relax=False","relax":false,"stream":1358954579},"effects":{"delivered":[],"egress":[[24,["9.9.9.9",62201]],[73,["9.9.9.9",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* Source IP and Port has changed for OBP from 82.65.127.86:62201 to 9.9.9.9:62201, updating"]],"quenched":[]}} -{"case":{"desc":"from 9.9.9.9:62201 relax=True","from":["9.9.9.9",62201],"kind":"bcka","name":"bcka from 9.9.9.9:62201 relax=True","relax":true,"stream":1358954580},"effects":{"delivered":[],"egress":[[24,["9.9.9.9",62201]],[73,["9.9.9.9",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* Source IP and Port has changed for OBP from 82.65.127.86:62201 to 9.9.9.9:62201, updating"]],"quenched":[]}} +{"case":{"desc":"from 82.65.127.86:40000 relax=False","from":["82.65.127.86",40000],"kind":"bcka","name":"bcka from 82.65.127.86:40000 relax=False","relax":false,"stream":1358954577},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* BCKA from 82.65.127.86:40000 is not this bridge's peer 82.65.127.86:62201 (keepalive only; a second instance of the peer, or another bridge sharing the passphrase)"]],"quenched":[]}} +{"case":{"desc":"from 82.65.127.86:40000 relax=True","from":["82.65.127.86",40000],"kind":"bcka","name":"bcka from 82.65.127.86:40000 relax=True","relax":true,"stream":1358954578},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* BCKA from 82.65.127.86:40000 is not this bridge's peer 82.65.127.86:62201 (keepalive only; a second instance of the peer, or another bridge sharing the passphrase)"]],"quenched":[]}} +{"case":{"desc":"from 9.9.9.9:62201 relax=False","from":["9.9.9.9",62201],"kind":"bcka","name":"bcka from 9.9.9.9:62201 relax=False","relax":false,"stream":1358954579},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* BCKA from 9.9.9.9:62201 is not this bridge's peer 82.65.127.86:62201 (keepalive only; a second instance of the peer, or another bridge sharing the passphrase)"]],"quenched":[]}} +{"case":{"desc":"from 9.9.9.9:62201 relax=True","from":["9.9.9.9",62201],"kind":"bcka","name":"bcka from 9.9.9.9:62201 relax=True","relax":true,"stream":1358954580},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* BCKA from 9.9.9.9:62201 is not this bridge's peer 82.65.127.86:62201 (keepalive only; a second instance of the peer, or another bridge sharing the passphrase)"]],"quenched":[]}} {"case":{"desc":"from 82.65.127.86:62201 relax=False","from":["82.65.127.86",62201],"kind":"bcsq","name":"bcsq from 82.65.127.86:62201 relax=False","relax":false,"stream":1358954581},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* BCSQ accepted: stream_id=1358954581 TGID=214 (peer quenched; forwarding on this OBP stops for this stream/TG)"]],"quenched":[]}} {"case":{"desc":"from 82.65.127.86:62201 relax=True","from":["82.65.127.86",62201],"kind":"bcsq","name":"bcsq from 82.65.127.86:62201 relax=True","relax":true,"stream":1358954582},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* BCSQ accepted: stream_id=1358954582 TGID=214 (peer quenched; forwarding on this OBP stops for this stream/TG)"]],"quenched":[]}} {"case":{"desc":"from 82.65.127.86:40000 relax=False","from":["82.65.127.86",40000],"kind":"bcsq","name":"bcsq from 82.65.127.86:40000 relax=False","relax":false,"stream":1358954583},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[20,"(OBP-FR) *BridgeControl* BCSQ accepted: stream_id=1358954583 TGID=214 (peer quenched; forwarding on this OBP stops for this stream/TG)"]],"quenched":[]}} @@ -94,3 +94,4 @@ {"case":{"desc":"from 82.65.127.86:40000 relax=True","from":["82.65.127.86",40000],"kind":"bcst","name":"bcst from 82.65.127.86:40000 relax=True","relax":true,"stream":1358954590},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[5,"(OBP-FR) *BridgeControl* BCST STUN request received"],[20,"(OBP-FR) Bridge STUNned, discarding"]],"quenched":[]}} {"case":{"desc":"from 9.9.9.9:62201 relax=False","from":["9.9.9.9",62201],"kind":"bcst","name":"bcst from 9.9.9.9:62201 relax=False","relax":false,"stream":1358954591},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[5,"(OBP-FR) *BridgeControl* BCST STUN request received"],[20,"(OBP-FR) Bridge STUNned, discarding"]],"quenched":[]}} {"case":{"desc":"from 9.9.9.9:62201 relax=True","from":["9.9.9.9",62201],"kind":"bcst","name":"bcst from 9.9.9.9:62201 relax=True","relax":true,"stream":1358954592},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: 82.65.127.86, TARGET_PORT: 62201"],[5,"(OBP-FR) *BridgeControl* BCST STUN request received"],[20,"(OBP-FR) Bridge STUNned, discarding"]],"quenched":[]}} +{"case":{"desc":"with no TARGET_IP configured","from":["82.65.127.86",62201],"kind":"bcka","name":"bcka with no TARGET_IP configured","no_peer":true,"stream":1358954593},"effects":{"delivered":[],"egress":[[24,["82.65.127.86",62201]],[73,["82.65.127.86",62201]]],"log":[[20,"(OBP-FR) Starting OBP. TARGET_IP: , TARGET_PORT: 62201"],[10,"(OBP-FR) *BridgeControl* not sending KeepAlive, TARGET not currently known"],[20,"(OBP-FR) *BridgeControl* OBP peer address learned from keepalive: 82.65.127.86:62201"]],"quenched":[]}} diff --git a/tests/harness/obp_ingress.py b/tests/harness/obp_ingress.py index 4100e29..b1171f9 100644 --- a/tests/harness/obp_ingress.py +++ b/tests/harness/obp_ingress.py @@ -122,6 +122,9 @@ def build_config(case: dict[str, Any]) -> dict[str, Any]: "_PEER_IDS": {}, "_LOCAL_SUBSCRIBER_IDS": {}, } + if case.get("no_peer"): # inbound-only bridge: the operator set no TARGET_IP + system["TARGET_IP"] = None + system["TARGET_SOCK"] = (None, PEER[1]) if case.get("stun"): config["STUN"] = True return config @@ -275,6 +278,7 @@ def _cases() -> list[dict[str, Any]]: ): for relax in (False, True): add(kind=kind, desc=f"from {addr[0]}:{addr[1]} relax={relax}", relax=relax, **{"from": list(addr)}) + add(kind="bcka", desc="with no TARGET_IP configured", no_peer=True, **{"from": list(PEER)}) return cases diff --git a/tests/infrastructure/test_obp_ingress_effects.py b/tests/infrastructure/test_obp_ingress_effects.py index d3946c7..d60743b 100644 --- a/tests/infrastructure/test_obp_ingress_effects.py +++ b/tests/infrastructure/test_obp_ingress_effects.py @@ -30,7 +30,16 @@ read it before regenerating with ``CAPTURE=1``. from __future__ import annotations import pytest -from tests.harness.obp_ingress import CASES, FIXTURE, as_json, capture_enabled, load, observe, record +from tests.harness.obp_ingress import ( + CASES, + FIXTURE, + PEER, + as_json, + capture_enabled, + load, + observe, + record, +) @pytest.fixture(scope="module") @@ -50,3 +59,28 @@ def test_corpus_covers_every_case(recorded: dict[str, dict]) -> None: def test_ingress_effects_match_the_recording(case: dict, recorded: dict[str, dict]) -> None: expected = recorded[case["name"]]["effects"] assert as_json(observe(case)) == expected + + +@pytest.mark.parametrize( + "case", + [c for c in CASES if c["kind"] == "bcka" and not c.get("no_peer")], + ids=lambda c: c["name"], +) +def test_a_keepalive_never_moves_egress_off_the_peer(case: dict) -> None: + """A keepalive carries no NETWORK_ID, so anyone holding the passphrase can send + one: a second instance of the peer, or another bridge on a shared-passphrase + mesh. It must not decide where this bridge transmits. Recorded above as well, + but asserted here so regenerating the corpus cannot drop it. + """ + for _size, addr in observe(case)["egress"]: + assert tuple(addr) == PEER + + +def test_a_keepalive_bootstraps_a_bridge_with_no_configured_peer() -> None: + """The other half: with no TARGET_IP there is nothing to protect and nothing to + steal, and the keepalive is the only way an inbound-only bridge learns where to + answer. Bootstrapping an unknown peer stays allowed. + """ + case = next(c for c in CASES if c.get("no_peer")) + for _size, addr in observe(case)["egress"]: + assert tuple(addr) == PEER diff --git a/tests/infrastructure/test_obp_proxy.py b/tests/infrastructure/test_obp_proxy.py index 4a804b1..49a5e2c 100644 --- a/tests/infrastructure/test_obp_proxy.py +++ b/tests/infrastructure/test_obp_proxy.py @@ -37,6 +37,7 @@ from adn_server.application.proxy.deployment import ( ) from adn_server.domain import bytes_4 from adn_server.domain.errors import ConfigError +from adn_server.domain.mesh_session import obp_session from adn_server.infrastructure.config_normalizer import normalize_obp_config from adn_server.infrastructure.config_validator import validate_config from adn_server.infrastructure.hbp_constants import DMRD @@ -417,7 +418,7 @@ def _shared_passphrase_registry() -> tuple[ObpBridgeRegistry, dict[str, _Recordi passphrase=_PASS, sink=InProcessObpSink(proto), reply_transport=ObpIngressReplyTransport(_RecordingTransport()), - sys_cfg={"TARGET_SOCK": peer, "TARGET_IP": peer[0], "TARGET_PORT": peer[1]}, + peer_hint=peer, ) ) return registry, protos @@ -548,9 +549,9 @@ def test_reload_keeps_egress_pinned_to_the_bridge_socket(monkeypatch) -> None: assert fake.transports[62032].sent == [] -def test_control_frame_prefers_configured_peer_over_relaxed_target() -> None: - """RELAX_CHECKS rewrites TARGET_SOCK in place; a bridge dragged onto another - bridge's address must not start stealing that peer's control frames.""" +def test_control_frame_prefers_configured_peer_over_learned_one() -> None: + """A bridge that learned another bridge's address must not start stealing that + peer's control frames: the configured peer outranks anything learned.""" config = _obp_mesh_config() protocols = {name: _RecordingObp() for name in ("OBP-FR", "OBP-PT")} registry = build_obp_bridge_registry( @@ -561,10 +562,7 @@ def test_control_frame_prefers_configured_peer_over_relaxed_target() -> None: primary_transport=_RecordingTransport(), ) peer_pt = ("85.241.222.7", 62268) - # What _obp_sync_target_sock_from_peer() does after a misattributed frame. - poisoned = config["SYSTEMS"]["OBP-FR"] - poisoned["TARGET_IP"], poisoned["TARGET_PORT"] = peer_pt - poisoned["TARGET_SOCK"] = peer_pt + obp_session(config, "OBP-FR").learn_peer(peer_pt, at=1.0) demux = ObpFanInDemux(registry) demux.deliver(build_bcka(_PASS), peer_pt, local_port=62032, transport=_RecordingTransport()) @@ -572,6 +570,28 @@ def test_control_frame_prefers_configured_peer_over_relaxed_target() -> None: assert not protocols["OBP-FR"].packets +def test_control_frame_follows_a_peer_that_moved() -> None: + """Voice (DMRD/DMRE, which carries NETWORK_ID) teaches the session where a peer + really is; control frames from there must resolve to it, not to whichever bridge + happens to be registered first.""" + config = _obp_mesh_config() + protocols = {name: _RecordingObp() for name in ("OBP-FR", "OBP-PT")} + registry = build_obp_bridge_registry( + config, + protocols, + bind_legacy_ports=True, + listen_port=62032, + primary_transport=_RecordingTransport(), + ) + moved_to = ("129.80.176.29", 62268) + obp_session(config, "OBP-PT").learn_peer(moved_to, at=1.0) + + demux = ObpFanInDemux(registry) + demux.deliver(build_bcka(_PASS), moved_to, local_port=62032, transport=_RecordingTransport()) + assert protocols["OBP-PT"].packets, "the keepalive did not follow the peer that moved" + assert not protocols["OBP-FR"].packets + + def test_debug_does_not_log_voice_packets(caplog) -> None: """DMRD/DMRE demux by NETWORK_ID, never ambiguous, and *CALL START*/*CALL END* (routing_use_cases.py) already give once-per-call visibility. Concurrent calls