Files
fips/testing/chaos/sim/netem.py
T
Arjen dca938c104 feat(peer): switchover between a peer's transports, with no handshake
A peer reachable over more than one transport keeps one Noise session
and moves its traffic between transports on failure or degradation.
Implements docs/design/fips-multi-path-switchover.md §4-§10 and closes
the two-interface case in #143.

Three inner link messages next to Heartbeat: `0x52 PathProbe` and
`0x53 PathAck`, carrying a probe id, the sender's path id and a
`remote_active` bit, and `0x54 PathClose`, naming the receiver's path
id and a reason. A probe is an ordinary encrypted frame sent on a
candidate transport; the receiver, having decrypted it against the
session found by index, adds the path as `Probing`, marks it `rx_live`
and answers on that same path. The prober's receipt of the ack marks
the path `Live`, `tx_live`, and takes an RTT sample. No handshake, no
key material, no index allocation. Old nodes drop the unknown types at
debug, so a path to one stays `Probing` and never becomes eligible.

The discovery gate changes shape: a live peer beaconing on a transport
we hold no path to it over becomes a path candidate rather than being
skipped, and the heartbeat tick probes it. The active path's first
probe is small — the handshake proved it and seeded its MTU — while a
standby's discovery probes are padded to the link MTU, as is one a
minute on every path, so a medium that passes small frames and drops
large ones never proves itself. A standby the peer never answers on is
given up after eight probes; the active path never is.

Detection is per path and takes each medium's own failure signal: a
carrier edge, an unreachable-on-send (`ENETUNREACH`/`EHOSTUNREACH`), an
interface going away, or two unanswered heartbeats on a path the peer
is also silent on. Any of them marks the path `Suspect` and selection
leaves it at once, because the standby is warm: heartbeats run at
`node.path.active_heartbeat_ms` where either side sends and
`standby_heartbeat_ms` elsewhere, both stretched by the path's own
round trip so a Tor or Nym path is neither flooded nor declared dead
every round trip. A peer holding one live path is not heartbeated here
at all — selection has nothing to move to, and the link heartbeat keeps
its liveness. Soft signals (the peer's `remote_active` flipping away,
silence here while a standby hears the peer) trigger a probe, never
`Suspect`: reading them as a verdict forces both sides onto one path
and loops under a one-way failure.

A node that loses a path tells the peer with a `PathClose` on a
surviving one, so the peer moves at once rather than after its own
timeout. A transport that returns inside the five-minute grace revives
its dead paths as `Probing` with their RTT window and ETX intact.

Selection is measured, not configured: each path scores
`quality_index(etx, min_rtt)`, and traffic moves when the active path
is no longer eligible (mandatory) or when a standby beats it by
`switch_margin` for `switch_dwell_secs` (discretionary). Min RTT over a
window rather than SRTT, because SRTT inflates under load on the path
carrying traffic while an idle standby looks pristine — a ping-pong
generator. A `role: backup` transport carries a peer's traffic only
while no normal path is eligible, and yields outright when one becomes
eligible. `fipsctl path pin` overrides both while its path is eligible.

A switch re-seeds the path MTU from the new path, tightens the
session MTUs and refreshes the MSS ceiling, so the first frames after a
switch are not black-holed. The link record and `addr_to_link` follow
the active path. The link cost the tree sees is held at its pre-switch
value for the dwell, and until the two receiver reports that span the
switch have arrived — the first counts every frame in flight on the old
path as lost and spikes the per-report ETX for one interval, the second
replaces it — so neither a short flap nor that spike ripples mesh-wide
through parent selection or the next-hop order.

Operator surface: `role: backup` on any transport, `node.path.*`
(`switch_margin`, validated finite and at least 1.0, `switch_dwell_secs`,
`min_samples`, `active_heartbeat_ms`, `standby_heartbeat_ms`), and
`fipsctl path show|pin|unpin` over the `path_show`, `path_pin` and
`path_unpin` control commands.

Two chaos scenarios calibrate the defaults and are wired into both
runners: `dual-path-flap` (a raw-Ethernet veth as the cable, the Docker
bridge over UDP as the wifi) and `dual-udp-flap` (two interface-bound
UDP instances). Each flaps one path under iperf and carries detectors
that can fail on a switchover that did not carry traffic — a per-node
ceiling on "Peer promoted to active" (a second is a re-peering), a
one-second ceiling from link-down to the first switch, and a two-second
ceiling on any zero-byte iperf interval run — alongside the
`path_switches` band. Neither has been run to calibrate; the defaults
are chosen, not derived, and the design doc says so.

Refs #143
2026-09-23 11:48:33 -03:00

551 lines
24 KiB
Python

"""Per-link network impairment via tc HTB + netem + u32 filters.
Each container gets an HTB root qdisc on eth0 with one class per peer.
Each class has a netem leaf qdisc for that specific link's impairment.
u32 filters match destination IP to direct traffic to the right class.
eth0 root (HTB 1:)
├── class 1:1 → netem 11: (peer 1) ← u32 filter: dst=<peer1_ip>
├── class 1:2 → netem 12: (peer 2) ← u32 filter: dst=<peer2_ip>
└── class 1:99 → pfifo (default, no impairment)
Optional ingress policing adds rate-based packet dropping on the
receive side via tc ingress qdisc + policer filters:
eth0 ingress (ffff:)
├── u32 filter: src=<peer1_ip> → police rate <R>kbit burst <B> drop
└── u32 filter: src=<peer2_ip> → police rate <R>kbit burst <B> drop
"""
from __future__ import annotations
import logging
import random
from dataclasses import dataclass, field
from .docker_exec import docker_exec_quiet, is_container_running
from .scenario import BandwidthConfig, IngressConfig, LinkPolicyOverride, NetemConfig, NetemPolicy
from .topology import SimTopology, veth_interface_name
log = logging.getLogger(__name__)
IFACE = "eth0"
@dataclass
class NetemParams:
"""Concrete netem parameters for one link direction."""
delay_ms: int = 0
jitter_ms: int = 0
loss_pct: float = 0.0
duplicate_pct: float = 0.0
reorder_pct: float = 0.0
corrupt_pct: float = 0.0
def to_tc_args(self) -> str:
"""Build the netem arguments string for tc."""
parts = []
if self.delay_ms > 0:
if self.jitter_ms > 0:
parts.append(f"delay {self.delay_ms}ms {self.jitter_ms}ms")
else:
parts.append(f"delay {self.delay_ms}ms")
if self.loss_pct > 0:
parts.append(f"loss {self.loss_pct:.1f}%")
if self.duplicate_pct > 0:
parts.append(f"duplicate {self.duplicate_pct:.1f}%")
if self.reorder_pct > 0 and self.delay_ms > 0:
parts.append(f"reorder {self.reorder_pct:.1f}%")
if self.corrupt_pct > 0:
parts.append(f"corrupt {self.corrupt_pct:.1f}%")
return " ".join(parts) if parts else "delay 0ms"
@dataclass
class LinkNetemState:
"""Tracks the netem state for a single link direction (one container's view)."""
container: str
dest_ip: str
class_id: str # e.g., "1:1"
netem_handle: str # e.g., "11:"
params: NetemParams = field(default_factory=NetemParams)
rate_mbit: int = 0 # 0 = unlimited (1gbit default)
ingress_rate_kbps: int = 0 # 0 = no ingress policing
ingress_burst_bytes: int = 32000
ingress_filter_prio: int = 0 # u32 filter priority for this peer
# Set while a link flap holds this direction down. `params` keeps the
# link's normal impairment, so every path that re-installs the qdisc asks
# tc_args() and a restart or mutation cannot end a flap early.
held_down: bool = False
def tc_args(self) -> str:
"""Return the netem arguments this direction should carry now."""
return "loss 100%" if self.held_down else self.params.to_tc_args()
@dataclass
class VethNetemState:
"""Tracks the netem state for a veth interface (Ethernet link direction)."""
container: str
iface: str # e.g., "ve-n01-n02"
params: NetemParams = field(default_factory=NetemParams)
# See LinkNetemState.held_down.
held_down: bool = False
def tc_args(self) -> str:
"""Return the netem arguments this direction should carry now."""
return "loss 100%" if self.held_down else self.params.to_tc_args()
class NetemManager:
"""Manages per-link netem impairment across all containers."""
def __init__(
self,
topology: SimTopology,
config: NetemConfig,
rng: random.Random,
bandwidth: BandwidthConfig | None = None,
ingress: IngressConfig | None = None,
):
self.topology = topology
self.config = config
self.rng = rng
# Per-container, per-dest-ip netem state (UDP links on eth0)
self.states: dict[str, dict[str, LinkNetemState]] = {}
# Per-container, per-veth netem state (Ethernet links)
self.veth_states: dict[str, dict[str, VethNetemState]] = {}
# Nodes currently down (updated by NodeManager) — skip tc ops on these
self.down_nodes: set[str] = set()
# Per-edge bandwidth: (node_a, node_b) -> rate in mbit
self._edge_rates: dict[tuple[str, str], int] = {}
if bandwidth and bandwidth.enabled:
for a, b in topology.edges:
rate = rng.choice(bandwidth.tiers_mbps)
self._edge_rates[(a, b)] = rate
self._edge_rates[(b, a)] = rate
# Per-edge ingress policing: (node_a, node_b) -> rate in kbps
self._ingress_config = ingress
self._ingress_rates: dict[tuple[str, str], int] = {}
if ingress and ingress.enabled:
for a, b in topology.edges:
rate = rng.choice(ingress.tiers_kbps)
self._ingress_rates[(a, b)] = rate
self._ingress_rates[(b, a)] = rate
# Per-edge policy overrides: canonical "nXX-nYY" -> NetemPolicy
# Build a set of canonical edge strings for validation
topo_edge_strs = {"-".join(sorted([a, b])) for a, b in topology.edges}
self._edge_overrides: dict[str, NetemPolicy] = {}
for override in config.link_policies:
policy = override.policy
if policy is None and override.policy_name:
policy = config.mutation.policies.get(override.policy_name)
if policy is None:
continue
for edge_str in override.edges:
if edge_str not in topo_edge_strs:
log.warning(
"link_policy edge %s not in topology — override ignored",
edge_str,
)
self._edge_overrides[edge_str] = policy
if self._edge_overrides:
log.info(
"Per-link policy overrides: %d edges",
len(self._edge_overrides),
)
# Edges the periodic mutation must never touch (canonical "nXX-nYY").
# Unlike a link_policy typo, which merely fails to shape a link, a typo
# here would silently leave an asserted link unprotected and reintroduce
# the exact flakiness the exclusion exists to remove — so an unknown
# edge is a hard error, not a warning.
self._mutation_exclude: set[str] = set()
for edge_str in config.mutation.exclude_edges:
canonical = "-".join(sorted(edge_str.split("-")))
if canonical not in topo_edge_strs:
raise ValueError(
f"netem.mutation.exclude_edges references {edge_str!r}, "
"which is not an edge in the topology"
)
self._mutation_exclude.add(canonical)
if self._mutation_exclude:
log.info(
"Mutation excludes %d edge(s): %s",
len(self._mutation_exclude),
", ".join(sorted(self._mutation_exclude)),
)
def _htb_rate(self, node_id: str, peer_id: str) -> str:
"""Return the HTB rate string for a link direction."""
rate = self._edge_rates.get((node_id, peer_id), 0)
return f"{rate}mbit" if rate > 0 else "1gbit"
def _policy_for_edge(self, node_a: str, node_b: str) -> NetemPolicy:
"""Return the netem policy for an edge, checking overrides first."""
canonical = "-".join(sorted([node_a, node_b]))
if canonical in self._edge_overrides:
return self._edge_overrides[canonical]
return self.config.default_policy
def setup_initial(self):
"""Set up HTB qdiscs and initial netem on all containers.
UDP peers use HTB + u32 filters on eth0. Ethernet peers use a
simple root netem qdisc on their dedicated veth interface.
"""
if self._edge_rates:
log.info("Bandwidth pacing enabled (%d edges with rate limits)",
len(self._edge_rates) // 2)
for node_id in sorted(self.topology.nodes):
node = self.topology.nodes[node_id]
container = self.topology.container_name(node_id)
# Split peers by transport type: IP-based (UDP/TCP) vs Ethernet (veth)
ip_peers = {}
eth_peers = []
for peer_id in sorted(node.peers):
transport = self.topology.transport_for_edge(node_id, peer_id)
if self.topology.is_veth_transport(transport):
eth_peers.append(peer_id)
# A dual edge also has a UDP half over the bridge, which
# gets its own HTB class like any IP peer.
if self.topology.is_dual_udp_edge(node_id, peer_id):
ip_peers[peer_id] = self.topology.nodes[peer_id].docker_ip
else:
ip_peers[peer_id] = self.topology.nodes[peer_id].docker_ip
# --- IP-based peers (UDP/TCP): HTB + u32 on eth0 ---
if ip_peers:
cmds = [f"tc qdisc del dev {IFACE} root 2>/dev/null || true"]
cmds.append(
f"tc qdisc add dev {IFACE} root handle 1: htb default 99"
)
cmds.append(
f"tc class add dev {IFACE} parent 1: classid 1:99 htb rate 1gbit"
)
container_states = {}
for idx, (peer_id, dest_ip) in enumerate(ip_peers.items(), start=1):
class_id = f"1:{idx}"
netem_handle = f"{idx + 10}:"
policy = self._policy_for_edge(node_id, peer_id)
if (
self.topology.is_dual_udp_edge(node_id, peer_id)
and self.config.dual_udp_policy is not None
):
policy = self.config.dual_udp_policy
params = self._sample_policy(policy)
rate = self._htb_rate(node_id, peer_id)
rate_mbit = self._edge_rates.get((node_id, peer_id), 0)
cmds.append(
f"tc class add dev {IFACE} parent 1: classid {class_id} htb rate {rate}"
)
cmds.append(
f"tc qdisc add dev {IFACE} parent {class_id} "
f"handle {netem_handle} netem {params.to_tc_args()}"
)
cmds.append(
f"tc filter add dev {IFACE} parent 1: protocol ip "
f"prio {idx} u32 match ip dst {dest_ip}/32 flowid {class_id}"
)
ingress_rate = self._ingress_rates.get((node_id, peer_id), 0)
ingress_burst = self._ingress_config.burst_bytes if self._ingress_config else 32000
state = LinkNetemState(
container=container,
dest_ip=dest_ip,
class_id=class_id,
netem_handle=netem_handle,
params=params,
rate_mbit=rate_mbit,
ingress_rate_kbps=ingress_rate,
ingress_burst_bytes=ingress_burst,
ingress_filter_prio=idx,
)
container_states[dest_ip] = state
# Ingress policing: add ingress qdisc + per-peer policer filters
if self._ingress_rates:
cmds.append(f"tc qdisc add dev {IFACE} ingress 2>/dev/null || true")
for idx, (peer_id, dest_ip) in enumerate(ip_peers.items(), start=1):
ingress_rate = self._ingress_rates.get((node_id, peer_id), 0)
if ingress_rate > 0:
ingress_burst = self._ingress_config.burst_bytes if self._ingress_config else 32000
cmds.append(
f"tc filter add dev {IFACE} parent ffff: protocol ip "
f"prio {idx} u32 match ip src {dest_ip}/32 "
f"police rate {ingress_rate}kbit burst {ingress_burst} drop"
)
full_cmd = " && ".join(cmds)
result = docker_exec_quiet(container, full_cmd, timeout=30)
if result is not None:
log.info(
"Configured per-link netem on %s (%d IP peers%s)",
container,
len(ip_peers),
", ingress policing" if self._ingress_rates else "",
)
else:
log.warning("Failed to configure netem on %s", container)
self.states[container] = container_states
# --- Ethernet peers: simple netem on veth ---
if eth_peers:
container_veth_states = {}
for peer_id in eth_peers:
iface = veth_interface_name(node_id, peer_id)
policy = self._policy_for_edge(node_id, peer_id)
params = self._sample_policy(policy)
cmd = (
f"tc qdisc del dev {iface} root 2>/dev/null || true && "
f"tc qdisc add dev {iface} root netem {params.to_tc_args()}"
)
result = docker_exec_quiet(container, cmd, timeout=10)
if result is not None:
log.debug("Veth netem on %s:%s -> %s", container, iface, params.to_tc_args())
else:
log.warning("Failed to configure veth netem on %s:%s", container, iface)
container_veth_states[iface] = VethNetemState(
container=container,
iface=iface,
params=params,
)
self.veth_states[container] = container_veth_states
log.info(
"Configured veth netem on %s (%d Ethernet peers)",
container,
len(eth_peers),
)
def setup_node(self, node_id: str):
"""Re-apply HTB/netem/filters for a single node (after container restart).
Uses the saved state from the initial setup so the node gets the same
class IDs and current netem params it had before going down.
"""
container = self.topology.container_name(node_id)
# Re-apply IP-based netem (eth0 HTB + u32) for UDP/TCP peers
container_states = self.states.get(container)
if container_states:
cmds = [f"tc qdisc del dev {IFACE} root 2>/dev/null || true"]
cmds.append(
f"tc qdisc add dev {IFACE} root handle 1: htb default 99"
)
cmds.append(
f"tc class add dev {IFACE} parent 1: classid 1:99 htb rate 1gbit"
)
for dest_ip, state in container_states.items():
rate = f"{state.rate_mbit}mbit" if state.rate_mbit > 0 else "1gbit"
cmds.append(
f"tc class add dev {IFACE} parent 1: classid {state.class_id} htb rate {rate}"
)
cmds.append(
f"tc qdisc add dev {IFACE} parent {state.class_id} "
f"handle {state.netem_handle} netem {state.tc_args()}"
)
prio = state.class_id.split(":")[1]
cmds.append(
f"tc filter add dev {IFACE} parent 1: protocol ip "
f"prio {prio} u32 match ip dst {dest_ip}/32 flowid {state.class_id}"
)
# Re-apply ingress policing
has_ingress = any(s.ingress_rate_kbps > 0 for s in container_states.values())
if has_ingress:
cmds.append(f"tc qdisc add dev {IFACE} ingress 2>/dev/null || true")
for dest_ip, state in container_states.items():
if state.ingress_rate_kbps > 0:
cmds.append(
f"tc filter add dev {IFACE} parent ffff: protocol ip "
f"prio {state.ingress_filter_prio} u32 match ip src {dest_ip}/32 "
f"police rate {state.ingress_rate_kbps}kbit "
f"burst {state.ingress_burst_bytes} drop"
)
full_cmd = " && ".join(cmds)
result = docker_exec_quiet(container, full_cmd, timeout=30)
if result is not None:
log.info(
"Re-applied IP netem on %s (%d peers%s)",
container,
len(container_states),
", ingress policing" if has_ingress else "",
)
else:
log.warning("Failed to re-apply IP netem on %s", container)
# Re-apply Ethernet veth netem
veth_states = self.veth_states.get(container)
if veth_states:
for state in veth_states.values():
self._apply_veth(state)
log.info(
"Re-applied veth netem on %s (%d Ethernet peers)",
container,
len(veth_states),
)
# And on each running neighbour's end of those links. The restore
# recreates the whole pair, so the survivor's end is a new interface
# with no qdisc, and without this that direction of every restored
# link ran unshaped for the rest of the run.
for peer_id in sorted(self.topology.nodes[node_id].peers):
if peer_id in self.down_nodes:
continue
if not self.topology.is_veth_transport(
self.topology.transport_for_edge(node_id, peer_id)
):
continue
peer_container = self.topology.container_name(peer_id)
state = self.veth_states.get(peer_container, {}).get(
veth_interface_name(peer_id, node_id)
)
if state is not None:
self._apply_veth(state)
def _apply_veth(self, state: VethNetemState):
"""Install a veth end's current netem parameters as its root qdisc."""
cmd = (
f"tc qdisc del dev {state.iface} root 2>/dev/null || true && "
f"tc qdisc add dev {state.iface} root netem {state.tc_args()}"
)
result = docker_exec_quiet(state.container, cmd, timeout=10)
if result is not None:
log.debug("Re-applied veth netem on %s:%s", state.container, state.iface)
else:
log.warning("Failed to re-apply veth netem on %s:%s", state.container, state.iface)
def mutate(self):
"""Randomly mutate netem params on a fraction of links."""
if not self.config.mutation.policies:
return
# Only consider edges where both endpoints are up, and never the
# explicitly excluded ones (links an assertion depends on).
live_edges = [
(a, b) for a, b in self.topology.edges
if a not in self.down_nodes and b not in self.down_nodes
and "-".join(sorted([a, b])) not in self._mutation_exclude
]
if not live_edges:
return
num_to_mutate = max(1, int(len(live_edges) * self.config.mutation.fraction))
edges_to_mutate = self.rng.sample(
live_edges, min(num_to_mutate, len(live_edges))
)
# Pick a random policy for this mutation round
policy_name = self.rng.choice(list(self.config.mutation.policies.keys()))
policy = self.config.mutation.policies[policy_name]
log.info(
"Netem mutation: %d links -> '%s' policy",
len(edges_to_mutate),
policy_name,
)
for a, b in edges_to_mutate:
params = self._sample_policy(policy)
self._update_link(a, b, params)
def _update_link(self, node_a: str, node_b: str, params: NetemParams):
"""Update netem on both directions of a link.
A direction held down by a link flap only records the new params;
LinkManager installs them when the flap ends.
"""
transport = self.topology.transport_for_edge(node_a, node_b)
for src, dst in [(node_a, node_b), (node_b, node_a)]:
if src in self.down_nodes:
continue
container = self.topology.container_name(src)
# Safety net: detect containers that crashed outside of NodeManager
if not is_container_running(container):
log.debug(
"Container %s not running (unexpected), marking %s as down",
container,
src,
)
self.down_nodes.add(src)
continue
if self.topology.is_veth_transport(transport):
# Ethernet, or a UDP instance bound to a veth: simple netem
# replace on the veth
iface = veth_interface_name(src, dst)
veth_states = self.veth_states.get(container, {})
state = veth_states.get(iface)
if state is None:
continue
if state.held_down:
state.params = params
log.debug("Deferred veth netem %s:%s until the flap ends", src, iface)
continue
cmd = f"tc qdisc replace dev {iface} root netem {params.to_tc_args()}"
result = docker_exec_quiet(container, cmd)
if result is not None:
state.params = params
log.debug("Updated veth netem %s:%s -> %s", src, iface, params.to_tc_args())
else:
# IP-based (UDP/TCP): HTB class-based netem on eth0
dest_ip = self.topology.nodes[dst].docker_ip
states = self.states.get(container, {})
state = states.get(dest_ip)
if state is None:
continue
if state.held_down:
state.params = params
log.debug("Deferred netem %s -> %s until the flap ends", src, dst)
continue
cmd = (
f"tc qdisc replace dev {IFACE} parent {state.class_id} "
f"handle {state.netem_handle} netem {params.to_tc_args()}"
)
result = docker_exec_quiet(container, cmd)
if result is not None:
state.params = params
log.debug(
"Updated netem %s -> %s: %s",
src,
dst,
params.to_tc_args(),
)
def _sample_policy(self, policy: NetemPolicy) -> NetemParams:
"""Sample concrete params from a policy's ranges."""
return NetemParams(
delay_ms=int(self.rng.uniform(policy.delay_ms[0], policy.delay_ms[1])),
jitter_ms=int(self.rng.uniform(policy.jitter_ms[0], policy.jitter_ms[1])),
loss_pct=round(
self.rng.uniform(policy.loss_pct[0], policy.loss_pct[1]), 1
),
duplicate_pct=round(
self.rng.uniform(policy.duplicate_pct[0], policy.duplicate_pct[1]), 1
),
reorder_pct=round(
self.rng.uniform(policy.reorder_pct[0], policy.reorder_pct[1]), 1
),
corrupt_pct=round(
self.rng.uniform(policy.corrupt_pct[0], policy.corrupt_pct[1]), 1
),
)