diff --git a/src/transport/ethernet/mod.rs b/src/transport/ethernet/mod.rs index 74df5caa..ba81dc64 100644 --- a/src/transport/ethernet/mod.rs +++ b/src/transport/ethernet/mod.rs @@ -1209,6 +1209,42 @@ fn report_sustained_absence(ctx: &Arc) -> bool { // Receive Loop // ============================================================================ +/// Why a received data frame was dropped by [`data_payload`]. +/// +/// Private and only logged at trace, so not a `thiserror` type. +#[derive(Debug, PartialEq, Eq)] +enum DataFrameError { + /// Fewer bytes than the type byte and the length field. + TooShort { len: usize }, + /// The length field names more payload than the frame carries. + LengthExceeds { + payload_len: usize, + available: usize, + }, +} + +/// Extract the payload of a received data frame. +/// +/// `frame` is the received bytes including the type byte, which the caller +/// has already matched as `FRAME_TYPE_DATA`. The layout is +/// `[type:1][length:2 LE][payload:N][padding]`: the length field is what +/// trims Ethernet minimum-frame padding, which would otherwise be handed to +/// AEAD verification as ciphertext. +fn data_payload(frame: &[u8]) -> Result<&[u8], DataFrameError> { + if frame.len() < 3 { + return Err(DataFrameError::TooShort { len: frame.len() }); + } + let payload_len = u16::from_le_bytes([frame[1], frame[2]]) as usize; + let available = frame.len() - 3; + if payload_len > available { + return Err(DataFrameError::LengthExceeds { + payload_len, + available, + }); + } + Ok(&frame[3..3 + payload_len]) +} + /// Ethernet receive loop — runs as a spawned task. /// /// Returns on a dead socket (see [`RECV_ERROR_EXIT_THRESHOLD`]); the binder @@ -1241,29 +1277,31 @@ async fn ethernet_receive_loop( let frame_type = buf[0]; match frame_type { FRAME_TYPE_DATA => { - // Data frame: [type:1][length:2 LE][payload:N] - // Use the length field to trim Ethernet minimum-frame padding. - if len < 3 { - trace!("Data frame too short ({len} bytes), ignoring"); - continue; - } - let payload_len = u16::from_le_bytes([buf[1], buf[2]]) as usize; - if payload_len > len - 3 { - trace!( - "Data frame length field ({payload_len}) exceeds \ - available bytes ({}), ignoring", - len - 3 - ); - continue; - } - let data = buf[3..3 + payload_len].to_vec(); + let data = match data_payload(&buf[..len]) { + Ok(payload) => payload.to_vec(), + Err(DataFrameError::TooShort { len }) => { + trace!("Data frame too short ({len} bytes), ignoring"); + continue; + } + Err(DataFrameError::LengthExceeds { + payload_len, + available, + }) => { + trace!( + "Data frame length field ({payload_len}) exceeds \ + available bytes ({available}), ignoring" + ); + continue; + } + }; + let bytes = data.len(); let addr = TransportAddr::from_bytes(&src_mac); let packet = ReceivedPacket::new(transport_id, addr, data); trace!( transport_id = %transport_id, remote_mac = %format_mac(&src_mac), - bytes = payload_len, + bytes, "Ethernet data frame received" ); @@ -1530,26 +1568,62 @@ mod tests { assert_eq!(&frame[3..], &[1, 2, 3, 4]); // payload } + /// Build a data frame as the sender does: type, LE length, payload. + fn data_frame(payload: &[u8]) -> Vec { + let mut frame = Vec::with_capacity(3 + payload.len()); + frame.push(FRAME_TYPE_DATA); + frame.extend_from_slice(&(payload.len() as u16).to_le_bytes()); + frame.extend_from_slice(payload); + frame + } + #[test] fn test_data_frame_padding_trimmed() { // Simulate Ethernet minimum-frame padding: a 4-byte payload produces // a 7-byte frame (type + len + payload), padded to 46 bytes by NIC. - let payload = vec![0xAA, 0xBB, 0xCC, 0xDD]; - let payload_len = payload.len() as u16; - - // Build frame as sender would - let mut frame = Vec::with_capacity(3 + payload.len()); - frame.push(FRAME_TYPE_DATA); - frame.extend_from_slice(&payload_len.to_le_bytes()); - frame.extend_from_slice(&payload); - - // Simulate NIC padding to 46 bytes + // Drives the receive loop's own parse, not a copy of it. + let mut frame = data_frame(&[0xAA, 0xBB, 0xCC, 0xDD]); frame.resize(46, 0x00); - // Receiver extracts using length field - let recv_len = u16::from_le_bytes([frame[1], frame[2]]) as usize; - let extracted = &frame[3..3 + recv_len]; - assert_eq!(extracted, &[0xAA, 0xBB, 0xCC, 0xDD]); + assert_eq!(data_payload(&frame), Ok(&[0xAA, 0xBB, 0xCC, 0xDD][..])); + } + + #[test] + fn test_data_payload_unpadded_frame_returns_whole_payload() { + let payload: Vec = (0..200).map(|i| i as u8).collect(); + let frame = data_frame(&payload); + + assert_eq!(data_payload(&frame), Ok(&payload[..])); + } + + #[test] + fn test_data_payload_zero_length_padded_frame_returns_empty() { + let mut frame = data_frame(&[]); + frame.resize(46, 0x00); + + assert_eq!(data_payload(&frame), Ok(&[][..])); + } + + #[test] + fn test_data_payload_frame_shorter_than_header_is_rejected() { + assert_eq!( + data_payload(&[FRAME_TYPE_DATA, 0x04]), + Err(DataFrameError::TooShort { len: 2 }) + ); + } + + #[test] + fn test_data_payload_length_field_beyond_received_bytes_is_rejected() { + let mut frame = data_frame(&[0xAA, 0xBB, 0xCC, 0xDD]); + frame.truncate(6); + + assert_eq!( + data_payload(&frame), + Err(DataFrameError::LengthExceeds { + payload_len: 4, + available: 3 + }) + ); } #[test] diff --git a/testing/chaos/sim/links.py b/testing/chaos/sim/links.py index 577ded5b..e6967957 100644 --- a/testing/chaos/sim/links.py +++ b/testing/chaos/sim/links.py @@ -4,6 +4,11 @@ Simulates link failures by setting netem to 100% packet loss on the specific tc class for that peer. Requires the NetemManager to have already set up per-link classful qdiscs. Includes connectivity protection to prevent graph partitioning. + +The flap is recorded in the NetemManager's per-direction state +(``held_down``) rather than here, so every path that re-installs a +qdisc (a node restart's re-apply, a mutation) keeps the link down for +the flap's declared length. """ from __future__ import annotations @@ -28,9 +33,6 @@ class LinkState: is_down: bool = False down_since: float | None = None restore_at: float | None = None - # Saved netem params to restore when link comes back up - saved_params_a: str | None = None # tc args for a->b direction - saved_params_b: str | None = None # tc args for b->a direction class LinkManager: @@ -106,9 +108,8 @@ class LinkManager: a, b = edge state = self.link_states[edge] - # Save current netem params and apply 100% loss - state.saved_params_a = self._set_loss(a, b, "loss 100%") - state.saved_params_b = self._set_loss(b, a, "loss 100%") + self._set_held(a, b, True) + self._set_held(b, a, True) now = time.time() state.is_down = True @@ -118,41 +119,64 @@ class LinkManager: log.info("Link DOWN: %s -- %s (restore in %.0fs)", a, b, duration) def _link_up(self, edge: tuple[str, str]): - """Restore link by reverting netem to saved params.""" + """Restore link by re-installing each direction's current netem params. + + The current params include any mutation recorded during the flap. + """ a, b = edge state = self.link_states[edge] - # Restore previous netem params - if state.saved_params_a: - self._set_loss(a, b, state.saved_params_a) - if state.saved_params_b: - self._set_loss(b, a, state.saved_params_b) + self._set_held(a, b, False) + self._set_held(b, a, False) down_for = time.time() - state.down_since if state.down_since else 0 state.is_down = False state.down_since = None state.restore_at = None - state.saved_params_a = None - state.saved_params_b = None log.info("Link UP: %s -- %s (was down %.0fs)", a, b, down_for) - def _set_loss(self, src_node: str, dst_node: str, netem_args: str) -> str | None: - """Set netem args on the link for src->dst. Returns the previous netem args. + def _set_held(self, src_node: str, dst_node: str, held: bool): + """Mark the src->dst direction held down (or not) and install it. + + The flag is set whether or not the container is running, so a node + that restarts during or after the flap re-installs the right params. + Whether to issue tc now is decided by the live running check alone: + the NetemManager's down_nodes set is not a reliable liveness signal + (its only writers are the safety nets here and in _update_link, and + nothing clears it when a node restarts). Transport-aware: UDP links use tc class on eth0, Ethernet links use tc qdisc replace on the dedicated veth interface. """ if not self.netem_mgr: - return None - - # Skip if the node's container is down - if src_node in self.netem_mgr.down_nodes: - return None + return container = self.topology.container_name(src_node) + transport = self.topology.transport_for_edge(src_node, dst_node) - # Safety net: detect containers that crashed outside of NodeManager + if transport == "ethernet": + iface = veth_interface_name(src_node, dst_node) + state = self.netem_mgr.veth_states.get(container, {}).get(iface) + if state is None: + log.warning("No veth netem state for %s -> %s (%s)", src_node, dst_node, iface) + return + cmd = f"tc qdisc replace dev {iface} root netem " + else: + dest_ip = self.topology.nodes[dst_node].docker_ip + state = self.netem_mgr.states.get(container, {}).get(dest_ip) + if state is None: + log.warning("No netem state for %s -> %s", src_node, dst_node) + return + cmd = ( + f"tc qdisc replace dev {IFACE} parent {state.class_id} " + f"handle {state.netem_handle} netem " + ) + + state.held_down = held + + # Safety net: detect containers that crashed outside of NodeManager. + # The restart's setup_node installs state.tc_args() instead. if not is_container_running(container): log.debug( "Container %s not running (unexpected), marking %s as down", @@ -160,40 +184,9 @@ class LinkManager: src_node, ) self.netem_mgr.down_nodes.add(src_node) - return None + return - transport = self.topology.transport_for_edge(src_node, dst_node) - - if transport == "ethernet": - # Ethernet: simple netem on veth - iface = veth_interface_name(src_node, dst_node) - veth_states = self.netem_mgr.veth_states.get(container, {}) - veth_state = veth_states.get(iface) - if veth_state is None: - log.warning("No veth netem state for %s -> %s (%s)", src_node, dst_node, iface) - return None - - prev_args = veth_state.params.to_tc_args() - cmd = f"tc qdisc replace dev {iface} root netem {netem_args}" - docker_exec_quiet(container, cmd) - return prev_args - else: - # IP-based (UDP/TCP): HTB class on eth0 - dest_ip = self.topology.nodes[dst_node].docker_ip - - states = self.netem_mgr.states.get(container, {}) - link_state = states.get(dest_ip) - if link_state is None: - log.warning("No netem state for %s -> %s", src_node, dst_node) - return None - - prev_args = link_state.params.to_tc_args() - cmd = ( - f"tc qdisc replace dev {IFACE} parent {link_state.class_id} " - f"handle {link_state.netem_handle} netem {netem_args}" - ) - docker_exec_quiet(container, cmd) - return prev_args + docker_exec_quiet(container, cmd + state.tc_args()) def _would_disconnect(self, edge: tuple[str, str]) -> bool: """Check if removing this edge (plus currently-down edges) disconnects the graph.""" diff --git a/testing/chaos/sim/netem.py b/testing/chaos/sim/netem.py index e95d8548..8085e2e0 100644 --- a/testing/chaos/sim/netem.py +++ b/testing/chaos/sim/netem.py @@ -75,6 +75,14 @@ class LinkNetemState: 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 @@ -84,6 +92,12 @@ class VethNetemState: 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: @@ -334,7 +348,7 @@ class NetemManager: ) cmds.append( f"tc qdisc add dev {IFACE} parent {state.class_id} " - f"handle {state.netem_handle} netem {state.params.to_tc_args()}" + f"handle {state.netem_handle} netem {state.tc_args()}" ) prio = state.class_id.split(":")[1] cmds.append( @@ -398,7 +412,7 @@ class NetemManager: """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.params.to_tc_args()}" + 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: @@ -441,7 +455,11 @@ class NetemManager: 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.""" + """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)]: @@ -466,6 +484,10 @@ class NetemManager: 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: @@ -478,6 +500,10 @@ class NetemManager: 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()}" diff --git a/testing/ci-local.sh b/testing/ci-local.sh index a7704955..c9d91bb2 100755 --- a/testing/ci-local.sh +++ b/testing/ci-local.sh @@ -102,6 +102,14 @@ # A preempting CI worker maps 130/143 → "cancelled" (discard, do not record a # failing commit), 0 → green, any other non-zero → red. # +# Host load annotation: every failed suite in the summary is followed by the +# CPU pressure stall (/proc/pressure/cpu) over that suite's run and the peak +# number of other CI runs' containers, labelled LOAD-DEGRADED at or above +# CI_LOAD_THRESHOLD percent (testing/lib/load-annotate.py). It is context for +# reading a red and changes no verdict or exit code. The samples go to +# $FIPS_CI_LOAD_LOG if set (kept after the run), else to a per-run file that +# teardown removes. +# # ── CI parity invariant ───────────────────────────────────────────────────── # This local default suite set and the GitHub integration matrix # (.github/workflows/ci.yml) MUST run the same integration suites, EXCEPT for @@ -317,6 +325,9 @@ OVERALL=0 record() { local name="$1" rc="$2" RESULTS["$name"]=$rc + # A chaos suite's window is written by run_chaos, which runs once per + # scenario; record() runs twice for a parallel one (child and parent). + [[ "$name" == chaos-* ]] || load_mark end "$name" if [[ $rc -ne 0 ]]; then OVERALL=1 fail "$name" @@ -395,6 +406,38 @@ readonly CI_RUN_NAME_SUFFIX # while adding the cross-run prefix that scopes the reap. ci_project() { printf '%s_%s' "$CI_PROJECT_PREFIX" "$1"; } +# ── Host load samples ───────────────────────────────────────────────────── +# +# Stall percent at or above which a failed suite is labelled LOAD-DEGRADED. +# Measured 2026-09-19 on core-vm (12 CPUs): green chaos sets with no foreign +# load ran at 4-11% CPU stall, and the one reproduced load-sensitive red +# (churn-mixed-10 answering 9 of 10) ran at 63%, under 24 CPU hogs plus I/O. +CI_LOAD_THRESHOLD=15 +if [[ -n "${FIPS_CI_LOAD_LOG:-}" ]]; then + CI_LOAD_LOG="$FIPS_CI_LOAD_LOG" + CI_LOAD_LOG_OWNED=0 +else + CI_LOAD_LOG="${TMPDIR:-/tmp}/fips-ci-load-${CI_RUN_ID}.log" + CI_LOAD_LOG_OWNED=1 +fi +CI_LOAD_PID="" + +# Append one window line (begin/end , or barrier) to the load log. +load_mark() { echo "$(date +%s.%N) $*" >> "$CI_LOAD_LOG" 2>/dev/null || true; } + +# Sample CPU pressure and foreign CI runs every 5 s until killed. +load_sampler() { + local psi runs + load_mark barrier + while true; do + psi="$(awk '/^some/ { for (i = 1; i <= NF; i++) if ($i ~ /^total=/) { sub("total=", "", $i); print $i } }' /proc/pressure/cpu 2>/dev/null)" + runs="$(docker ps --filter "label=$CI_LABEL" --format '{{.Label "com.corganlabs.fips-ci.run"}}' 2>/dev/null \ + | sort -u | grep -cvx -e "$CI_RUN_ID" -e '')" + echo "$(date +%s.%N) psi ${psi:-none} runs ${runs:-0} load1 $(cut -d' ' -f1 /proc/loadavg)" >> "$CI_LOAD_LOG" + sleep 5 + done +} + # The name suffix one chaos scenario runs under. Its container names, its # generated-config directory and the token in its host veth names all derive # from this, so teardown recomputes it to know which interfaces are ours. Reads @@ -414,6 +457,13 @@ ci_teardown() { CI_CLEANED=1 local run_status="${1:-1}" + # 0. The load sampler first, so it stops writing before its log goes. + if [[ -n "$CI_LOAD_PID" ]]; then + kill "$CI_LOAD_PID" 2>/dev/null || true + wait "$CI_LOAD_PID" 2>/dev/null || true + fi + [[ "$CI_LOAD_LOG_OWNED" -eq 1 ]] && rm -f "$CI_LOAD_LOG" + # 1. Propagate to parallel chaos children and reap them (bounded). if [[ ${#CI_CHAOS_PIDS[@]} -gt 0 ]]; then kill -TERM "${CI_CHAOS_PIDS[@]}" 2>/dev/null || true @@ -749,6 +799,7 @@ run_chaos() { local results="$SCRIPT_DIR/chaos/sim-results/ci$suffix" local -x FIPS_SIM_OUTPUT="$results" + load_mark begin "chaos-$name" info "[chaos/$name] Running simulation" if bash testing/chaos/scripts/chaos.sh "$@" 2>&1; then rc=0 @@ -757,6 +808,7 @@ run_chaos() { chaos_dump "$name" "$results" fi + load_mark end "chaos-$name" record "chaos-$name" $rc # record() ends in pass()/fail(), which are echoes, so it returns 0 for any @@ -1476,6 +1528,9 @@ run_integration() { # All chaos children have been waited on; clear so a later signal does # not try to kill already-reaped PIDs. CI_CHAOS_PIDS=() + # The next sequential suite's window starts here, not at the last + # suite before the chaos block. + load_mark barrier fi # Sidecar @@ -1573,10 +1628,13 @@ print_summary() { else failed=$((failed + 1)) echo -e " ${RED}✗${RESET} $name" + python3 "$SCRIPT_DIR/lib/load-annotate.py" "$CI_LOAD_LOG" \ + --threshold "$CI_LOAD_THRESHOLD" "$name" 2>&1 | sed 's/^/ /' fi done echo "" + python3 "$SCRIPT_DIR/lib/load-annotate.py" "$CI_LOAD_LOG" --run 2>&1 | sed 's/^/ /' echo -e " ${BOLD}Total: $total Passed: $passed Failed: $failed${RESET}" echo "" @@ -1691,6 +1749,9 @@ main() { stage "FIPS Local CI" info "Project root: $PROJECT_ROOT" + load_sampler & + CI_LOAD_PID=$! + # Above the mode branches deliberately, so --only, --test-only and # --build-only are gated too: a divergence invalidates any claim that a # local run means what a GitHub run means, whichever subset was asked for. diff --git a/testing/lib/load-annotate.py b/testing/lib/load-annotate.py new file mode 100644 index 00000000..d58e46a6 --- /dev/null +++ b/testing/lib/load-annotate.py @@ -0,0 +1,181 @@ +#!/usr/bin/env python3 +"""Annotate failed local-CI suites with the host load during their run. + +ci-local.sh samples CPU pressure (/proc/pressure/cpu, "some" total stall +microseconds) and the count of other CI runs' containers every few +seconds into a load log, and writes window lines around each suite. This +reads that log and prints, for each failed suite named on the command +line, the stall percentage over the suite's window, the peak number of +foreign runs, and a label: + + LOAD-DEGRADED stall at or above the threshold + not load-degraded stall below it + load unknown (...) the window or the samples needed to say are missing + +The label is context for a human reading a red. It changes no verdict: +a red stays a red. Unknown is never reported as "not load-degraded". + +Log lines (whitespace-separated, first field a Unix epoch): + + psi runs load1 + barrier + begin (chaos suites only) + end + +A chaos suite's window runs from its begin to its end. Any other suite's +window runs from the latest preceding "end" of a non-chaos suite or +"barrier" line up to its own end. + +Usage: + load-annotate.py --threshold ... + load-annotate.py --run +""" + +from __future__ import annotations + +import argparse +import sys +from dataclasses import dataclass + + +@dataclass +class Sample: + """One sampler line.""" + + t: float + psi: int | None + runs: int | None + + +def parse(path: str): + """Return (samples, events) from a load log; events keep file order.""" + samples: list[Sample] = [] + events: list[tuple[float, str, str | None]] = [] + with open(path) as f: + for line in f: + parts = line.split() + if len(parts) < 2: + continue + try: + t = float(parts[0]) + except ValueError: + continue + kind = parts[1] + if kind == "psi": + fields = dict(zip(parts[1::2], parts[2::2])) + psi = fields.get("psi") + runs = fields.get("runs") + samples.append(Sample( + t, + int(psi) if psi and psi.isdigit() else None, + int(runs) if runs and runs.isdigit() else None, + )) + elif kind == "barrier": + events.append((t, "barrier", None)) + elif kind in ("begin", "end") and len(parts) >= 3: + events.append((t, kind, parts[2])) + return samples, events + + +def window(name: str, events) -> tuple[float, float] | str: + """Return (start, end) for a suite, or the reason there is none.""" + if name.startswith("chaos-"): + begins = [e for e in events if e[1] == "begin" and e[2] == name] + ends = [e for e in events if e[1] == "end" and e[2] == name] + if not begins and not ends: + return "no window" + if len(begins) != 1 or len(ends) != 1 or ends[0][0] < begins[0][0]: + return "ambiguous window" + return begins[0][0], ends[0][0] + idx = [i for i, e in enumerate(events) if e[1] == "end" and e[2] == name] + if not idx: + return "no window" + if len(idx) > 1: + return "ambiguous window" + end_i = idx[0] + start = None + for e in reversed(events[:end_i]): + if e[1] == "barrier" or (e[1] == "end" and not e[2].startswith("chaos-")): + start = e[0] + break + if start is None: + return "no window" + return start, events[end_i][0] + + +def stall(samples: list[Sample], start: float, end: float) -> float | str: + """Stall percent over [start, end], or the reason it cannot be given.""" + psi = [s for s in samples if s.psi is not None] + if not psi: + return "no PSI samples" + before = [s for s in psi if s.t <= start] + after = [s for s in psi if s.t >= end] + s0 = before[-1] if before else next((s for s in psi if s.t >= start), None) + s1 = after[0] if after else next((s for s in reversed(psi) if s.t <= end), None) + if s0 is None or s1 is None or s1.t - s0.t <= 0: + return "too few PSI samples in window" + return (s1.psi - s0.psi) / ((s1.t - s0.t) * 1e6) * 100 + + +def peak_runs(samples: list[Sample], start: float, end: float) -> int | None: + """Largest foreign-run count sampled inside the window.""" + runs = [s.runs for s in samples if start <= s.t <= end and s.runs is not None] + return max(runs) if runs else None + + +def annotate(name: str, samples, events, threshold: float) -> str: + """One annotation line for a failed suite.""" + w = window(name, events) + if isinstance(w, str): + return f"{name}: load unknown ({w})" + start, end = w + pct = stall(samples, start, end) + if isinstance(pct, str): + return f"{name}: load unknown ({pct})" + runs = peak_runs(samples, start, end) + label = "LOAD-DEGRADED" if pct >= threshold else "not load-degraded" + runs_txt = "?" if runs is None else str(runs) + return (f"{name}: {label} (cpu stall {pct:.1f}% over {end - start:.0f}s, " + f"threshold {threshold:g}%; peak foreign CI runs {runs_txt})") + + +def run_line(samples) -> str: + """Run-level summary: peak interval stall and peak foreign runs.""" + psi = [s for s in samples if s.psi is not None] + peaks = [ + (b.psi - a.psi) / ((b.t - a.t) * 1e6) * 100 + for a, b in zip(psi, psi[1:]) if b.t > a.t + ] + runs = [s.runs for s in samples if s.runs is not None] + stall_txt = f"{max(peaks):.1f}%" if peaks else "unknown (no PSI samples)" + runs_txt = str(max(runs)) if runs else "unknown" + return f"host load: peak cpu stall {stall_txt}; peak foreign CI runs {runs_txt}" + + +def main(argv: list[str]) -> int: + """Entry point.""" + ap = argparse.ArgumentParser(description=__doc__.splitlines()[0]) + ap.add_argument("log") + ap.add_argument("--threshold", type=float) + ap.add_argument("--run", action="store_true") + ap.add_argument("suites", nargs="*") + a = ap.parse_intermixed_args(argv) + try: + samples, events = parse(a.log) + except OSError as exc: + for name in a.suites: + print(f"{name}: load unknown (load log unreadable: {exc.strerror})") + if a.run: + print("host load: unknown (load log unreadable)") + return 0 + if a.run: + print(run_line(samples)) + if a.suites and a.threshold is None: + ap.error("--threshold is required with suite names") + for name in a.suites: + print(annotate(name, samples, events, a.threshold)) + return 0 + + +if __name__ == "__main__": + sys.exit(main(sys.argv[1:]))