From 5c24ecb48bb7bd52dc962b0e48a64b61467f3c40 Mon Sep 17 00:00:00 2001 From: yo Date: Sun, 20 Sep 2026 09:50:39 +0200 Subject: [PATCH] feat(obp): replay a capture through the ingress to see why frames were refused MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The engine from phase 3 takes plain values and answers with effects, so it does not need the server to run. ``adn-server --replay capture.pcap`` uses that: the frames from a tcpdump capture go through the same ingress, against the operator's own adn-server.yaml, and each one comes back with the bridge it belongs to and either a delivery or the reason it was refused. 12:04:31 82.65.127.86:62201 OBP-FR DMRD v1 2130001 -> 214 delivered 12:04:31 85.241.222.7:62268 OBP-PT DMRE v5 2680015 -> 9 dropped (tg-filter-server) +BCSQ 12:04:32 203.0.113.9:50000 - DMRD v1 unmatched Nothing is sent and no port is bound, so it runs beside a live master. Which bridge a frame belongs to is decided on the evidence the server has — the port it arrived on, then the configured peer, then whoever can verify it — which is also the answer to "whose keepalive is this?" on a mesh where every bridge shares one passphrase. ``--system`` narrows it to one link, ``--replay-limit`` stops early and ``--replay-summary`` prints the tally alone. ``infrastructure/pcap.py`` reads classic pcap with no dependencies: both endiannesses, microsecond and nanosecond timestamps, Ethernet (VLAN tags included), Linux cooked v1 and v2 (``tcpdump -i any`` writes SLL2, found while running this against a real capture), raw IP and loopback, IPv4 and IPv6 UDP. pcapng says which command converts it. Docs: the OBP proxy page gains a "why did that call not cross" section in both languages, and its RELAX_CHECKS note now says what phase 2 made true — what the wire teaches lives in the session, TARGET_IP stays as written. Tests: 27 new, 97% of the replay module and 87% of the pcap reader, plus an end-to-end run of the real command against a real YAML. Full suite 1025 passed, 2 skipped (the 2 failures are this machine's and fail on develop too). Co-Authored-By: Claude Opus 5 --- docs/en/server/user-guide/obp-proxy.md | 34 +- docs/es/server/user-guide/obp-proxy.md | 34 +- src/adn_server/infrastructure/obp_replay.py | 339 +++++++++++++++++ src/adn_server/infrastructure/pcap.py | 180 +++++++++ src/adn_server/main.py | 41 +- tests/infrastructure/test_obp_replay.py | 400 ++++++++++++++++++++ 6 files changed, 1025 insertions(+), 3 deletions(-) create mode 100644 src/adn_server/infrastructure/obp_replay.py create mode 100644 src/adn_server/infrastructure/pcap.py create mode 100644 tests/infrastructure/test_obp_replay.py diff --git a/docs/en/server/user-guide/obp-proxy.md b/docs/en/server/user-guide/obp-proxy.md index b79b41f..64e6f81 100644 --- a/docs/en/server/user-guide/obp-proxy.md +++ b/docs/en/server/user-guide/obp-proxy.md @@ -46,6 +46,38 @@ Example: migrate `OBP-CL2` to the shared fan-in while `OBP-EU` keeps `PORT: 6299 - `NETWORK_ID` must be unique among enabled OPENBRIDGE systems. - `LISTEN_PORT` must not collide with any OPENBRIDGE `PORT` when `BIND_LEGACY_PORTS` is true. -- `RELAX_CHECKS: true` is recommended so `TARGET_SOCK` is learned from the first valid packet. +- `RELAX_CHECKS: true` is recommended so the peer address is learned from the first valid packet. What is learned lives in the bridge's session, not in the config: `TARGET_IP` / `TARGET_PORT` stay as written, and a reload puts the link back on them. + +## Why did that call not cross? + +Replay a capture through the same ingress the server runs, offline and against +your own `adn-server.yaml`. Nothing is sent and no port is bound, so the server +can keep running: + +```bash +tcpdump -i any -n -w /tmp/obp.pcap udp port 62201 or udp port 62268 # a minute is plenty +adn-server -c adn-server.yaml --replay /tmp/obp.pcap +``` + +Each frame comes back with the bridge it belongs to and a verdict: + +``` +12:04:31 82.65.127.86:62201 OBP-FR DMRD v1 2130001 -> 214 delivered +12:04:31 85.241.222.7:62268 OBP-PT DMRE v5 2680015 -> 9 dropped (tg-filter-server) +BCSQ +12:04:32 203.0.113.9:50000 - DMRD v1 unmatched + +3 datagram(s) + OBP-FR + 1 delivered + OBP-PT + 1 dropped: tg-filter-server + (no bridge) + 1 unmatched +``` + +`unmatched` means no enabled bridge could verify the frame with its passphrase. +Add `--system OBP-FR` to look at one link, `--replay-limit N` to stop early and +`--replay-summary` for the tally alone. Classic pcap only; convert a pcapng with +`editcap -F pcap in.pcapng out.pcap`. See also: [OpenBridge protocol](../protocols/openbridge.md). diff --git a/docs/es/server/user-guide/obp-proxy.md b/docs/es/server/user-guide/obp-proxy.md index d3f4337..10d7bcf 100644 --- a/docs/es/server/user-guide/obp-proxy.md +++ b/docs/es/server/user-guide/obp-proxy.md @@ -44,6 +44,38 @@ Ejemplo: migrar `OBP-CL2` al fan-in compartido mientras `OBP-EU` conserva `PORT: - `NETWORK_ID` único entre OPENBRIDGE habilitados. - `LISTEN_PORT` sin colisión con ningún `PORT` de sección si `BIND_LEGACY_PORTS` es true. -- `RELAX_CHECKS: true` recomendado para aprender `TARGET_SOCK` del primer paquete válido. +- `RELAX_CHECKS: true` recomendado para aprender la dirección del peer del primer paquete válido. Lo aprendido vive en la sesión del bridge, no en la config: `TARGET_IP` / `TARGET_PORT` se quedan como están escritos, y una recarga devuelve el enlace a ellos. + +## ¿Por qué no cruzó esa llamada? + +Pasa una captura por el mismo ingress que corre el servidor, sin conexión y +contra tu propio `adn-server.yaml`. No se envía nada ni se abre ningún puerto, +así que el servidor puede seguir funcionando: + +```bash +tcpdump -i any -n -w /tmp/obp.pcap udp port 62201 or udp port 62268 # con un minuto sobra +adn-server -c adn-server.yaml --replay /tmp/obp.pcap +``` + +Cada trama vuelve con el bridge al que pertenece y un veredicto: + +``` +12:04:31 82.65.127.86:62201 OBP-FR DMRD v1 2130001 -> 214 delivered +12:04:31 85.241.222.7:62268 OBP-PT DMRE v5 2680015 -> 9 dropped (tg-filter-server) +BCSQ +12:04:32 203.0.113.9:50000 - DMRD v1 unmatched + +3 datagram(s) + OBP-FR + 1 delivered + OBP-PT + 1 dropped: tg-filter-server + (no bridge) + 1 unmatched +``` + +`unmatched` significa que ningún bridge habilitado pudo verificar la trama con +su passphrase. `--system OBP-FR` mira un solo enlace, `--replay-limit N` corta +antes y `--replay-summary` deja solo el recuento. Solo pcap clásico; un pcapng +se convierte con `editcap -F pcap in.pcapng out.pcap`. Ver también: [protocolo OpenBridge](../protocols/openbridge.md). diff --git a/src/adn_server/infrastructure/obp_replay.py b/src/adn_server/infrastructure/obp_replay.py new file mode 100644 index 0000000..41e3a20 --- /dev/null +++ b/src/adn_server/infrastructure/obp_replay.py @@ -0,0 +1,339 @@ +# ADN DMR Peer Server - infrastructure obp replay +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Replay a capture through the OpenBridge ingress and say what it would do. + +Offline answer to "why did that call not cross": the frames from a ``tcpdump`` +capture go through the same engine the server runs, against the operator's own +adn-server.yaml, and each one comes back with the bridge it belongs to and +either a delivery or the reason it was refused. Nothing is sent, nothing is +bound; the server can keep running while this reads the capture. + +This is what the engine being pure buys — the decisions can be taken anywhere. +""" + +from __future__ import annotations + +import sys +import time +from collections import Counter +from dataclasses import dataclass +from pathlib import Path +from typing import Any, TextIO + +from adn_server.application.proxy.deployment import obp_bridge_legacy_listen_port +from adn_server.domain import int_id +from adn_server.domain.mesh_admission import AclRules, AdmissionContext, server_prefix +from adn_server.domain.mesh_engine import ( + BridgePolicy, + Deliver, + Reject, + accepts_source, + ingest_dmrd_v1, + ingest_dmre_v5, + reject_v1_protocol, + server_id_bytes, +) +from adn_server.domain.mesh_routing import PeerMeshConfig +from adn_server.domain.mesh_session import MeshSessionStore, ObpBridgeSession +from adn_server.infrastructure.acl_router import InMemoryAclRouter +from adn_server.infrastructure.hbp_constants import BC, BCKA, BCSQ, BCST, BCVE, DMRD, DMRE +from adn_server.infrastructure.mesh.dmre_v5 import parse_dmre_trailer +from adn_server.infrastructure.mesh.registry import MeshCodecRegistry +from adn_server.infrastructure.pcap import CaptureError, CapturedDatagram, read_udp + +_CONTROL_NAMES = {BCKA: "BCKA", BCSQ: "BCSQ", BCST: "BCST", BCVE: "BCVE"} + + +@dataclass +class Verdict: + """What the ingress would do with one captured datagram.""" + + datagram: CapturedDatagram + system: str | None = None + kind: str = "?" + rf_src: int = 0 + dst_id: int = 0 + outcome: str = "unmatched" + reason: str = "" + quench: bool = False + + def line(self) -> str: + when = time.strftime("%H:%M:%S", time.localtime(self.datagram.timestamp)) + source = f"{self.datagram.source[0]}:{self.datagram.source[1]}" + system = self.system or "-" + call = f"{self.rf_src} -> {self.dst_id}" if self.rf_src or self.dst_id else "" + verdict = self.outcome if not self.reason else f"{self.outcome} ({self.reason})" + if self.quench: + verdict += " +BCSQ" + return f"{when} {source:<24} {system:<12} {self.kind:<9} {call:<24} {verdict}" + + +def _openbridge_systems(config: dict[str, Any]) -> dict[str, dict[str, Any]]: + return { + name: sys_cfg + for name, sys_cfg in config.get("SYSTEMS", {}).items() + if isinstance(sys_cfg, dict) + and sys_cfg.get("MODE") == "OPENBRIDGE" + and sys_cfg.get("ENABLED", True) + } + + +def _policy(config: dict[str, Any], name: str, sys_cfg: dict[str, Any], router: Any) -> BridgePolicy: + global_cfg = config.get("GLOBAL", {}) + admission = AdmissionContext( + stunned="STUN" in config, + acl_check=router.acl_check, + global_rules=AclRules( + enabled=bool(global_cfg.get("USE_ACL")), + sub_acl=global_cfg.get("SUB_ACL", (True, [])), + tg1_acl=global_cfg.get("TG1_ACL", (True, [])), + ), + system_rules=AclRules( + enabled=bool(sys_cfg.get("USE_ACL")), + sub_acl=sys_cfg.get("SUB_ACL", (True, [])), + tg1_acl=sys_cfg.get("TG1_ACL", (True, [])), + ), + server_id=server_prefix(global_cfg.get("SERVER_ID", 0)), + validate_server_ids=bool(global_cfg.get("VALIDATE_SERVER_IDS")), + known_server_prefixes=config.get("_SERVER_IDS", set()), + resolve_server_id=lambda _sid: True, # alias tables are not loaded offline + ) + return BridgePolicy( + system=name, + network_id=sys_cfg.get("NETWORK_ID", b""), + proto_ver=sys_cfg.get("VER"), + relax_checks=bool(sys_cfg.get("RELAX_CHECKS")), + server_id=server_id_bytes(global_cfg.get("SERVER_ID", 0)), + admission=admission, + ) + + +def _mesh_config(sys_cfg: dict[str, Any], server_id: bytes) -> PeerMeshConfig: + passphrase = sys_cfg.get("PASSPHRASE") or b"" + if isinstance(passphrase, str): + passphrase = (passphrase.strip().encode("utf-8") + b"\x00" * 20)[:20] + ver = sys_cfg.get("VER") + return PeerMeshConfig( + passphrase=passphrase, + server_id=server_id, + wire_ver=int(ver) if ver is not None else None, + ) + + +def _candidates( + systems: dict[str, dict[str, Any]], + datagram: CapturedDatagram, + *, + listen_port: int, +) -> list[str]: + """Bridges to try for this datagram, the likeliest first. + + Same evidence the server has: the port it arrived on, then the configured + peer address, then everyone else — a shared passphrase makes the rest + ambiguous, which is exactly why the order matters. + """ + host, port = datagram.source + local_port = datagram.destination[1] + exact: list[str] = [] + same_host: list[str] = [] + rest: list[str] = [] + for name, sys_cfg in systems.items(): + bridge_port = obp_bridge_legacy_listen_port( + sys_cfg, listen_port=listen_port, bind_legacy_ports=True + ) + if bridge_port == local_port: + exact.insert(0, name) + continue + sock = sys_cfg.get("TARGET_SOCK") + target_host = sock[0] if isinstance(sock, tuple) else sys_cfg.get("TARGET_IP") + target_port = sock[1] if isinstance(sock, tuple) else sys_cfg.get("TARGET_PORT") + if target_host == host and target_port == port: + exact.append(name) + elif target_host == host: + same_host.append(name) + else: + rest.append(name) + return exact + same_host + rest + + +def _control_verdict(datagram: CapturedDatagram, payload: bytes, system: str | None) -> Verdict: + name = _CONTROL_NAMES.get(payload[:4], "BC?") + return Verdict( + datagram=datagram, + system=system, + kind=name, + outcome="control", + ) + + +def replay( + config: dict[str, Any], + datagrams: list[CapturedDatagram], + *, + system: str | None = None, + now: float | None = None, +) -> list[Verdict]: + """Run captured datagrams through the ingress engine. Sends nothing.""" + systems = _openbridge_systems(config) + if system is not None: + systems = {name: cfg for name, cfg in systems.items() if name == system} + if not systems: + raise KeyError(f"no enabled OPENBRIDGE system named {system!r}") + listen_port = int(config.get("OBP_PROXY", {}).get("LISTEN_PORT", 62032) or 62032) + router = InMemoryAclRouter() + registry = MeshCodecRegistry() + store = MeshSessionStore() + store.sync(config) + server_id = server_id_bytes(config.get("GLOBAL", {}).get("SERVER_ID", 0)) + + verdicts: list[Verdict] = [] + for datagram in datagrams: + payload = datagram.payload + if len(payload) < 4: + continue + opcode = payload[:4] + if opcode[:2] == BC and opcode in _CONTROL_NAMES: + names = _candidates(systems, datagram, listen_port=listen_port) + verdicts.append(_control_verdict(datagram, payload, names[0] if names else None)) + continue + if opcode not in (DMRD, DMRE): + continue + verdict = Verdict(datagram=datagram, kind="DMRD v1" if opcode == DMRD else "DMRE v5") + for name in _candidates(systems, datagram, listen_port=listen_port): + sys_cfg = systems[name] + ingress = registry.decode_auto(payload, _mesh_config(sys_cfg, server_id)) + if ingress is None: + continue + session = store.session(name, sys_cfg) + policy = _policy(config, name, sys_cfg, router) + verdict.system = name + frame = ingress.voice_frame + verdict.rf_src = int(int_id(frame[5:8])) + verdict.dst_id = int(int_id(frame[8:11])) + effects = _effects_for( + opcode, ingress, payload, datagram, policy=policy, session=session, now=now + ) + _fill(verdict, effects) + break + verdicts.append(verdict) + return verdicts + + +def _effects_for( + opcode: bytes, + ingress: Any, + payload: bytes, + datagram: CapturedDatagram, + *, + policy: BridgePolicy, + session: ObpBridgeSession, + now: float | None, +) -> list[Any] | None: + moment = datagram.timestamp if now is None else now + if opcode == DMRD: + if policy.rejects_v1: + return reject_v1_protocol(ingress.voice_frame[16:20], policy=policy) + if not accepts_source(datagram.source, policy=policy, session=session): + return None + return ingest_dmrd_v1(ingress, datagram.source, policy=policy, session=session, now=moment) + trailer = parse_dmre_trailer(payload) + timestamp = trailer.timestamp if trailer is not None else b"\x00" * 8 + if not accepts_source(datagram.source, policy=policy, session=session): + return None + return ingest_dmre_v5( + ingress, + datagram.source, + policy=policy, + session=session, + timestamp_ns=int.from_bytes(timestamp, "big"), + now=moment, + ) + + +def _fill(verdict: Verdict, effects: list[Any] | None) -> None: + if effects is None: + verdict.outcome = "refused" + verdict.reason = "source not accepted" + return + for effect in effects: + if isinstance(effect, Reject): + verdict.outcome = "dropped" + verdict.reason = effect.reason + verdict.quench = effect.rejection.quench + return + if isinstance(effect, Deliver): + verdict.outcome = "delivered" + return + verdict.outcome = "ignored" + + +def format_report(verdicts: list[Verdict], *, capture: str, verbose: bool = True) -> str: + """The per-frame lines and the tally the sysop actually reads.""" + lines = [f"OBP replay of {capture}", ""] + if verbose: + lines.extend(verdict.line() for verdict in verdicts) + lines.append("") + by_system: dict[str, Counter] = {} + for verdict in verdicts: + label = verdict.system or "(no bridge)" + key = verdict.outcome if not verdict.reason else f"{verdict.outcome}: {verdict.reason}" + by_system.setdefault(label, Counter())[key] += 1 + lines.append(f"{len(verdicts)} datagram(s)") + for label in sorted(by_system): + lines.append(f" {label}") + for key, count in sorted(by_system[label].items(), key=lambda item: (-item[1], item[0])): + lines.append(f" {count:>6} {key}") + return "\n".join(lines) + + +def run_replay( + config: dict[str, Any], + capture_path: str, + *, + system: str | None = None, + limit: int | None = None, + summary_only: bool = False, + out: TextIO | None = None, +) -> int: + """Read a capture, replay it, print the report. Returns 0 when it ran.""" + stream = out or sys.stdout + try: + datagrams = list(read_udp(Path(capture_path))) + except (CaptureError, OSError) as exc: + print(f"ERROR capture: {exc}", file=sys.stderr) + return 1 + if limit is not None: + datagrams = datagrams[:limit] + try: + verdicts = replay(config, datagrams, system=system) + except KeyError as exc: + print(f"ERROR system: {exc}", file=sys.stderr) + return 1 + print(format_report(verdicts, capture=capture_path, verbose=not summary_only), file=stream) + return 0 + + +__all__ = [ + "Verdict", + "format_report", + "replay", + "run_replay", +] diff --git a/src/adn_server/infrastructure/pcap.py b/src/adn_server/infrastructure/pcap.py new file mode 100644 index 0000000..ea95d1c --- /dev/null +++ b/src/adn_server/infrastructure/pcap.py @@ -0,0 +1,180 @@ +# ADN DMR Peer Server - infrastructure pcap reader +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Read UDP datagrams out of a classic pcap file, without dependencies. + +Enough of the format to walk a ``tcpdump -w`` capture and hand back what a +socket would have received: the payload, who sent it, which port it arrived on +and when. Anything that is not IPv4/IPv6 UDP is skipped. +""" + +from __future__ import annotations + +import struct +from collections.abc import Iterator +from dataclasses import dataclass +from pathlib import Path + +PCAP_MAGIC_US = 0xA1B2C3D4 # timestamps in microseconds +PCAP_MAGIC_NS = 0xA1B23C4D # timestamps in nanoseconds +PCAPNG_MAGIC = 0x0A0D0D0A + +LINKTYPE_NULL = 0 +LINKTYPE_ETHERNET = 1 +LINKTYPE_RAW = 101 +LINKTYPE_LINUX_SLL = 113 +LINKTYPE_LINUX_SLL2 = 276 +LINKTYPE_IPV4 = 228 +LINKTYPE_IPV6 = 229 + + +class CaptureError(Exception): + """The file is not a capture this reader can walk.""" + + +@dataclass(frozen=True) +class CapturedDatagram: + """One UDP datagram as it appeared on the wire.""" + + timestamp: float + source: tuple[str, int] + destination: tuple[str, int] + payload: bytes + + +def _ipv4(raw: bytes) -> str: + return ".".join(str(b) for b in raw) + + +def _ipv6(raw: bytes) -> str: + parts = [f"{raw[i] << 8 | raw[i + 1]:x}" for i in range(0, 16, 2)] + return ":".join(parts) + + +def _strip_link_layer(frame: bytes, linktype: int) -> tuple[bytes, int] | None: + """Return the network-layer payload and its ethertype-ish family.""" + if linktype == LINKTYPE_ETHERNET: + if len(frame) < 14: + return None + ethertype = int.from_bytes(frame[12:14], "big") + offset = 14 + while ethertype in (0x8100, 0x88A8): # VLAN tags + if len(frame) < offset + 4: + return None + ethertype = int.from_bytes(frame[offset + 2 : offset + 4], "big") + offset += 4 + return frame[offset:], ethertype + if linktype == LINKTYPE_LINUX_SLL: + if len(frame) < 16: + return None + return frame[16:], int.from_bytes(frame[14:16], "big") + if linktype == LINKTYPE_LINUX_SLL2: + # `tcpdump -i any` on a recent libpcap: protocol first, 20-byte header + if len(frame) < 20: + return None + return frame[20:], int.from_bytes(frame[:2], "big") + if linktype == LINKTYPE_NULL: + if len(frame) < 4: + return None + family = int.from_bytes(frame[:4], "little") + return frame[4:], 0x0800 if family == 2 else 0x86DD + if linktype in (LINKTYPE_RAW, LINKTYPE_IPV4, LINKTYPE_IPV6): + if not frame: + return None + version = frame[0] >> 4 + return frame, 0x0800 if version == 4 else 0x86DD + return None + + +def _udp_from_ip(packet: bytes, ethertype: int) -> tuple[str, str, bytes] | None: + """Return ``(src_ip, dst_ip, udp_segment)`` for an IPv4/IPv6 UDP packet.""" + if ethertype == 0x0800: + if len(packet) < 20 or packet[0] >> 4 != 4: + return None + header_len = (packet[0] & 0x0F) * 4 + if packet[9] != 17 or len(packet) < header_len + 8: # 17 = UDP + return None + return _ipv4(packet[12:16]), _ipv4(packet[16:20]), packet[header_len:] + if ethertype == 0x86DD: + if len(packet) < 40 or packet[0] >> 4 != 6: + return None + if packet[6] != 17 or len(packet) < 48: # no extension-header walking + return None + return _ipv6(packet[8:24]), _ipv6(packet[24:40]), packet[40:] + return None + + +def read_udp(path: str | Path) -> Iterator[CapturedDatagram]: + """Walk a capture and yield its UDP datagrams in order.""" + path = Path(path) + with path.open("rb") as fh: + header = fh.read(24) + if len(header) < 24: + raise CaptureError(f"{path}: too short to be a capture") + magic = int.from_bytes(header[:4], "big") + if magic == PCAPNG_MAGIC: + raise CaptureError( + f"{path}: pcapng is not supported; convert it first " + "(editcap -F pcap in.pcapng out.pcap)" + ) + if magic in (PCAP_MAGIC_US, PCAP_MAGIC_NS): + endian = ">" + else: + magic = int.from_bytes(header[:4], "little") + if magic not in (PCAP_MAGIC_US, PCAP_MAGIC_NS): + raise CaptureError(f"{path}: not a pcap file") + endian = "<" + divisor = 1_000_000_000 if magic == PCAP_MAGIC_NS else 1_000_000 + linktype = struct.unpack(endian + "I", header[20:24])[0] + + record = struct.Struct(endian + "IIII") + while True: + raw = fh.read(record.size) + if len(raw) < record.size: + return + seconds, fraction, captured_len, _original_len = record.unpack(raw) + frame = fh.read(captured_len) + if len(frame) < captured_len: + return + stripped = _strip_link_layer(frame, linktype) + if stripped is None: + continue + packet, ethertype = stripped + addresses = _udp_from_ip(packet, ethertype) + if addresses is None: + continue + source_ip, destination_ip, segment = addresses + if len(segment) < 8: + continue + source_port, destination_port, length = struct.unpack(">HHH", segment[:6]) + payload = segment[8 : max(8, length)] if length >= 8 else segment[8:] + yield CapturedDatagram( + timestamp=seconds + fraction / divisor, + source=(source_ip, source_port), + destination=(destination_ip, destination_port), + payload=payload, + ) + + +__all__ = [ + "CaptureError", + "CapturedDatagram", + "read_udp", +] diff --git a/src/adn_server/main.py b/src/adn_server/main.py index 2704100..899e799 100644 --- a/src/adn_server/main.py +++ b/src/adn_server/main.py @@ -28,6 +28,7 @@ ADN DMR Peer Server entrypoint. Run: python -m adn_server.main [-c adn-server.yaml] [--logging LEVEL] python -m adn_server.main --echo [-c adn-echo.yaml] python -m adn_server.main --doctor [-c adn-server.yaml] + python -m adn_server.main --replay capture.pcap [--system OBP-FR] Config default: adn-server.yaml (or adn-echo.yaml with --echo). """ @@ -46,8 +47,9 @@ if str(_ROOT) not in sys.path: from adn_server.domain.errors import ConfigError from adn_server.infrastructure import YamlConfigLoader, setup_logging from adn_server.infrastructure.bootstrap import run_peer_server -from adn_server.infrastructure.config_normalizer import apply_talker_alias_defaults +from adn_server.infrastructure.config_normalizer import apply_talker_alias_defaults, normalize_obp_config from adn_server.infrastructure.doctor import run_doctor +from adn_server.infrastructure.obp_replay import run_replay from adn_server.infrastructure.echo import run_echo @@ -79,6 +81,31 @@ def _parse_args() -> argparse.Namespace: action="store_true", help="Validate config, ports, and peers; exit non-zero on errors", ) + parser.add_argument( + "--replay", + dest="REPLAY_CAPTURE", + default=None, + metavar="CAPTURE.pcap", + help="Replay OpenBridge frames from a pcap through the ingress and report what it would do", + ) + parser.add_argument( + "--system", + dest="REPLAY_SYSTEM", + default=None, + help="With --replay: only this OPENBRIDGE system", + ) + parser.add_argument( + "--replay-limit", + dest="REPLAY_LIMIT", + type=int, + default=None, + help="With --replay: stop after this many datagrams", + ) + parser.add_argument( + "--replay-summary", + action="store_true", + help="With --replay: print the tally only, not one line per frame", + ) parser.add_argument("--version", action="version", version=f"adn-server {__version__}") return parser.parse_args() @@ -119,6 +146,18 @@ def main() -> None: print(f"(CONFIG) {exc}", file=sys.stderr) sys.exit(1) + if args.REPLAY_CAPTURE: + normalize_obp_config(config) + sys.exit( + run_replay( + config, + args.REPLAY_CAPTURE, + system=args.REPLAY_SYSTEM, + limit=args.REPLAY_LIMIT, + summary_only=args.replay_summary, + ) + ) + if not args.echo: apply_talker_alias_defaults(config) diff --git a/tests/infrastructure/test_obp_replay.py b/tests/infrastructure/test_obp_replay.py new file mode 100644 index 0000000..b565b6a --- /dev/null +++ b/tests/infrastructure/test_obp_replay.py @@ -0,0 +1,400 @@ +# ADN DMR Peer Server - tests infrastructure obp replay +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""Replaying a capture: read the pcap, ask the engine, report the verdict.""" + +from __future__ import annotations + +import struct +import time + +import pytest + +from adn_server.domain import bytes_3, bytes_4 +from adn_server.infrastructure.hbp_constants import BCKA, DMRD +from adn_server.infrastructure.mesh.dmre_v5 import build_dmre +from adn_server.infrastructure.mesh.obp_v1 import build_bcka, build_dmrd_v1 +from adn_server.infrastructure.obp_replay import format_report, replay, run_replay +from adn_server.infrastructure.pcap import CaptureError, read_udp + +_PASSPHRASE = (b"test-passphrase" + b"\x00" * 20)[:20] +_FR = ("82.65.127.86", 62201) +_PT = ("85.241.222.7", 62268) +_NOW = 1_800_000_000.0 + + +# --- building a capture on the fly ------------------------------------------- + + +def _udp_packet(source: tuple[str, int], destination: tuple[str, int], payload: bytes) -> bytes: + src_ip = bytes(int(part) for part in source[0].split(".")) + dst_ip = bytes(int(part) for part in destination[0].split(".")) + udp = struct.pack(">HHHH", source[1], destination[1], 8 + len(payload), 0) + payload + total = 20 + len(udp) + ip = struct.pack(">BBHHHBBH", 0x45, 0, total, 0, 0, 64, 17, 0) + src_ip + dst_ip + ethernet = b"\x02" * 6 + b"\x03" * 6 + b"\x08\x00" + return ethernet + ip + udp + + +def write_pcap(path, frames: list[tuple[float, tuple, tuple, bytes]]) -> None: + with open(path, "wb") as fh: + fh.write(struct.pack(" bytes: + body = b"".join( + [ + DMRD, + bytes([1]), + bytes_3(src), + bytes_3(dst), + bytes_4(20840), + bytes([bits]), + bytes_4(stream), + b"\x00" * 33, + ] + ) + return build_dmrd_v1(body, bytes_4(20840), _PASSPHRASE) + + +def _voice_v5(*, dst: int = 214, age: float = 0.0, now: float = _NOW) -> bytes: + body = b"".join( + [ + DMRD, + bytes([1]), + bytes_3(2130001), + bytes_3(dst), + bytes_4(20840), + bytes([0x21]), + bytes_4(0x4321), + b"\x00" * 33, + ] + ) + packet = build_dmre( + body, + server_id=bytes_4(20840), + ber=b"\x00", + rssi=b"\x00", + embedded_ver=5, + timestamp_ns=int((now - age) * 1_000_000_000), + source_server=bytes_4(2084), + source_rptr=bytes_4(0), + hops=b"\x00", + passphrase=_PASSPHRASE, + extended_layout=True, + ) + assert packet is not None + return packet + + +def _config() -> dict: + return { + "GLOBAL": { + "SERVER_ID": bytes_4(21310), + "USE_ACL": False, + "SUB_ACL": (True, []), + "TG1_ACL": (True, []), + }, + "SYSTEMS": { + "OBP-FR": { + "MODE": "OPENBRIDGE", + "ENABLED": True, + "NETWORK_ID": bytes_4(20840), + "PASSPHRASE": _PASSPHRASE, + "TARGET_IP": _FR[0], + "TARGET_PORT": _FR[1], + "TARGET_SOCK": _FR, + "RELAX_CHECKS": False, + "VER": 1, + "ENHANCED_OBP": True, + "_REPORT_PORT": 62201, + }, + "HOTSPOT": {"MODE": "MASTER", "ENABLED": True}, + }, + "_SERVER_IDS": set(), + } + + +# --- the pcap reader ---------------------------------------------------------- + + +def test_reading_udp_out_of_a_capture(tmp_path) -> None: + capture = tmp_path / "obp.pcap" + write_pcap(capture, [(_NOW, _FR, ("10.0.0.1", 62201), b"hello")]) + datagrams = list(read_udp(capture)) + assert len(datagrams) == 1 + assert datagrams[0].source == _FR + assert datagrams[0].destination == ("10.0.0.1", 62201) + assert datagrams[0].payload == b"hello" + assert int(datagrams[0].timestamp) == int(_NOW) + + +def test_a_file_that_is_not_a_capture(tmp_path) -> None: + path = tmp_path / "nope.pcap" + path.write_bytes(b"not a capture at all, really" * 2) + with pytest.raises(CaptureError): + list(read_udp(path)) + + +def test_pcapng_says_how_to_convert_it(tmp_path) -> None: + path = tmp_path / "new.pcapng" + path.write_bytes(b"\x0a\x0d\x0d\x0a" + b"\x00" * 40) + with pytest.raises(CaptureError, match="editcap"): + list(read_udp(path)) + + +def test_reading_a_linux_cooked_v2_capture(tmp_path) -> None: + """``tcpdump -i any`` on a recent libpcap writes SLL2, not Ethernet.""" + capture = tmp_path / "any.pcap" + packet = _udp_packet(_FR, ("10.0.0.1", 62201), b"hello")[14:] # drop the ethernet header + sll2 = struct.pack(">HHIHBB", 0x0800, 0, 2, 1, 0, 6) + b"\x02" * 8 + with open(capture, "wb") as fh: + fh.write(struct.pack(" None: + with open(path, "wb") as fh: + fh.write(struct.pack(" bytes: + return _udp_packet(_FR, ("10.0.0.1", 62201), b"hello")[14:] + + +def test_reading_a_raw_ip_capture(tmp_path) -> None: + capture = tmp_path / "raw.pcap" + _write_raw(capture, 101, _ip_udp()) + assert [d.payload for d in read_udp(capture)] == [b"hello"] + + +def test_reading_a_loopback_capture(tmp_path) -> None: + capture = tmp_path / "null.pcap" + _write_raw(capture, 0, struct.pack(" None: + capture = tmp_path / "vlan.pcap" + tagged = b"\x02" * 6 + b"\x03" * 6 + b"\x81\x00" + b"\x00\x64" + b"\x08\x00" + _ip_udp() + _write_raw(capture, 1, tagged) + assert [d.payload for d in read_udp(capture)] == [b"hello"] + + +def test_reading_an_ipv6_datagram(tmp_path) -> None: + capture = tmp_path / "v6.pcap" + payload = b"hello" + udp = struct.pack(">HHHH", _FR[1], 62201, 8 + len(payload), 0) + payload + header = struct.pack(">IHBB", 0x60000000, len(udp), 17, 64) + source = bytes.fromhex("20010db8000000000000000000000001") + destination = bytes.fromhex("20010db8000000000000000000000002") + _write_raw(capture, 101, header + source + destination + udp) + datagrams = list(read_udp(capture)) + assert datagrams[0].payload == payload + assert datagrams[0].source == ("2001:db8:0:0:0:0:0:1", _FR[1]) + + +def test_frames_that_are_not_udp_are_skipped(tmp_path) -> None: + capture = tmp_path / "tcp.pcap" + packet = bytearray(_ip_udp()) + packet[9] = 6 # TCP + _write_raw(capture, 101, bytes(packet)) + assert list(read_udp(capture)) == [] + + +def test_an_unknown_link_layer_is_skipped(tmp_path) -> None: + capture = tmp_path / "weird.pcap" + _write_raw(capture, 999, _ip_udp()) + assert list(read_udp(capture)) == [] + + +def test_a_truncated_record_ends_the_walk(tmp_path) -> None: + capture = tmp_path / "cut.pcap" + _write_raw(capture, 101, _ip_udp()) + with open(capture, "ab") as fh: + fh.write(struct.pack(" None: + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice())]) + assert [v.outcome for v in verdicts] == ["delivered"] + assert verdicts[0].system == "OBP-FR" + assert (verdicts[0].rf_src, verdicts[0].dst_id) == (2130001, 214) + assert verdicts[0].kind == "DMRD v1" + + +def test_a_local_talkgroup_is_reported_with_its_reason_and_quench() -> None: + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice(dst=9))]) + assert verdicts[0].outcome == "dropped" + assert verdicts[0].reason == "tg-filter" + assert verdicts[0].quench is True + + +def test_a_frame_from_an_unexpected_address_is_refused_when_checks_are_strict() -> None: + verdicts = _replay([(_NOW, ("9.9.9.9", 40000), ("10.0.0.1", 62201), _voice())]) + assert verdicts[0].outcome == "refused" + assert "source" in verdicts[0].reason + + +def test_a_frame_no_bridge_can_verify_is_left_unmatched() -> None: + config = _config() + config["SYSTEMS"]["OBP-FR"]["PASSPHRASE"] = (b"another" + b"\x00" * 20)[:20] + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice())], config) + assert verdicts[0].outcome == "unmatched" + assert verdicts[0].system is None + + +def test_a_v5_frame_is_replayed_with_the_time_it_was_captured() -> None: + config = _config() + config["SYSTEMS"]["OBP-FR"]["VER"] = 5 + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice_v5())], config) + assert verdicts[0].kind == "DMRE v5" + assert verdicts[0].outcome == "delivered" + + +def test_a_v5_frame_that_was_already_late_when_captured_is_dropped() -> None: + config = _config() + config["SYSTEMS"]["OBP-FR"]["VER"] = 5 + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice_v5(age=9.0))], config) + assert verdicts[0].reason == "stale-packet" + + +def test_a_v1_frame_on_a_v5_link_asks_the_peer_for_its_version() -> None: + config = _config() + config["SYSTEMS"]["OBP-FR"]["VER"] = 5 + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice())], config) + assert verdicts[0].reason == "proto-version" + + +def test_control_frames_are_named_not_judged() -> None: + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), build_bcka(_PASSPHRASE))]) + assert verdicts[0].kind == "BCKA" + assert verdicts[0].outcome == "control" + assert verdicts[0].datagram.payload[:4] == BCKA + + +def test_the_bridge_is_chosen_by_the_port_the_frame_arrived_on() -> None: + """Two bridges, one passphrase: the local port is the evidence that tells them apart.""" + config = _config() + config["SYSTEMS"]["OBP-PT"] = { + **config["SYSTEMS"]["OBP-FR"], + "TARGET_IP": _PT[0], + "TARGET_PORT": _PT[1], + "TARGET_SOCK": _PT, + "RELAX_CHECKS": True, + "_REPORT_PORT": 62268, + } + verdicts = _replay([(_NOW, ("203.0.113.9", 50000), ("10.0.0.1", 62268), _voice())], config) + assert verdicts[0].system == "OBP-PT" + + +def test_only_the_system_that_was_asked_for() -> None: + config = _config() + config["SYSTEMS"]["OBP-PT"] = {**config["SYSTEMS"]["OBP-FR"], "_REPORT_PORT": 62268} + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice())], config, system="OBP-PT") + assert verdicts[0].system == "OBP-PT" + with pytest.raises(KeyError): + _replay([], config, system="OBP-NOPE") + + +def test_the_report_tallies_what_happened() -> None: + verdicts = _replay( + [ + (_NOW, _FR, ("10.0.0.1", 62201), _voice(stream=1)), + (_NOW, _FR, ("10.0.0.1", 62201), _voice(dst=9, stream=2)), + (_NOW, _FR, ("10.0.0.1", 62201), _voice(dst=9990, stream=3)), + ] + ) + report = format_report(verdicts, capture="obp.pcap") + assert "3 datagram(s)" in report + assert "OBP-FR" in report + assert "delivered" in report + assert "dropped: tg-filter" in report + + +def test_the_report_can_skip_the_per_frame_lines() -> None: + verdicts = _replay([(_NOW, _FR, ("10.0.0.1", 62201), _voice())]) + summary = format_report(verdicts, capture="obp.pcap", verbose=False) + assert "delivered" in summary + assert time.strftime("%H:%M:%S", time.localtime(_NOW)) not in summary + + +# --- the command --------------------------------------------------------------- + + +def test_run_replay_prints_a_report(tmp_path, capsys) -> None: + capture = tmp_path / "obp.pcap" + write_pcap( + capture, + [ + (_NOW, _FR, ("10.0.0.1", 62201), _voice(stream=1)), + (_NOW + 1, _FR, ("10.0.0.1", 62201), _voice(dst=9, stream=2)), + ], + ) + assert run_replay(_config(), str(capture)) == 0 + printed = capsys.readouterr().out + assert "2 datagram(s)" in printed + assert "tg-filter" in printed + + +def test_run_replay_stops_at_the_limit(tmp_path, capsys) -> None: + capture = tmp_path / "obp.pcap" + write_pcap(capture, [(_NOW + i, _FR, ("10.0.0.1", 62201), _voice(stream=i)) for i in range(5)]) + assert run_replay(_config(), str(capture), limit=2) == 0 + assert "2 datagram(s)" in capsys.readouterr().out + + +def test_run_replay_reports_a_bad_capture(tmp_path, capsys) -> None: + path = tmp_path / "broken.pcap" + path.write_bytes(b"x" * 64) + assert run_replay(_config(), str(path)) == 1 + assert "ERROR capture" in capsys.readouterr().err + + +def test_run_replay_reports_an_unknown_system(tmp_path, capsys) -> None: + capture = tmp_path / "obp.pcap" + write_pcap(capture, [(_NOW, _FR, ("10.0.0.1", 62201), _voice())]) + assert run_replay(_config(), str(capture), system="NOPE") == 1 + assert "ERROR system" in capsys.readouterr().err