Merge maint into master: hold chaos link flaps across restarts, test the Ethernet data-frame parse, annotate failed local CI suites with host load

Conflicts resolved:

- testing/chaos/sim/netem.py: master re-applies veth netem through
  _apply_veth, on the restarted node and on each surviving neighbour's end
  of a recreated pair. _apply_veth now installs state.tc_args(), so both
  ends of a flapped Ethernet link stay at 100% loss across a node restart.
  On maint only the restarted end was held.
- testing/chaos/sim/runner.py: kept master's version. Master already waits
  for three agreeing tree reads before the final snapshot, which spans the
  restored nodes' control-socket start-up, so maint's separate
  control-socket wait is not added and the final snapshot keeps a single
  wait.
- testing/ci-local.sh: run_chaos keeps master's per-run results directory
  and adds maint's load-window begin mark after it.
This commit is contained in:
Johnathan Corgan
2026-09-19 18:17:17 +00:00
5 changed files with 424 additions and 89 deletions
+105 -31
View File
@@ -1209,6 +1209,42 @@ fn report_sustained_absence(ctx: &Arc<BinderContext>) -> 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<u8> {
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<u8> = (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]
+48 -55
View File
@@ -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."""
+29 -3
View File
@@ -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()}"
+61
View File
@@ -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 <suite>, 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.
+181
View File
@@ -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):
<t> psi <some_total_us|none> runs <n> load1 <x>
<t> barrier
<t> begin <suite> (chaos suites only)
<t> end <suite>
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 <load-log> --threshold <pct> <failed-suite>...
load-annotate.py <load-log> --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:]))