Files
fips/testing/chaos/sim/runner.py
T
Johnathan Corgan 57a21b3ab8 Assert datagram delivery over Ethernet in the ethernet-mesh chaos scenario
No chaos scenario sent a datagram across an Ethernet link: both Ethernet
scenarios asserted only that the tree formed, so framing, AEAD over
AF_PACKET and delivery through the Ethernet transport were untested.

A new delivery assertion pings node pairs at each configured payload size
at teardown, after flapped links and stopped nodes are restored, one packet
at a time until min_replies replies arrive or deadline_secs passes. The
verdict is whether that happened, not a loss ratio, so a busy host slows the
probe without failing it. With require_transport and transport_node it also
reads that node's peers before and after the probe and fails if the node has
no peers or any peer on another transport, so traversal is checked rather
than assumed. A probe that never ran fails as a harness failure.

ethernet-mesh asserts delivery to and from n04, whose only edges are
Ethernet, at payloads 0 and 1200: n04 to n06 and back cross Ethernet then
UDP, and n04 to n05 is one Ethernet hop. Break-checks: dropping echo
requests on n06's TUN fails n04->n06 with zero replies while baseline
passes; adding a UDP edge to n04 fails the traversal check.

The payload sizes cannot exercise the receiver's padding trim on veth: an
empty echo is a 226-byte frame and a 1200-byte echo 1342 bytes, and the
veth does not pad shorter frames (54-byte control frames and 48-byte beacons
arrived unpadded; no received data frame carried bytes past its length
field). The trim is covered by the unit tests on the receive loop's parse.
2026-09-26 18:59:43 +00:00

1021 lines
42 KiB
Python

"""Main simulation orchestration."""
from __future__ import annotations
import json
import logging
import os
import random
import signal
import subprocess
import sys
import time
from datetime import datetime
from .assertions import (
AssertionOutcome,
BloomSendRateMonitor,
evaluate_baseline,
evaluate_congestion_signals,
evaluate_delivery,
evaluate_max_errors,
evaluate_max_parent_switches,
evaluate_min_parent_switches,
evaluate_tree_parents,
)
from .compose import generate_compose
from .config_gen import write_configs
from .control import (
query_peers,
query_status,
snapshot_all_congestion,
snapshot_all_mmp,
snapshot_all_trees,
)
from .docker_exec import docker_compose, existing_containers, force_remove
from .link_swap import LinkSwapManager
from .links import LinkManager
from .logs import AnalysisResult, analyze_logs, collect_logs, write_sim_metadata
from .naming import name_suffix
from .netclaim import claim_network, remove_network
from .netem import NetemManager
from .nodes import NodeManager
from .peer_churn import PeerChurnManager
from .scenario import Scenario
from .topology import SimTopology, generate_topology
from .traffic import TrafficManager
from .veth import VethManager
log = logging.getLogger(__name__)
# Upper bound on the wait for restored nodes' control sockets before the
# final snapshot. See _wait_control_sockets.
CONTROL_WAIT_SECS = 60
class SimRunner:
def __init__(self, scenario: Scenario):
self.scenario = scenario
self.rng = random.Random(scenario.seed)
self.topology: SimTopology | None = None
self.compose_file: str | None = None
# Claimed in _setup; the compose file refers to it as external, so
# `compose down` does not remove it and teardown must.
self.network_name: str | None = None
self.output_dir: str = self._resolve_output_dir(scenario)
self._interrupted = False
# Set when setup, warmup or the simulation loop raises. The exception
# is logged rather than propagated, so nothing else on the return path
# of run() carries the fact that the simulation did not complete.
self.aborted = False
# Whether this run ever owned containers. Teardown runs from a
# finally, so it runs after a failed setup too, and container names
# are global: without this it would harvest whatever currently
# answers to them. topology and compose_file are both set well
# before the containers start and so cannot stand in for it.
self._containers_started = False
# Shared set of currently-down node IDs (updated by NodeManager,
# read by NetemManager, LinkManager, TrafficManager)
self._down_nodes: set[str] = set()
# Managers (initialized during setup)
self.veth_mgr: VethManager | None = None
self.netem_mgr: NetemManager | None = None
self.link_mgr: LinkManager | None = None
self.link_swap_mgr: LinkSwapManager | None = None
self.traffic_mgr: TrafficManager | None = None
self.node_mgr: NodeManager | None = None
self.peer_churn_mgr: PeerChurnManager | None = None
# Post-run assertion monitors (sampled near end of run).
self.bloom_rate_monitor: BloomSendRateMonitor | None = None
self.assertion_outcomes: list[AssertionOutcome] = []
# Set by the final snapshot. None means the snapshot never ran,
# which the congestion assertion must treat as a failure rather
# than as an absence of congestion.
self.final_congestion: dict | None = None
self.final_tree: dict | None = None
# Set by _probe_delivery. None means the probe never ran, which the
# delivery assertion reports as a harness failure.
self.delivery_probe: dict | None = None
def _evaluate_max_parent_switches(
self, cfg, parent_switches: list[tuple[str, str]]
) -> AssertionOutcome:
"""Count parent switches in the configured scope and apply the ceiling.
A per-node scope is resolved through the topology and fails
explicitly if the node id is not in it. That is the whole reason
this lives here rather than in assertions.py: filtering log lines
by an unknown container name yields zero matches, and zero
trivially satisfies a ceiling, so a typo in ``node:`` would turn
the assertion into an unconditional pass.
Note the limit of that guard. It covers an unresolvable node id
and nothing else. A node that is in the topology but whose log is
empty, or whose log level suppressed the event, still counts zero
and still passes. Only the misspelling is caught here; the log
level is guarded at load time, and an empty log for a live node
is not guarded at all.
"""
if cfg.node is None:
return evaluate_max_parent_switches(
cfg, len(parent_switches), "mesh-wide"
)
if cfg.node not in self.topology.nodes:
known = ", ".join(sorted(self.topology.nodes))
return AssertionOutcome(
name="max_parent_switches",
passed=False,
detail=(
f"FAIL max_parent_switches: node '{cfg.node}' is not in "
f"this topology (nodes: {known}). Nothing was counted, so "
f"the ceiling was never actually tested."
),
)
source = self.topology.container_name(cfg.node)
count = sum(1 for src, _ in parent_switches if src == source)
return evaluate_max_parent_switches(cfg, count, f"node {cfg.node}")
@staticmethod
def _resolve_output_dir(scenario: Scenario) -> str:
"""Build a timestamped output directory path.
Format: {base}/{scenario_name}-{YYYYMMDD-HHMMSS}/
The base path is determined by (in priority order):
1. FIPS_SIM_OUTPUT environment variable
2. The scenario YAML's logging.output_dir (default: ./sim-results)
"""
base = os.environ.get("FIPS_SIM_OUTPUT", scenario.logging.output_dir)
timestamp = datetime.now().strftime("%Y%m%d-%H%M%S")
return os.path.join(base, f"{timestamp}-{scenario.name}")
def run(self) -> AnalysisResult | None:
"""Run the full simulation lifecycle."""
signal.signal(signal.SIGINT, self._handle_sigint)
signal.signal(signal.SIGTERM, self._handle_sigint)
result = None
try:
self._setup()
self._warmup()
self._simulation_loop()
except Exception:
# Log rather than re-raise: the default excepthook writes to stderr
# and would not reach the file handler installed in _setup(), so a
# re-raise would lose the traceback from runner.log, which is the
# artifact that survives into CI. The flag carries the abort out.
log.exception("Simulation failed")
self.aborted = True
finally:
result = self._teardown()
return result
def _handle_sigint(self, signum, frame):
if self._interrupted:
log.warning("Force exit")
sys.exit(1)
log.info("Interrupt received, shutting down gracefully...")
self._interrupted = True
def _setup(self):
"""Generate topology, configs, compose file. Start containers."""
s = self.scenario
mesh_name = f"sim-{s.name}-{s.seed}"
# Set up runner log file so all sim orchestration output is captured
# alongside the per-node FIPS logs for post-run analysis.
os.makedirs(self.output_dir, exist_ok=True)
runner_log_path = os.path.join(self.output_dir, "runner.log")
fh = logging.FileHandler(runner_log_path, mode="w")
fh.setLevel(logging.DEBUG)
fh.setFormatter(logging.Formatter(
"%(asctime)s %(levelname)-5s %(name)s: %(message)s",
datefmt="%H:%M:%S",
))
logging.getLogger().addHandler(fh)
log.info("Runner log: %s", runner_log_path)
# 0. Claim this run's network range, before anything derives from it.
#
# The ordering is not obvious and it matters: node IPs are computed
# from the subnet inside generate_topology, and traffic shaping keys
# its filters on those addresses. Claiming after the topology exists
# would give a network on one range and `tc` filters on another, which
# does not fail at bring-up — it silently leaves the shaping matching
# nothing.
self.network_name = f"fips-sim{name_suffix()}-net"
s.topology.subnet = claim_network(
self.network_name,
labels={
"com.corganlabs.fips-ci": "1",
"com.corganlabs.fips-ci.run": os.environ.get(
"FIPS_CI_RUN_ID", "manual"
),
},
candidates=(
[s.topology.pinned_subnet] if s.topology.pinned_subnet else None
),
)
# 1. Generate topology
log.info(
"Generating %d-node %s topology (seed=%d)...",
s.topology.num_nodes,
s.topology.algorithm,
s.seed,
)
self.topology = generate_topology(s.topology, self.rng, mesh_name)
log.info(
"Topology: %d nodes, %d edges",
len(self.topology.nodes),
len(self.topology.edges),
)
# Log adjacency summary
for nid in sorted(self.topology.nodes):
peers = sorted(self.topology.nodes[nid].peers)
log.info(" %s: peers=%s", nid, ",".join(peers))
# 2. Generate configs
#
# The directory carries the same suffix as the container names, so
# parallel scenarios cannot overwrite each other's compose file
# between writing it and starting containers from it.
docker_network_dir = os.path.join(os.path.dirname(__file__), "..")
config_dir = os.path.normpath(
os.path.join(
docker_network_dir,
"generated-configs",
f"sim{self.topology.name_suffix}",
)
)
# Select ephemeral identity nodes (if peer churn enabled)
self._ephemeral_nodes: set[str] = set()
if s.peer_churn.enabled and s.peer_churn.ephemeral_fraction > 0:
all_nodes = sorted(self.topology.nodes.keys())
count = int(len(all_nodes) * s.peer_churn.ephemeral_fraction)
self._ephemeral_nodes = set(self.rng.sample(all_nodes, count))
log.info(
"Ephemeral identity nodes (%d/%d): %s",
len(self._ephemeral_nodes),
len(all_nodes),
", ".join(sorted(self._ephemeral_nodes)),
)
write_configs(
self.topology, config_dir, self.scenario.fips_overrides,
ephemeral_nodes=self._ephemeral_nodes,
)
log.info("Wrote node configs to %s", config_dir)
# 3. Generate docker-compose.yml
self.compose_file = generate_compose(
self.topology, self.scenario, config_dir, self.network_name
)
log.info("Wrote %s", self.compose_file)
# 4. Obtain the test image (once, rather than per-service at scale).
#
# Building it is right for a bare run and wrong under a harness. When
# FIPS_TEST_IMAGE is set the image belongs to the caller, and every
# scenario of a parallel run would otherwise rebuild it into one shared
# name from one shared context — so assert it exists and fail loudly if
# it does not, rather than manufacture a substitute nobody asked for.
from .compose import FIPS_SIM_IMAGE
if os.environ.get("FIPS_TEST_IMAGE"):
log.info("Using caller-supplied image %s", FIPS_SIM_IMAGE)
probe = subprocess.run(
["docker", "image", "inspect", FIPS_SIM_IMAGE],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
)
if probe.returncode != 0:
raise RuntimeError(
f"FIPS_TEST_IMAGE names {FIPS_SIM_IMAGE}, which is not present. "
"The harness that set it is expected to have built it."
)
else:
log.info("Building Docker image...")
docker_dir = os.environ.get("FIPS_BUILD_CONTEXT") or os.path.join(
os.path.dirname(__file__), "..", "..", "docker"
)
subprocess.run(
["docker", "build", "-t", FIPS_SIM_IMAGE, docker_dir],
check=True,
)
# 5. Start containers
log.info("Starting %d containers...", len(self.topology.nodes))
docker_compose(self.compose_file, ["up", "-d"])
self._containers_started = True
# 6. Set up veth pairs for Ethernet edges (before netem)
#
# The entrypoint script waits for configured Ethernet interfaces
# to appear before starting FIPS, so we just need to create the
# veth pairs promptly after containers are running.
if self.topology.has_ethernet():
self.veth_mgr = VethManager(self.topology)
log.info("Setting up Ethernet veth pairs...")
self.veth_mgr.setup_all()
# 7. Initialize managers
if s.netem.enabled:
bw = s.bandwidth if s.bandwidth.enabled else None
ig = s.ingress if s.ingress.enabled else None
self.netem_mgr = NetemManager(self.topology, s.netem, self.rng, bandwidth=bw, ingress=ig)
self.netem_mgr.down_nodes = self._down_nodes
log.info("Applying initial per-link netem...")
self.netem_mgr.setup_initial()
if s.link_flaps.enabled:
self.link_mgr = LinkManager(
self.topology, s.link_flaps, self.rng, netem_mgr=self.netem_mgr
)
if s.link_swap.enabled:
if not self.netem_mgr:
raise RuntimeError(
"link_swap requires netem.enabled (depends on per-link tc state)"
)
self.link_swap_mgr = LinkSwapManager(
self.topology, s.link_swap, self.netem_mgr, self.rng,
)
if s.assertions.bloom_send_rate is not None:
self.bloom_rate_monitor = BloomSendRateMonitor(
self.topology, s.assertions.bloom_send_rate,
)
if s.traffic.enabled:
self.traffic_mgr = TrafficManager(
self.topology, s.traffic, self.rng, down_nodes=self._down_nodes
)
if s.node_churn.enabled:
self.node_mgr = NodeManager(
self.topology, s.node_churn, self.rng,
netem_mgr=self.netem_mgr, down_nodes=self._down_nodes,
veth_mgr=self.veth_mgr,
on_node_restart=self._handle_node_restart,
)
if s.peer_churn.enabled:
self.peer_churn_mgr = PeerChurnManager(
self.topology, s.peer_churn, self.rng,
down_nodes=self._down_nodes,
ephemeral_nodes=self._ephemeral_nodes,
)
# Share npub cache with traffic manager so iperf3 targets
# use runtime npubs (critical for ephemeral identity nodes).
if self.traffic_mgr:
self.traffic_mgr.npub_cache = self.peer_churn_mgr.npub_cache
def _warmup(self):
"""Wait for mesh convergence."""
n = len(self.topology.nodes)
wait = max(10, n) # Heuristic: ~1s per node, minimum 10s
log.info("Waiting %ds for mesh convergence...", wait)
self._sleep(wait)
self._take_snapshot("warmup")
# Populate npub cache after convergence (nodes must be running)
if self.peer_churn_mgr:
self.peer_churn_mgr.refresh_all_npubs()
# Initial link swap policy assignment (after warmup so the netem
# tc state is fully set up and the daemons have already
# discovered each other under the calm baseline).
if self.link_swap_mgr:
self.link_swap_mgr.setup_initial()
def _handle_node_restart(self, node_id: str):
"""Called after a node container is restarted.
For ephemeral identity nodes, waits briefly for the daemon to
start, then queries its new npub and updates the peer churn
manager's cache.
"""
if not self.peer_churn_mgr:
return
if node_id not in self.peer_churn_mgr.ephemeral_nodes:
return
# Brief delay for daemon startup before querying control socket
time.sleep(2)
new_npub = self.peer_churn_mgr.refresh_npub(node_id)
if new_npub:
log.info("Ephemeral node %s new identity: %s...%s", node_id, new_npub[:12], new_npub[-6:])
else:
log.warning("Failed to refresh npub for ephemeral node %s", node_id)
def _simulation_loop(self):
"""Main event loop driving stochastic behavior."""
start = time.time()
s = self.scenario
duration = s.duration_secs
log.info("Simulation running for %ds...", duration)
# Schedule first events
next_netem = self._schedule_next(start, s.netem.mutation.interval_secs) if self.netem_mgr else float("inf")
next_flap = self._schedule_next(start, s.link_flaps.interval_secs) if self.link_mgr else float("inf")
next_traffic = self._schedule_next(start, s.traffic.interval_secs) if self.traffic_mgr else float("inf")
next_churn = self._schedule_next(start, s.node_churn.interval_secs) if self.node_mgr else float("inf")
next_peer_churn = self._schedule_next(start, s.peer_churn.interval_secs) if self.peer_churn_mgr else float("inf")
# Bloom-send-rate assertion: sample at window_secs before end.
bloom_window_start_at = float("inf")
if self.bloom_rate_monitor is not None:
bloom_window_start_at = (
start + duration - s.assertions.bloom_send_rate.window_secs
)
bloom_window_started = False
while not self._interrupted:
now = time.time()
elapsed = now - start
if elapsed >= duration:
break
# Bloom-rate window-start sampling
if (
self.bloom_rate_monitor is not None
and not bloom_window_started
and now >= bloom_window_start_at
):
log.info(
"Sampling bloom-rate window start (last %ds of run)...",
s.assertions.bloom_send_rate.window_secs,
)
self.bloom_rate_monitor.sample_window_start()
bloom_window_started = True
# Deterministic link swap (before netem mutation so a
# mutation round can't clobber the swap mid-tick).
if self.link_swap_mgr:
self.link_swap_mgr.maybe_swap(now)
# Netem mutation
if self.netem_mgr and now >= next_netem:
self.netem_mgr.mutate()
next_netem = self._schedule_next(now, s.netem.mutation.interval_secs)
# Link flaps
if self.link_mgr:
if now >= next_flap:
self.link_mgr.maybe_flap()
next_flap = self._schedule_next(now, s.link_flaps.interval_secs)
self.link_mgr.restore_expired()
# Traffic generation
if self.traffic_mgr:
if now >= next_traffic:
self.traffic_mgr.maybe_spawn()
next_traffic = self._schedule_next(now, s.traffic.interval_secs)
self.traffic_mgr.cleanup_expired()
# Node churn
if self.node_mgr:
if now >= next_churn:
self.node_mgr.maybe_kill()
next_churn = self._schedule_next(now, s.node_churn.interval_secs)
self.node_mgr.restore_expired()
# Peer churn (topology mutation)
if self.peer_churn_mgr:
if now >= next_peer_churn:
self.peer_churn_mgr.maybe_churn()
next_peer_churn = self._schedule_next(now, s.peer_churn.interval_secs)
# Status line
down_links = self.link_mgr.down_count if self.link_mgr else 0
down_nodes = self.node_mgr.down_count if self.node_mgr else 0
active = self.traffic_mgr.active_count if self.traffic_mgr else 0
peer_churns = self.peer_churn_mgr.churn_count if self.peer_churn_mgr else 0
status_extra = f" peer_churns={peer_churns}" if self.peer_churn_mgr else ""
if self.link_swap_mgr:
status_extra += f" swaps={self.link_swap_mgr.swap_count}"
print(
f"\r [{elapsed:.0f}s/{duration}s] "
f"nodes={len(self.topology.nodes)} "
f"edges={len(self.topology.edges)} "
f"links_down={down_links} "
f"nodes_down={down_nodes} "
f"traffic={active}"
f"{status_extra} ",
end="",
flush=True,
)
self._sleep(1)
print() # Clear status line
# Bloom-rate window-end sampling.
# Done before teardown so containers are still running.
if self.bloom_rate_monitor is not None:
if not bloom_window_started:
# Loop exited before the window-start mark (e.g.,
# interrupted). Take both samples now so we still
# produce a finite outcome rather than an empty dict.
log.warning(
"Bloom-rate window did not start during run; "
"sampling both endpoints at end (delta will be 0)."
)
self.bloom_rate_monitor.sample_window_start()
log.info("Sampling bloom-rate window end...")
self.bloom_rate_monitor.sample_end()
def _evaluate_assertions(self) -> None:
"""Evaluate post-run assertions and stash outcomes on self.
Called from teardown while containers are still running so the
assertion outcomes (which include per-node detail) can be
written alongside the run artifacts.
"""
if self.bloom_rate_monitor is not None:
outcome = self.bloom_rate_monitor.evaluate()
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
@property
def assertions_failed(self) -> bool:
return any(not o.passed for o in self.assertion_outcomes)
def _run_status(self) -> str:
"""Name for how the run itself ended, before teardown is considered."""
if self.aborted:
return "aborted"
if self._interrupted:
return "interrupted"
return "completed"
def _write_status(self, status: str) -> None:
"""Record how the run ended, for a reader who has only this directory.
One field per line so it can be grepped. The container names are
included because they are the handle on everything else the run
touched, and because a directory that names them cannot be mistaken
for a directory assembled from some other scenario's mesh.
Never raises. Both callers run this immediately before stopping the
containers, so letting a full disk or a vanished output directory
out of here would leak the mesh it was about to tear down.
"""
try:
os.makedirs(self.output_dir, exist_ok=True)
path = os.path.join(self.output_dir, "status.txt")
with open(path, "w") as f:
f.write(f"status={status}\n")
f.write(f"scenario={self.scenario.name}\n")
f.write(f"seed={self.scenario.seed}\n")
f.write(f"containers=fips-node-nNN{name_suffix()}\n")
except Exception:
log.exception("Could not write status file")
def _own_containers(self) -> list[str]:
"""Name every container this run asked compose to create.
Empty before the topology exists, which is the only window in which
a teardown can run with nothing of its own on the host.
"""
if not self.topology:
return []
return [self.topology.container_name(nid) for nid in self.topology.nodes]
def _stop_mesh(self) -> None:
"""Take the containers down, then check that they actually went.
`docker compose down` exits 0 while leaving containers behind when it
races an `up -d` that failed part way through: compose enumerates the
project before the stragglers have registered, finds nothing to
remove, and says so successfully. Nothing downstream looks at what
survived, so the leak is invisible at the one moment it is cheap to
see.
"""
log.info("Stopping containers...")
docker_compose(self.compose_file, ["down"], check=False)
self._check_teardown()
def _check_teardown(self) -> None:
"""Report, then clear, any of this run's containers that outlived `down`.
Scoped to names this run generated. A check that reasoned about the
compose project label, or about container names in general, would
answer for whatever a concurrent scenario happens to own, and the CI
matrix runs plenty of those at once.
"""
wanted = self._own_containers()
if not wanted:
return
survivors = existing_containers(wanted)
if survivors is None:
# Detection failed, which is not the same as detecting nothing.
log.error("Could not ask docker what survived `down`; leak undetectable")
return
if not survivors:
return
log.warning(
"`compose down` exited 0 but left %d of this run's %d containers: %s",
len(survivors), len(wanted), ", ".join(survivors),
)
# Remedy, kept apart from the detection above on purpose: deleting
# everything from here down leaves the report intact.
force_remove(survivors)
leaked = existing_containers(survivors)
if not leaked:
return
# Either docker could not be asked a second time, or the containers
# are still there after a forced removal. Neither is the transient
# race, and neither clears itself, so this is the run's verdict.
detail = "unknown" if leaked is None else ", ".join(leaked)
log.error("Containers survived forced removal, leaked to the host: %s", detail)
self._write_leak(leaked if leaked is not None else wanted)
self.aborted = True
def _write_leak(self, names: list[str]) -> None:
"""Record the leaked names beside the run's other artifacts.
Never raises, for the same reason `_write_status` does not: this runs
on a teardown path that has already gone wrong. Written alongside
status.txt rather than into it: the status is the simulation's own
outcome, and a leak on the way out does not retract it.
"""
try:
os.makedirs(self.output_dir, exist_ok=True)
path = os.path.join(self.output_dir, "leaked-containers.txt")
with open(path, "w") as f:
for name in names:
f.write(name + "\n")
except Exception:
log.exception("Could not write leaked-containers file")
def _release_network(self) -> None:
"""Give this run's claimed /24 back.
The compose file declares the network `external:`, so `compose down`
leaves it alone and nothing else reclaims the range. Called from both
teardown paths, including the setup-failed one, because a run that
fell over after claiming still holds a range.
"""
if self.network_name:
remove_network(self.network_name)
def _teardown(self) -> AnalysisResult | None:
"""Stop dynamic elements, collect logs, analyze, stop containers."""
if not self._containers_started:
# There is no mesh to harvest. Snapshots, logs, analysis and
# assertions all address containers by a name that is global to
# the host, so running them here would describe whoever holds
# those names now and leave a plausible report of a run that
# never happened. Leaving them out makes the presence of
# analysis.txt proof that this scenario's own mesh existed.
self._write_status("setup-failed")
if self.compose_file:
# `up -d` can fail part way through, so this run may own
# containers or a network even with no mesh to speak of.
self._stop_mesh()
self._release_network()
return None
result = None
status = self._run_status()
try:
# Evaluate post-run assertions before doing any teardown so
# control sockets are still reachable.
self._evaluate_assertions()
# Stop traffic
if self.traffic_mgr:
log.info("Stopping traffic sessions...")
self.traffic_mgr.stop_all()
# Restore links
if self.link_mgr:
log.info("Restoring downed links...")
self.link_mgr.restore_all()
# Restore stopped nodes (needed for snapshots and log collection)
if self.node_mgr:
log.info("Restoring stopped nodes...")
self.node_mgr.restore_all()
# A node restarted just above can take several seconds to open
# its control socket, and a snapshot taken before then records
# it as absent. Wait for every node to answer first.
self._wait_control_sockets(CONTROL_WAIT_SECS)
# Every flapped link and stopped node is back, so the delivery
# probe measures the restored mesh.
if self.scenario.assertions.delivery is not None:
try:
self._probe_delivery()
except Exception:
# Leaves delivery_probe None, which the assertion
# reports as a harness failure; the rest of teardown
# still has to run.
log.exception("Delivery probe failed")
# Collect iperf3 throughput results before containers stop
if self.traffic_mgr:
iperf_results = self.traffic_mgr.collect_results()
if iperf_results:
iperf_path = os.path.join(self.output_dir, "iperf3-results.json")
with open(iperf_path, "w") as f:
json.dump(iperf_results, f, indent=2)
log.info("Saved %d iperf3 results to %s", len(iperf_results), iperf_path)
# Take final tree snapshot while nodes are still running
self._take_snapshot("final")
# Collect logs before stopping containers
container_names = [
self.topology.container_name(nid) for nid in sorted(self.topology.nodes)
]
log.info("Collecting logs from %d containers...", len(container_names))
logs = collect_logs(container_names, self.output_dir)
# Analyze
result = analyze_logs(logs)
analysis_path = os.path.join(self.output_dir, "analysis.txt")
with open(analysis_path, "w") as f:
f.write(result.summary())
print(result.summary())
# Log-derived assertions (evaluated after analyze_logs so
# parent_switches and similar are populated). See
# _evaluate_max_parent_switches for why the per-node variant
# resolves its node against the topology rather than filtering
# optimistically.
mps_cfg = self.scenario.assertions.min_parent_switches
if mps_cfg is not None:
outcome = evaluate_min_parent_switches(
mps_cfg, len(result.parent_switches)
)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
xps_cfg = self.scenario.assertions.max_parent_switches
if xps_cfg is not None:
outcome = self._evaluate_max_parent_switches(
xps_cfg, result.parent_switches
)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
# Applied to every scenario by default, so this is the one
# assertion that is present even when the YAML declares no
# assertions block at all.
err_cfg = self.scenario.assertions.max_errors
if err_cfg is not None:
outcome = evaluate_max_errors(err_cfg, result.errors)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
bl_cfg = self.scenario.assertions.baseline
if bl_cfg is not None:
outcome = evaluate_baseline(
bl_cfg, self.final_tree, len(result.sessions_established)
)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
tp_cfg = self.scenario.assertions.tree_parents
if tp_cfg is not None:
outcome = evaluate_tree_parents(tp_cfg, self.final_tree)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
dl_cfg = self.scenario.assertions.delivery
if dl_cfg is not None:
outcome = evaluate_delivery(dl_cfg, self.delivery_probe)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
cong_cfg = self.scenario.assertions.congestion_signals
if cong_cfg is not None:
outcome = evaluate_congestion_signals(
cong_cfg, self.final_congestion
)
self.assertion_outcomes.append(outcome)
if outcome.passed:
log.info("%s", outcome.detail)
else:
log.error("%s", outcome.detail)
# Write assertion outcomes
if self.assertion_outcomes:
assertions_path = os.path.join(self.output_dir, "assertions.txt")
with open(assertions_path, "w") as f:
for o in self.assertion_outcomes:
f.write(o.detail + "\n")
print("=== Assertions ===")
for o in self.assertion_outcomes:
print(o.detail)
print()
# Write metadata
write_sim_metadata(
self.output_dir,
scenario_name=self.scenario.name,
seed=self.scenario.seed,
num_nodes=len(self.topology.nodes),
num_edges=len(self.topology.edges),
duration_secs=self.scenario.duration_secs,
topology=self.topology,
)
# Clean up veth pairs
if self.veth_mgr:
log.info("Cleaning up veth pairs...")
self.veth_mgr.teardown_all()
except Exception:
# Same reasoning as run(): logging keeps the traceback in
# runner.log, where re-raising would send it to stderr instead
# and past the exit ladder, which would report a harvest that
# fell over as a bad command line.
log.exception("Teardown failed")
self.aborted = True
status = "teardown-failed"
result = None
finally:
# Status first: it is the one artifact that must exist whatever
# else happens, and stopping containers can still time out.
self._write_status(status)
self._stop_mesh()
self._release_network()
return result
def _wait_control_sockets(self, timeout: float) -> None:
"""Wait until every node's control socket answers, up to timeout.
Bounded, and it never fails the run itself: a node that still does
not answer is left to the final snapshot, where the assertions see
it as absent.
"""
start = time.monotonic()
waiting = sorted(self.topology.nodes)
while True:
waiting = [
nid for nid in waiting
if query_status(self.topology.container_name(nid)) is None
]
elapsed = time.monotonic() - start
if not waiting or elapsed >= timeout:
break
time.sleep(2)
if waiting:
log.warning(
"Control socket wait: %s still not answering after %.1fs",
", ".join(waiting), elapsed,
)
else:
log.info("Control socket wait: all nodes answered after %.1fs", elapsed)
def _read_peers(self, node_id: str, when: str, wait_secs: float) -> dict:
"""Read a node's peers, waiting up to wait_secs for at least one.
A node restored at teardown may not have re-peered yet, which is
not the failure the transport check is for, so an empty list is
re-read until the wait runs out. A failed read is retried the same
way and reported as None if it never succeeds.
"""
container = self.topology.container_name(node_id)
start = time.monotonic()
while True:
data = query_peers(container)
peers = None if data is None else data.get("peers", [])
waited = time.monotonic() - start
if peers or waited >= wait_secs:
return {"when": when, "peers": peers, "waited_s": waited}
time.sleep(2)
def _ping_until(self, src: str, dst: str, size: int, cfg) -> dict:
"""Ping dst from src until cfg.min_replies replies or the deadline."""
container = self.topology.container_name(src)
target = f"{self.topology.nodes[dst].npub}.fips"
cmd = ["docker", "exec", container, "ping6", "-c", "1", "-W", "2",
"-s", str(size), target]
start = time.monotonic()
replies = attempts = 0
last = ""
while replies < cfg.min_replies and time.monotonic() - start < cfg.deadline_secs:
attempts += 1
try:
res = subprocess.run(cmd, capture_output=True, text=True, timeout=15)
lines = (res.stdout + res.stderr).strip().splitlines()
last = lines[-1] if lines else ""
if res.returncode == 0:
replies += 1
continue
except subprocess.TimeoutExpired:
last = "docker exec ping6 timed out"
time.sleep(1)
return {
"src": src, "dst": dst, "size": size, "replies": replies,
"attempts": attempts, "elapsed_s": time.monotonic() - start,
"last_output": last,
}
def _probe_delivery(self):
"""Run the delivery assertion's probes and keep the result."""
cfg = self.scenario.assertions.delivery
log.info("Probing delivery: %d pair(s) x %d size(s)",
len(cfg.pairs), len(cfg.payload_bytes))
transport = None
if cfg.require_transport:
transport = {
"node": cfg.transport_node, "want": cfg.require_transport,
"reads": [self._read_peers(cfg.transport_node, "before",
cfg.deadline_secs)],
}
results = [
self._ping_until(src, dst, size, cfg)
for src, dst in cfg.pairs
for size in cfg.payload_bytes
]
if transport is not None:
transport["reads"].append(
self._read_peers(cfg.transport_node, "after", 0)
)
self.delivery_probe = {"transport": transport, "results": results}
with open(os.path.join(self.output_dir, "delivery-probe.json"), "w") as f:
json.dump(self.delivery_probe, f, indent=2)
def _take_snapshot(self, label: str):
"""Query all nodes via control socket and save tree/MMP/congestion snapshots."""
if not self.topology:
return
log.info("Taking %s snapshot...", label)
tree_snap = snapshot_all_trees(self.topology)
mmp_snap = snapshot_all_mmp(self.topology)
congestion_snap = snapshot_all_congestion(self.topology)
tree_path = os.path.join(self.output_dir, f"tree-snapshot-{label}.json")
mmp_path = os.path.join(self.output_dir, f"mmp-snapshot-{label}.json")
congestion_path = os.path.join(self.output_dir, f"congestion-snapshot-{label}.json")
# Retained for the congestion assertion, which needs the same
# responses the file gets. Reading the file back would let a write
# failure present as an assertion that saw nothing.
if label == "final":
self.final_congestion = congestion_snap
self.final_tree = tree_snap
os.makedirs(self.output_dir, exist_ok=True)
with open(tree_path, "w") as f:
json.dump(tree_snap, f, indent=2)
with open(mmp_path, "w") as f:
json.dump(mmp_snap, f, indent=2)
with open(congestion_path, "w") as f:
json.dump(congestion_snap, f, indent=2)
log.info(
"Snapshot %s: %d/%d tree, %d/%d mmp, %d/%d congestion responses",
label,
len(tree_snap),
len(self.topology.nodes),
len(mmp_snap),
len(self.topology.nodes),
len(congestion_snap),
len(self.topology.nodes),
)
def _schedule_next(self, now: float, interval) -> float:
"""Schedule the next event using a Range interval."""
return now + self.rng.uniform(interval.min, interval.max)
def _sleep(self, seconds: float):
"""Sleep in small increments so SIGINT can break out."""
end = time.time() + seconds
while time.time() < end and not self._interrupted:
time.sleep(min(0.5, end - time.time()))