use serde::Deserialize;
#[derive(Debug, Clone, PartialEq, Eq, Deserialize)]
pub struct NodeSpec {
pub name: String,
pub host: String,
pub repo: String,
pub devices: Vec<String>,
#[serde(rename = "runnerPort", default)]
pub runner_port: Option<u16>,
}
#[derive(Debug, thiserror::Error)]
pub enum NodesError {
#[error("nodes yaml is malformed: {message}")]
Malformed { message: String },
#[error("nodes yaml lists no nodes")]
Empty,
#[error("node '{node}' lists no devices")]
EmptyDevices { node: String },
#[error("duplicate node name '{name}'")]
DuplicateName { name: String },
}
#[derive(Deserialize)]
struct NodesFile {
nodes: Vec<NodeSpec>,
}
pub fn parse_nodes(yaml: &str) -> Result<Vec<NodeSpec>, NodesError> {
let file: NodesFile = serde_norway::from_str(yaml).map_err(|e| NodesError::Malformed {
message: e.to_string(),
})?;
if file.nodes.is_empty() {
return Err(NodesError::Empty);
}
let mut seen = std::collections::HashSet::new();
for node in &file.nodes {
if node.devices.is_empty() {
return Err(NodesError::EmptyDevices {
node: node.name.clone(),
});
}
if !seen.insert(node.name.as_str()) {
return Err(NodesError::DuplicateName {
name: node.name.clone(),
});
}
}
Ok(file.nodes)
}
#[must_use]
pub fn expand_slots(nodes: &[NodeSpec]) -> Vec<(usize, String)> {
nodes
.iter()
.enumerate()
.flat_map(|(node, spec)| spec.devices.iter().map(move |d| (node, d.clone())))
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SlotAssignment {
pub node: usize,
pub device_ref: String,
pub flows: Vec<usize>,
}
#[must_use]
pub fn assign_flows(flow_count: usize, slots: &[(usize, String)]) -> Vec<SlotAssignment> {
crate::parallel::shard_flows(flow_count, slots.len())
.into_iter()
.zip(slots.iter().cloned())
.map(|(flows, (node, device_ref))| SlotAssignment {
node,
device_ref,
flows,
})
.collect()
}
pub const SSH_TRANSPORT_EXIT: u8 = 255;
#[must_use]
pub fn is_transport_failure(code: u8) -> bool {
code == SSH_TRANSPORT_EXIT
}
#[must_use]
pub fn shell_quote(s: &str) -> String {
format!("'{}'", s.replace('\'', "'\\''"))
}
#[must_use]
pub fn remote_argv(
node: &NodeSpec,
flows: &[String],
device_ref: &str,
passthrough: &[String],
) -> Vec<String> {
let quoted_flows: Vec<String> = flows.iter().map(|f| shell_quote(f)).collect();
let smix_argv =
crate::parallel::child_argv("ed_flows, &shell_quote(device_ref), passthrough);
let remote = format!(
"cd {} && target/release/smix {} --format json",
shell_quote(&node.repo),
smix_argv.join(" ")
);
vec![
"-o".to_string(),
"BatchMode=yes".to_string(),
node.host.clone(),
remote,
]
}
pub const FED_BUILD_STAMP: &str = "target/.smix-fed-stamp";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FlowReport {
pub flow: String,
pub outcome: String,
pub raw: serde_json::Value,
}
#[derive(Debug, thiserror::Error)]
pub enum ReportError {
#[error("remote stdout line is not JSON (protocol violation): {line}")]
NotJson { line: String },
#[error("remote report line is missing '{field}': {line}")]
MissingField { field: &'static str, line: String },
}
pub fn parse_report_lines(stdout: &str) -> Result<Vec<FlowReport>, ReportError> {
let mut reports = Vec::new();
for line in stdout.lines() {
if line.trim().is_empty() {
continue;
}
let raw: serde_json::Value =
serde_json::from_str(line).map_err(|_| ReportError::NotJson {
line: line.to_string(),
})?;
let field = |name: &'static str| -> Result<String, ReportError> {
raw.get(name)
.and_then(|v| v.as_str())
.map(str::to_string)
.ok_or(ReportError::MissingField {
field: name,
line: line.to_string(),
})
};
let flow = field("flow")?;
let outcome = field("runOutcome")?;
reports.push(FlowReport { flow, outcome, raw });
}
Ok(reports)
}
#[must_use]
pub fn readiness_argv(node: &NodeSpec) -> Vec<String> {
let remote = format!(
"cd {} && test -f {FED_BUILD_STAMP} && test -x target/release/smix && \
[ -z \"$(find crates -name '*.rs' -newer {FED_BUILD_STAMP})\" ]",
shell_quote(&node.repo),
);
vec![
"-o".to_string(),
"BatchMode=yes".to_string(),
node.host.clone(),
remote,
]
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RemoteOutput {
pub exit: u8,
pub stdout: String,
pub stderr: String,
}
pub fn run_ssh(argv: &[String]) -> std::io::Result<RemoteOutput> {
capture("ssh", argv)
}
fn capture(program: &str, argv: &[String]) -> std::io::Result<RemoteOutput> {
let out = std::process::Command::new(program).args(argv).output()?;
Ok(RemoteOutput {
exit: out.status.code().map_or(1, |c| c.clamp(0, 255) as u8),
stdout: String::from_utf8_lossy(&out.stdout).into_owned(),
stderr: String::from_utf8_lossy(&out.stderr).into_owned(),
})
}
pub const FED_ARTIFACT_DIR: &str = ".smix/fed-artifacts";
#[must_use]
pub fn artifact_pull_argv(node: &NodeSpec, remote_dir: &str, local_dir: &str) -> Vec<String> {
vec![
"-a".to_string(),
format!(
"{}:{}",
node.host,
shell_quote(&format!("{}/{}/", node.repo, remote_dir))
),
format!("{}/{}/", local_dir, node.name),
]
}
pub fn run_rsync(argv: &[String]) -> std::io::Result<RemoteOutput> {
capture("rsync", argv)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SlotResult {
pub node: usize,
pub exit: u8,
pub stdout: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct NodeResult {
pub name: String,
pub exit: u8,
pub reports: Vec<FlowReport>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
pub struct MergedNode {
pub node: String,
pub exit: u8,
pub flows: Vec<serde_json::Value>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
#[serde(rename_all = "camelCase")]
pub struct MergedReport {
pub nodes: Vec<MergedNode>,
pub aggregate_exit: u8,
}
pub fn fold_slot_results(
nodes: &[NodeSpec],
slots: &[SlotResult],
) -> Result<Vec<NodeResult>, ReportError> {
nodes
.iter()
.enumerate()
.map(|(index, spec)| {
let mut exits = Vec::new();
let mut reports = Vec::new();
for slot in slots.iter().filter(|s| s.node == index) {
exits.push(slot.exit);
if !is_transport_failure(slot.exit) {
reports.extend(parse_report_lines(&slot.stdout)?);
}
}
Ok(NodeResult {
name: spec.name.clone(),
exit: crate::parallel::aggregate_exit(&exits),
reports,
})
})
.collect()
}
#[must_use]
pub fn merge_reports(results: &[NodeResult]) -> MergedReport {
let exits: Vec<u8> = results.iter().map(|r| r.exit).collect();
MergedReport {
nodes: results
.iter()
.map(|r| MergedNode {
node: r.name.clone(),
exit: r.exit,
flows: r.reports.iter().map(|f| f.raw.clone()).collect(),
})
.collect(),
aggregate_exit: crate::parallel::aggregate_exit(&exits),
}
}
#[derive(Debug, thiserror::Error)]
pub enum FederationRunError {
#[error("node '{node}' failed the readiness gate (stale or unreachable): {stderr}")]
Gate { node: String, stderr: String },
#[error("spawning ssh for node '{node}': {source}")]
Spawn {
node: String,
source: std::io::Error,
},
#[error(transparent)]
Report(#[from] ReportError),
#[error("artifact rsync for node '{node}' failed (exit {exit}): {stderr}")]
ArtifactPull {
node: String,
exit: u8,
stderr: String,
},
}
pub fn run_federation(
nodes: &[NodeSpec],
assignments: &[SlotAssignment],
flows: &[String],
passthrough: &[String],
pull_to: Option<&std::path::Path>,
) -> Result<MergedReport, FederationRunError> {
for node in nodes {
let gate = run_ssh(&readiness_argv(node)).map_err(|e| FederationRunError::Spawn {
node: node.name.clone(),
source: e,
})?;
if gate.exit != 0 {
return Err(FederationRunError::Gate {
node: node.name.clone(),
stderr: gate.stderr,
});
}
}
let mut running = Vec::new();
for assignment in assignments {
if assignment.flows.is_empty() {
continue;
}
let node = &nodes[assignment.node];
let slot_flows: Vec<String> = assignment.flows.iter().map(|&i| flows[i].clone()).collect();
let mut slot_passthrough = passthrough.to_vec();
if let Some(port) = node.runner_port {
slot_passthrough.push(shell_quote("--runner-port"));
slot_passthrough.push(shell_quote(&port.to_string()));
}
let argv = remote_argv(node, &slot_flows, &assignment.device_ref, &slot_passthrough);
let child = std::process::Command::new("ssh")
.args(&argv)
.stdout(std::process::Stdio::piped())
.spawn()
.map_err(|e| FederationRunError::Spawn {
node: node.name.clone(),
source: e,
})?;
running.push((assignment, child));
}
let mut slot_results = Vec::new();
for (assignment, child) in running {
let node = &nodes[assignment.node];
let out = child
.wait_with_output()
.map_err(|e| FederationRunError::Spawn {
node: node.name.clone(),
source: e,
})?;
let exit = out.status.code().map_or(1, |c| c.clamp(0, 255) as u8);
let note = if is_transport_failure(exit) {
" (ssh transport failure)"
} else {
""
};
eprintln!(
"smix run --nodes: node {} device {} exited {exit}{note}",
node.name, assignment.device_ref
);
slot_results.push(SlotResult {
node: assignment.node,
exit,
stdout: String::from_utf8_lossy(&out.stdout).into_owned(),
});
}
let results = fold_slot_results(nodes, &slot_results)?;
if let Some(local_dir) = pull_to {
let local = local_dir.display().to_string();
for node in nodes {
let pull =
run_rsync(&artifact_pull_argv(node, FED_ARTIFACT_DIR, &local)).map_err(|e| {
FederationRunError::Spawn {
node: node.name.clone(),
source: e,
}
})?;
if pull.exit != 0 {
return Err(FederationRunError::ArtifactPull {
node: node.name.clone(),
exit: pull.exit,
stderr: pull.stderr,
});
}
}
}
Ok(merge_reports(&results))
}
#[cfg(test)]
mod tests {
use super::*;
const DOCUMENTED_NODES_YAML: &str = "\
nodes:
- name: mini
host: mini
repo: /Users/doracawl/workspace/goliajp/smix
devices: [sim-smix-001]
- name: studio
host: studio.local
repo: /Users/doracawl/smix
devices: [sim-smix-002, sim-smix-003]
";
#[test]
fn parses_the_documented_nodes_yaml_shape() {
let nodes = parse_nodes(DOCUMENTED_NODES_YAML).unwrap();
assert_eq!(nodes.len(), 2);
assert_eq!(nodes[0].name, "mini");
assert_eq!(nodes[0].host, "mini");
assert_eq!(nodes[0].repo, "/Users/doracawl/workspace/goliajp/smix");
assert_eq!(nodes[0].devices, vec!["sim-smix-001"]);
assert_eq!(nodes[1].name, "studio");
assert_eq!(nodes[1].host, "studio.local");
assert_eq!(nodes[1].repo, "/Users/doracawl/smix");
assert_eq!(nodes[1].devices, vec!["sim-smix-002", "sim-smix-003"]);
}
#[test]
fn rejects_a_node_without_devices() {
let yaml = "\
nodes:
- name: mini
host: mini
repo: /Users/doracawl/workspace/goliajp/smix
devices: []
";
let err = parse_nodes(yaml).unwrap_err();
assert!(matches!(&err, NodesError::EmptyDevices { node } if node == "mini"));
assert!(err.to_string().contains("mini"));
}
fn roster(specs: &[(&str, &[&str])]) -> Vec<NodeSpec> {
specs
.iter()
.map(|(name, devices)| NodeSpec {
name: (*name).to_string(),
host: (*name).to_string(),
repo: format!("/repo/{name}"),
devices: devices.iter().map(|d| (*d).to_string()).collect(),
runner_port: None,
})
.collect()
}
#[test]
fn parses_an_optional_per_node_runner_port() {
let yaml = "\
nodes:
- name: studio
host: localhost
repo: /repo/studio
devices: [sim-1]
runnerPort: 22097
- name: mini
host: mini
repo: /repo/mini
devices: [sim-2]
";
let nodes = parse_nodes(yaml).unwrap();
assert_eq!(nodes[0].runner_port, Some(22097));
assert_eq!(nodes[1].runner_port, None);
}
#[test]
fn slots_flatten_nodes_in_listing_order() {
let nodes = roster(&[("a", &["a1", "a2"]), ("b", &["b1"])]);
assert_eq!(
expand_slots(&nodes),
vec![
(0, "a1".to_string()),
(0, "a2".to_string()),
(1, "b1".to_string())
]
);
}
#[test]
fn flows_round_robin_over_all_slots_across_nodes() {
let nodes = roster(&[("a", &["a1", "a2"]), ("b", &["b1"])]);
let slots = expand_slots(&nodes);
let assignments = assign_flows(5, &slots);
assert_eq!(
assignments,
vec![
SlotAssignment {
node: 0,
device_ref: "a1".to_string(),
flows: vec![0, 3]
},
SlotAssignment {
node: 0,
device_ref: "a2".to_string(),
flows: vec![1, 4]
},
SlotAssignment {
node: 1,
device_ref: "b1".to_string(),
flows: vec![2]
},
]
);
let buckets: Vec<Vec<usize>> = assignments.into_iter().map(|a| a.flows).collect();
assert_eq!(buckets, crate::parallel::shard_flows(5, 3));
}
#[test]
fn single_node_single_device_degenerates_to_the_sequential_order() {
let nodes = roster(&[("a", &["a1"])]);
let assignments = assign_flows(3, &expand_slots(&nodes));
assert_eq!(
assignments,
vec![SlotAssignment {
node: 0,
device_ref: "a1".to_string(),
flows: vec![0, 1, 2]
}]
);
}
#[test]
fn remote_argv_wraps_child_argv_in_ssh_with_explicit_json_format() {
let node = NodeSpec {
name: "mini".to_string(),
host: "mini".to_string(),
repo: "/Users/doracawl/workspace/goliajp/smix".to_string(),
devices: vec!["sim-smix-001".to_string()],
runner_port: None,
};
let argv = remote_argv(
&node,
&["a.yaml".to_string()],
"sim-smix-001",
&["--no-launch".to_string()],
);
assert_eq!(
argv,
vec![
"-o".to_string(),
"BatchMode=yes".to_string(),
"mini".to_string(),
"cd '/Users/doracawl/workspace/goliajp/smix' && target/release/smix \
run 'a.yaml' --device 'sim-smix-001' --no-launch --format json"
.to_string(),
]
);
let remote = &argv[3];
assert!(remote.contains("--format json"));
assert!(!remote.contains("--parallel"));
assert!(!remote.contains("--also-device"));
}
#[test]
fn shell_quoting_survives_spaces_and_single_quotes() {
assert_eq!(shell_quote("flows/a b.yaml"), "'flows/a b.yaml'");
assert_eq!(shell_quote("it's"), "'it'\\''s'");
assert_eq!(shell_quote(""), "''");
assert_eq!(shell_quote("a.yaml"), "'a.yaml'");
}
#[test]
fn transport_failure_255_wins_the_aggregate() {
assert!(is_transport_failure(SSH_TRANSPORT_EXIT));
for smix_code in [0u8, 1, 2, 3, 4, 5, 6, 130, 143] {
assert!(!is_transport_failure(smix_code));
}
assert_eq!(crate::parallel::aggregate_exit(&[0, 255, 2]), 255);
}
#[test]
fn parses_one_report_line_per_flow() {
let stdout = concat!(
r#"{"flow":"scripts/release/stress-corpus/launch-and-capture.yaml","runOutcome":"success","warnings":[],"steps":[]}"#,
"\n",
r#"{"flow":"scripts/release/stress-corpus/screenshot-twice.yaml","runOutcome":"failure","failure":{"code":"NotVisible","message":"no match","selector":null,"suggestions":[],"visibleCount":3}}"#,
"\n",
);
let reports = parse_report_lines(stdout).unwrap();
assert_eq!(reports.len(), 2);
assert_eq!(
reports[0].flow,
"scripts/release/stress-corpus/launch-and-capture.yaml"
);
assert_eq!(reports[0].outcome, "success");
assert_eq!(
reports[1].flow,
"scripts/release/stress-corpus/screenshot-twice.yaml"
);
assert_eq!(reports[1].outcome, "failure");
assert_eq!(reports[1].raw["failure"]["code"], "NotVisible");
}
#[test]
fn rejects_a_non_json_stdout_line() {
let stdout = concat!(
r#"{"flow":"a.yaml","runOutcome":"success","warnings":[],"steps":[]}"#,
"\n\n",
"kevy: AOF 3 entries replayed\n",
);
let err = parse_report_lines(stdout).unwrap_err();
assert!(
matches!(&err, ReportError::NotJson { line } if line == "kevy: AOF 3 entries replayed")
);
assert!(err.to_string().contains("kevy: AOF 3 entries replayed"));
}
#[test]
fn readiness_argv_pins_the_gate_command() {
let node = NodeSpec {
name: "mini".to_string(),
host: "mini".to_string(),
repo: "/Users/doracawl/workspace/goliajp/smix".to_string(),
devices: vec!["sim-smix-001".to_string()],
runner_port: None,
};
assert_eq!(
readiness_argv(&node),
vec![
"-o".to_string(),
"BatchMode=yes".to_string(),
"mini".to_string(),
"cd '/Users/doracawl/workspace/goliajp/smix' && \
test -f target/.smix-fed-stamp && \
test -x target/release/smix && \
[ -z \"$(find crates -name '*.rs' -newer target/.smix-fed-stamp)\" ]"
.to_string(),
]
);
}
#[test]
#[ignore]
fn federation_e2e_single_node_runs_flows_on_mini() {
let nodes_path = std::env::var("SMIX_FED_E2E_NODES")
.expect("SMIX_FED_E2E_NODES unset — this test is driven by the C3 e2e script");
let flows_env = std::env::var("SMIX_FED_E2E_FLOWS")
.expect("SMIX_FED_E2E_FLOWS unset — this test is driven by the C3 e2e script");
let flows: Vec<String> = flows_env.split(',').map(str::to_string).collect();
assert!(!flows.is_empty());
let yaml = std::fs::read_to_string(&nodes_path).unwrap();
let nodes = parse_nodes(&yaml).unwrap();
let slots = expand_slots(&nodes);
assert_eq!(slots.len(), 1, "single-node e2e expects exactly one slot");
let assignments = assign_flows(flows.len(), &slots);
assert_eq!(assignments[0].flows, (0..flows.len()).collect::<Vec<_>>());
let node = &nodes[assignments[0].node];
let device_ref = &assignments[0].device_ref;
let gate = run_ssh(&readiness_argv(node)).unwrap();
assert_eq!(
gate.exit, 0,
"readiness gate failed — node is stale\nstderr: {}",
gate.stderr
);
let out = run_ssh(&remote_argv(node, &flows, device_ref, &[])).unwrap();
assert!(
!is_transport_failure(out.exit),
"ssh transport failure\nstderr: {}",
out.stderr
);
assert_eq!(out.exit, 0, "remote run failed\nstderr: {}", out.stderr);
let reports = parse_report_lines(&out.stdout).unwrap();
assert_eq!(reports.len(), flows.len());
for (report, flow) in reports.iter().zip(&flows) {
assert_eq!(&report.flow, flow);
assert_eq!(
report.outcome, "success",
"flow {flow} failed: {}",
report.raw
);
}
}
#[test]
fn artifact_pull_argv_pins_the_rsync_command() {
let node = NodeSpec {
name: "mini".to_string(),
host: "mini".to_string(),
repo: "/Users/doracawl/workspace/goliajp/smix".to_string(),
devices: vec!["sim-smix-001".to_string()],
runner_port: None,
};
assert_eq!(
artifact_pull_argv(&node, FED_ARTIFACT_DIR, "/tmp/pull"),
vec![
"-a".to_string(),
"mini:'/Users/doracawl/workspace/goliajp/smix/.smix/fed-artifacts/'".to_string(),
"/tmp/pull/mini/".to_string(),
]
);
}
#[test]
#[ignore]
fn federation_e2e_two_nodes_merge_reports_and_recover_artifacts() {
let nodes_path = std::env::var("SMIX_FED_E2E_NODES")
.expect("SMIX_FED_E2E_NODES unset — this test is driven by the C4 e2e script");
let flows_env = std::env::var("SMIX_FED_E2E_FLOWS")
.expect("SMIX_FED_E2E_FLOWS unset — this test is driven by the C4 e2e script");
let pull_dir = std::env::var("SMIX_FED_E2E_PULL_DIR")
.expect("SMIX_FED_E2E_PULL_DIR unset — this test is driven by the C4 e2e script");
let ports_env = std::env::var("SMIX_FED_E2E_RUNNER_PORTS").unwrap_or_default();
let ports: std::collections::HashMap<String, String> = ports_env
.split(',')
.filter(|s| !s.is_empty())
.map(|pair| {
let (name, port) = pair
.split_once('=')
.expect("SMIX_FED_E2E_RUNNER_PORTS entry is name=port");
(name.to_string(), port.to_string())
})
.collect();
let flows: Vec<String> = flows_env.split(',').map(str::to_string).collect();
assert_eq!(flows.len(), 2, "two-node e2e expects exactly two flows");
let yaml = std::fs::read_to_string(&nodes_path).unwrap();
let nodes = parse_nodes(&yaml).unwrap();
assert_eq!(nodes.len(), 2, "two-node e2e expects exactly two nodes");
let slots = expand_slots(&nodes);
assert_eq!(slots.len(), 2, "two-node e2e expects exactly two slots");
let assignments = assign_flows(flows.len(), &slots);
let mut results = Vec::new();
for assignment in &assignments {
assert_eq!(
assignment.flows.len(),
1,
"each slot carries exactly one flow"
);
let node = &nodes[assignment.node];
let node_flows: Vec<String> =
assignment.flows.iter().map(|&i| flows[i].clone()).collect();
let gate = run_ssh(&readiness_argv(node)).unwrap();
assert_eq!(
gate.exit, 0,
"readiness gate failed on {} — node is stale\nstderr: {}",
node.name, gate.stderr
);
let mut passthrough = vec!["--debug-output".to_string(), FED_ARTIFACT_DIR.to_string()];
if let Some(port) = ports.get(&node.name) {
passthrough.push("--runner-port".to_string());
passthrough.push(port.clone());
}
let out = run_ssh(&remote_argv(
node,
&node_flows,
&assignment.device_ref,
&passthrough,
))
.unwrap();
assert!(
!is_transport_failure(out.exit),
"ssh transport failure on {}\nstderr: {}",
node.name,
out.stderr
);
assert_eq!(
out.exit, 0,
"remote run failed on {}\nstderr: {}",
node.name, out.stderr
);
let reports = parse_report_lines(&out.stdout).unwrap();
assert_eq!(reports.len(), 1);
assert_eq!(&reports[0].flow, &node_flows[0]);
assert_eq!(
reports[0].outcome, "success",
"flow {} failed on {}: {}",
node_flows[0], node.name, reports[0].raw
);
results.push(NodeResult {
name: node.name.clone(),
exit: out.exit,
reports,
});
}
let merged = merge_reports(&results);
assert_eq!(merged.nodes.len(), 2);
assert_eq!(merged.aggregate_exit, 0);
let doc = serde_json::to_string(&merged).unwrap();
for node in &nodes {
assert!(
doc.contains(&node.name),
"merged report is missing node {}",
node.name
);
}
for node in &nodes {
let pull = run_rsync(&artifact_pull_argv(node, FED_ARTIFACT_DIR, &pull_dir)).unwrap();
assert_eq!(
pull.exit, 0,
"artifact rsync failed for {}\nstderr: {}",
node.name, pull.stderr
);
let summary = std::path::Path::new(&pull_dir)
.join(&node.name)
.join("run-summary.json");
assert!(
summary.is_file(),
"recovered artifact missing: {}",
summary.display()
);
}
}
#[test]
fn merges_two_nodes_into_one_wrapped_report_pinning_the_json() {
let leaf_a = r#"{"flow":"scripts/release/stress-corpus/launch-and-capture.yaml","runOutcome":"success","warnings":[],"steps":[]}"#;
let leaf_b = r#"{"flow":"scripts/release/stress-corpus/screenshot-twice.yaml","runOutcome":"success","warnings":[],"steps":[]}"#;
let results = vec![
NodeResult {
name: "a".to_string(),
exit: 0,
reports: parse_report_lines(leaf_a).unwrap(),
},
NodeResult {
name: "b".to_string(),
exit: 0,
reports: parse_report_lines(leaf_b).unwrap(),
},
];
let merged = merge_reports(&results);
assert_eq!(
serde_json::to_string(&merged).unwrap(),
concat!(
r#"{"nodes":["#,
r#"{"node":"a","exit":0,"flows":[{"flow":"scripts/release/stress-corpus/launch-and-capture.yaml","runOutcome":"success","steps":[],"warnings":[]}]},"#,
r#"{"node":"b","exit":0,"flows":[{"flow":"scripts/release/stress-corpus/screenshot-twice.yaml","runOutcome":"success","steps":[],"warnings":[]}]}"#,
r#"],"aggregateExit":0}"#,
)
);
}
#[test]
fn transport_255_node_merges_empty_and_wins_the_aggregate() {
let leaf = r#"{"flow":"a.yaml","runOutcome":"success","warnings":[],"steps":[]}"#;
let results = vec![
NodeResult {
name: "a".to_string(),
exit: 0,
reports: parse_report_lines(leaf).unwrap(),
},
NodeResult {
name: "b".to_string(),
exit: SSH_TRANSPORT_EXIT,
reports: vec![],
},
];
let merged = merge_reports(&results);
assert_eq!(merged.nodes.len(), 2);
assert_eq!(merged.nodes[1].node, "b");
assert_eq!(merged.nodes[1].flows, Vec::<serde_json::Value>::new());
assert_eq!(merged.aggregate_exit, 255);
assert!(is_transport_failure(merged.aggregate_exit));
}
#[test]
fn a_node_with_all_failed_flows_keeps_failure_leaves_verbatim() {
let stdout = concat!(
r#"{"flow":"a.yaml","runOutcome":"failure","failure":{"code":"NotVisible","message":"no match","selector":null,"suggestions":[],"visibleCount":3}}"#,
"\n",
r#"{"flow":"b.yaml","runOutcome":"failure","failure":{"code":"Timeout","message":"gave up","selector":null,"suggestions":[],"visibleCount":0}}"#,
"\n",
);
let results = vec![NodeResult {
name: "mini".to_string(),
exit: 3,
reports: parse_report_lines(stdout).unwrap(),
}];
let merged = merge_reports(&results);
assert_eq!(merged.nodes[0].flows.len(), 2);
assert_eq!(merged.nodes[0].flows[0]["failure"]["code"], "NotVisible");
assert_eq!(merged.nodes[0].flows[1]["failure"]["code"], "Timeout");
assert_eq!(merged.aggregate_exit, 3);
}
#[test]
fn empty_inputs_merge_to_exit_zero() {
let merged = merge_reports(&[]);
assert_eq!(merged.nodes, Vec::<MergedNode>::new());
assert_eq!(merged.aggregate_exit, 0);
let merged = merge_reports(&[NodeResult {
name: "idle".to_string(),
exit: 0,
reports: vec![],
}]);
assert_eq!(merged.nodes.len(), 1);
assert_eq!(merged.nodes[0].flows, Vec::<serde_json::Value>::new());
assert_eq!(merged.aggregate_exit, 0);
}
#[test]
fn fold_groups_slots_by_node_with_max_exit_and_concatenated_reports() {
let nodes = roster(&[("a", &["a1", "a2"]), ("b", &["b1"])]);
let leaf_ok = r#"{"flow":"a.yaml","runOutcome":"success","warnings":[],"steps":[]}"#;
let leaf_fail = r#"{"flow":"b.yaml","runOutcome":"failure","failure":{"code":"NotVisible","message":"no match","selector":null,"suggestions":[],"visibleCount":3}}"#;
let leaf_other = r#"{"flow":"c.yaml","runOutcome":"success","warnings":[],"steps":[]}"#;
let slots = vec![
SlotResult {
node: 0,
exit: 0,
stdout: format!("{leaf_ok}\n"),
},
SlotResult {
node: 0,
exit: 3,
stdout: format!("{leaf_fail}\n"),
},
SlotResult {
node: 1,
exit: 0,
stdout: format!("{leaf_other}\n"),
},
];
let results = fold_slot_results(&nodes, &slots).unwrap();
assert_eq!(results.len(), 2);
assert_eq!(results[0].name, "a");
assert_eq!(results[0].exit, 3);
assert_eq!(results[0].reports.len(), 2);
assert_eq!(results[0].reports[0].flow, "a.yaml");
assert_eq!(results[0].reports[1].flow, "b.yaml");
assert_eq!(results[1].name, "b");
assert_eq!(results[1].exit, 0);
assert_eq!(results[1].reports.len(), 1);
assert_eq!(results[1].reports[0].flow, "c.yaml");
}
#[test]
fn fold_keeps_a_transport_lost_slot_empty_without_parsing_its_stdout() {
let nodes = roster(&[("a", &["a1"])]);
let slots = vec![SlotResult {
node: 0,
exit: SSH_TRANSPORT_EXIT,
stdout: "ssh: connect to host a port 22: Connection refused".to_string(),
}];
let results = fold_slot_results(&nodes, &slots).unwrap();
assert_eq!(results.len(), 1);
assert_eq!(results[0].exit, SSH_TRANSPORT_EXIT);
assert_eq!(results[0].reports, Vec::<FlowReport>::new());
}
#[test]
fn fold_surfaces_a_protocol_violation_from_a_healthy_slot() {
let nodes = roster(&[("a", &["a1"])]);
let slots = vec![SlotResult {
node: 0,
exit: 0,
stdout: "kevy: AOF 3 entries replayed\n".to_string(),
}];
let err = fold_slot_results(&nodes, &slots).unwrap_err();
assert!(
matches!(&err, ReportError::NotJson { line } if line == "kevy: AOF 3 entries replayed")
);
}
#[test]
fn rejects_duplicate_node_names() {
let yaml = "\
nodes:
- name: mini
host: mini-a
repo: /a
devices: [sim-1]
- name: mini
host: mini-b
repo: /b
devices: [sim-2]
";
let err = parse_nodes(yaml).unwrap_err();
assert!(matches!(&err, NodesError::DuplicateName { name } if name == "mini"));
}
}