feat(obp): replay a capture through the ingress to see why frames were refused

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 <noreply@anthropic.com>
pull/81/head
yo 1 week ago
parent e105d493f0
commit 5c24ecb48b

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

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

@ -0,0 +1,339 @@
# ADN DMR Peer Server - infrastructure obp replay
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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",
]

@ -0,0 +1,180 @@
# ADN DMR Peer Server - infrastructure pcap reader
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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",
]

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

@ -0,0 +1,400 @@
# ADN DMR Peer Server - tests infrastructure obp replay
#
# Copyright (C) 2026 Rodrigo Pérez, CE5RPY <ce5rpy@qmd.cl>
#
###############################################################################
# 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("<IHHiIII", 0xA1B2C3D4, 2, 4, 0, 0, 262144, 1))
for when, source, destination, payload in frames:
packet = _udp_packet(source, destination, payload)
fh.write(struct.pack("<IIII", int(when), int(when % 1 * 1_000_000), len(packet), len(packet)))
fh.write(packet)
def _voice(*, src: int = 2130001, dst: int = 214, bits: int = 0x21, stream: int = 0x1234) -> 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("<IHHiIII", 0xA1B2C3D4, 2, 4, 0, 0, 262144, 276))
frame = sll2 + packet
fh.write(struct.pack("<IIII", int(_NOW), 0, len(frame), len(frame)))
fh.write(frame)
datagrams = list(read_udp(capture))
assert [d.payload for d in datagrams] == [b"hello"]
assert datagrams[0].source == _FR
def _write_raw(path, linktype: int, frame: bytes) -> None:
with open(path, "wb") as fh:
fh.write(struct.pack("<IHHiIII", 0xA1B2C3D4, 2, 4, 0, 0, 262144, linktype))
fh.write(struct.pack("<IIII", int(_NOW), 0, len(frame), len(frame)))
fh.write(frame)
def _ip_udp() -> 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("<I", 2) + _ip_udp())
assert [d.payload for d in read_udp(capture)] == [b"hello"]
def test_reading_a_vlan_tagged_frame(tmp_path) -> 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("<IIII", int(_NOW), 0, 200, 200) + b"\x45" * 10)
assert len(list(read_udp(capture))) == 1
# --- the replay ---------------------------------------------------------------
def _replay(frames: list[tuple], config: dict | None = None, **kwargs):
from adn_server.infrastructure.pcap import CapturedDatagram
datagrams = [
CapturedDatagram(timestamp=when, source=src, destination=dst, payload=payload)
for when, src, dst, payload in frames
]
return replay(config or _config(), datagrams, **kwargs)
def test_a_good_frame_from_the_configured_peer_is_delivered() -> 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
Loading…
Cancel
Save

Powered by TurnKey Linux.