mirror of
https://github.com/jmcorgan/fips.git
synced 2026-10-05 11:08:25 +00:00
Two daemons whose only transports are interface-bound, run against a veth pair the harness creates, downs, deletes and recreates underneath them. Asserts the boot race (a daemon whose only interface is missing starts, reports the transport absent and the node Degraded, rather than exiting on NoTransports or skipping the transport for the life of the process), the late attach and discovery over it, the flap in both directions, destroy-and-recreate, and that an optional interface which never appears never moves node health. Also the log policy, which is the half that is easy to regress silently: absence is logged once on the edge and not once per retry; a required interface still absent past the ten-second bring-up window errors exactly once, while the optional one — absent just as long — stays silent; and that error is not repeated on a schedule. The detach edge is checked not to error, guarded by how long detection actually took, so a slow runner skips the check rather than failing on the harness's own latency. The containers run FIPS_TEST_MODE=default, not chaos. The chaos entrypoint waits up to 30 s for every configured Ethernet interface before starting the daemon, which is precisely the workaround under test — the daemon has to do its own waiting here or the suite proves nothing. Host-namespace ip(8) runs in a short-lived privileged container sharing the host network and PID namespaces, for the reason chaos/sim/veth.py documents: on macOS the containers live in the Docker VM, so ip(8) run on the macOS host could never reach them. Chaos ethernet transports are marked optional: true. In that harness a neighbour's interface disappearing is the scenario, not a fault — node_churn stops a container, which destroys its netns and with it both ends of every veth it held, so a surviving node watches a required interface vanish for the 30-90 s the neighbour is down, once per churn event. Reporting that at error is right for a deployment and wrong for a harness that tears the interface down on purpose; the mesh-wide zero-ERROR ceiling would have failed on injected chaos rather than on a defect. test(iface-binding): cover an interface present before the daemon starts Every scenario in the suite created its interface after the daemons were already running — that ordering is the boot race the suite was written for. But it means both nodes could only ever reach Present through binder_loop, so the inline bind in start_async, which is the ordinary case on a booted router, had no end-to-end coverage at all. That is where the churn guard went unseeded and the first detach stopped reaching node health, and no existing case could reach it: they all detach from a binding the loop created, which seeds the guard as a side effect. Case (f) adds a third node whose single required interface exists before its daemon does. The gate is what buys that ordering — the harness needs a running container to have a netns to move a veth into, but the daemon must not start until after the move, so node-c comes up parked on a file and the harness releases it once the interface is in place. Then one detach, on a binding the loop did not create, and the node must degrade. Verified against the defect rather than only against the fix: with the guard seed reverted, cases (a) through (e) all still pass and (f) is the only failure. A regression test that has never been seen to fail is a claim, not a test. It also asserts the reverse edge, so Degraded stays a level rather than a latch on this path too. test(iface-binding): assert the fast path and the churn guard Three gaps, two of them in tests that existed and asserted nothing. **The netlink path was never asserted to be in use.** The 1 s poll is a complete fallback and covers every wait in the suite, so the whole thing passed with `open_link_socket()` hardcoded to Err — the fast path could have been dead for a release and no test would have said so. The binder reports which backing it got at startup, so case (g) asks it directly rather than inferring from timing the poll would also satisfy, and the unit test that used to write `let _ = w.is_event_driven();` now asserts it on Linux, where the source is an unprivileged `AF_NETLINK` socket and falling back is a real loss rather than a sandbox's prerogative. **Churn damping had no end-to-end coverage**, which now matters twice over: it bounds the recovery announcements, and since the detach edge withdraws peers it is also the only thing bounding how often that withdrawal fires. Every flap elsewhere in the suite is a single down/up with long settles either side — exactly the shape the damper ignores. Case (h) drives four bindings that each die inside `MIN_STABLE_BINDING`, asserts the guard engages, asserts it then *suppresses* rather than merely counting, and asserts it is not a latch. **`a_poisoned_binding_does_not_strand_the_transport` discarded its result.** `let _ = eth.binding.tasks_alive();` left the entire point unasserted: reading a poisoned lock as "alive" would have the binder believe a dead binding healthy and never rebind, and treating it as an error would strand the transport. `false` is what routes it back through detach and rebind, so say so. `a_stop_racing_a_bind_leaves_nothing_behind` now asserts the error *kind*. `bind_and_spawn` refuses at its presence probe long before the post-store shutdown check, so `is_err()` alone passed on absence and would still pass with that check deleted. The test keeps the coverage it genuinely has — stop raises the flag before teardown, teardown leaves no socket and no loops — and says plainly that the race it is named for needs a bind that succeeds, which needs privilege no unit test has. Both new cases were verified against the defect: with the netlink source forced to Err, (g) fails; with `CHURN_THRESHOLD` raised out of reach, (h) fails. Nothing else in the suite notices either. One case was attempted and removed rather than shipped: `"interface replaced"` cannot be produced deterministically, because the delete that changes an ifindex fires a netlink event the binder acts on within microseconds, so `gone` wins the race. It passed about one run in three. reference/notes.md records the measurement and the two approaches that could work. Also fixes a real bug in the harness: `grep -q` under `set -o pipefail` exits on its first match, `docker logs` takes SIGPIPE, and the pipeline reports failure even though the line was found. That cost two false failures before it was spotted; `log_count` reads the stream to the end. test(chaos): cover an Ethernet rebind under active traffic The one case dynamic interface binding had no coverage for anywhere: a datagram crossing an Ethernet link while the interface underneath it goes away and comes back. No existing scenario could reach it, for two separate reasons. `ethernet-only` and `ethernet-mesh` both run with `traffic.enabled: false`, so no datagram crosses an Ethernet link in any test — `ethernet-only`'s own comment says exactly that, and names framing, the length field that trims NIC minimum-frame padding, and AEAD over Ethernet as unexercised because of it. And `link_flaps` cannot produce a rebind whatever it is pointed at: it simulates a down link with netem 100% loss, so the interface stays IFF_UP and the presence machine never sees an edge. `ethernet-mesh` has had link flaps enabled all along without once exercising a rebind. `node_churn` is what actually moves an interface. Stopping a container destroys its network namespace, deleting every veth in it — and deleting one end of a veth deletes its peer — so a *surviving* node watches its Ethernet interface disappear outright, and watches it return when the harness recreates the pair on restart. That is a real detach and a real rebind, driven from outside the daemon. The new scenario is a 4-node Ethernet ring with traffic on and one node churned at a time, with link flaps deliberately off so the only outage is a genuine interface removal and a traffic shortfall cannot be ambiguous between the two. Measured across four runs: 206-388 MB moved over Ethernet links while interfaces were being taken away underneath. It also needed an assertion that did not exist. Traffic results have always been written to `iperf3-results.json` and never read, so a scenario carrying `traffic.enabled: true` could have every session fail and still exit 0 on a green control plane — and a rebind under load is precisely what a tree snapshot cannot see. `min_traffic` counts sessions that finished with bytes actually received, treating iperf3's top-level `error` and a missing `end` block as zero, so a session only counts when it moved data. The baseline is calibrated against four runs rather than assumed: `max_roots` starts at the observed maximum plus one, and the site records the sample, its size, and why four runs is thin. The first draft asserted a single root and failed every run — the harness restores stopped nodes immediately before the final snapshot, so a just-restarted node has not re-parented yet and is briefly its own root. That is the scenario working. Wired into both runners, since a chaos scenario on one side only makes "local green" and "GitHub green" stop meaning the same thing; check-ci-parity was confirmed to fail on a one-sided addition before this was committed. The iface-binding suite's entry in the GitHub workflow's integration matrix moves here from the commit that introduced the presence machine. That commit declared the suite on GitHub before testing/iface-binding/ existed and before testing/ci-local.sh knew about it, so testing/check-ci-parity.sh failed there and the three workflow steps named files that were not yet in the tree. Registering both runners in the commit that adds the suite settles both.
895 lines
37 KiB
Python
895 lines
37 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_max_errors,
|
|
evaluate_max_parent_switches,
|
|
evaluate_min_parent_switches,
|
|
evaluate_min_traffic,
|
|
evaluate_tree_parents,
|
|
)
|
|
from .compose import generate_compose
|
|
from .config_gen import write_configs
|
|
from .control import 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__)
|
|
|
|
|
|
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
|
|
|
|
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()
|
|
|
|
# Collect iperf3 throughput results before containers stop
|
|
iperf_results: list[dict] = []
|
|
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)
|
|
|
|
# Traffic. Evaluated even when no session completed, because
|
|
# "nothing ran" is the failure this exists to catch.
|
|
traffic_cfg = self.scenario.assertions.min_traffic
|
|
if traffic_cfg is not None:
|
|
outcome = evaluate_min_traffic(traffic_cfg, iperf_results)
|
|
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)
|
|
|
|
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 _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()))
|