From 66fec04e75e721360df258460e6916c061def01e Mon Sep 17 00:00:00 2001 From: Michel Paulissen <862400+MichelPaulissen@users.noreply.github.com> Date: Sat, 4 Jul 2026 10:58:00 +0200 Subject: [PATCH] Update public locality failure reporting --- DISASMER_PUBLIC_TREE.json | 4 +- crates/disasmer-cli/src/main.rs | 156 ++++++++++++++++++++++++++-- scripts/cli-first-contract-smoke.js | 9 ++ 3 files changed, 157 insertions(+), 12 deletions(-) diff --git a/DISASMER_PUBLIC_TREE.json b/DISASMER_PUBLIC_TREE.json index 02dda6f..cf45e5a 100644 --- a/DISASMER_PUBLIC_TREE.json +++ b/DISASMER_PUBLIC_TREE.json @@ -1,7 +1,7 @@ { "kind": "disasmer-filtered-public-tree", - "source_commit": "cb99b1e5f0baa07af19211bd164241b5776f91fe", - "release_name": "dryrun-cb99b1e5f0ba", + "source_commit": "69c1074abe5cb9ccb3f31ab39327bfcbfd7b74b3", + "release_name": "dryrun-69c1074abe5c", "filtered_out": [ "private/**", "experiments/**", diff --git a/crates/disasmer-cli/src/main.rs b/crates/disasmer-cli/src/main.rs index a80fb44..4efa9b2 100644 --- a/crates/disasmer-cli/src/main.rs +++ b/crates/disasmer-cli/src/main.rs @@ -2228,10 +2228,12 @@ fn human_report(value: &Value) -> String { } if let Some(tasks) = value.get("current_tasks").and_then(Value::as_array) { push_task_placement_reasons(&mut lines, tasks); + push_task_locality_failures(&mut lines, tasks); } if let Some(tasks) = value.get("tasks").and_then(Value::as_array) { lines.push(format!("tasks: {}", tasks.len())); push_task_placement_reasons(&mut lines, tasks); + push_task_locality_failures(&mut lines, tasks); } if let Some(log_entries) = value.get("log_entries").and_then(Value::as_array) { lines.push(format!("log entries: {}", log_entries.len())); @@ -2477,6 +2479,38 @@ fn push_task_placement_reasons(lines: &mut Vec, tasks: &[Value]) { } } +fn push_task_locality_failures(lines: &mut Vec, tasks: &[Value]) { + for task in tasks { + let Some(locality) = task.get("locality_failure") else { + continue; + }; + if locality.is_null() { + continue; + } + let task_name = task + .get("task") + .and_then(Value::as_str) + .unwrap_or("unknown"); + let affected = locality + .get("affected_data") + .and_then(Value::as_str) + .unwrap_or("direct_transfer"); + let reason = locality + .get("reason") + .and_then(Value::as_str) + .unwrap_or("direct transfer or locality failed"); + lines.push(format!("locality {task_name}: {affected} ({reason})")); + let actions = locality + .get("safe_next_actions") + .and_then(Value::as_array) + .map(|actions| actions.iter().filter_map(Value::as_str).collect::>()) + .unwrap_or_default(); + if !actions.is_empty() { + lines.push(format!("locality next {task_name}: {}", actions.join("; "))); + } + } +} + fn push_nested_string_field(lines: &mut Vec, value: &Value, key: &str, label: &str) { if let Some(text) = value.get(key).and_then(Value::as_str) { lines.push(format!("{label}: {text}")); @@ -2613,6 +2647,9 @@ fn classify_cli_error_message(message: &str) -> &'static str { { return "active_process"; } + if message_mentions_locality_failure(&message) { + return "connectivity"; + } if message.contains("no capable node") || message.contains("missing capability") || message.contains("capability") @@ -2690,6 +2727,16 @@ fn classify_cli_error_message(message: &str) -> &'static str { "unknown" } +fn message_mentions_locality_failure(message: &str) -> bool { + message.contains("direct connectivity unavailable") + || message.contains("direct transfer") + || message.contains("source snapshot unavailable") + || (message.contains("required artifact") + && message.contains("unavailable") + && message.contains("direct connectivity")) + || message.contains("locality assumption") +} + fn cli_error_exit_code(category: &str) -> i64 { match category { "authentication" => 20, @@ -3375,6 +3422,9 @@ fn task_summaries(task_events: Option<&Value>) -> Value { .and_then(|placement| placement.get("score")) .cloned() .unwrap_or(Value::Null); + let failure_reason = task_failure_reason(event); + let locality_failure = task_locality_failure_from_reason(&failure_reason); + let machine_error = task_failure_machine_error_from_reason(event, &failure_reason); json!({ "process": event_string(event, "process"), "task": task, @@ -3388,8 +3438,9 @@ fn task_summaries(task_events: Option<&Value>) -> Value { "reasons": placement_reasons, "explanation_available": placement.is_some(), }, - "failure_reason": task_failure_reason(event), - "machine_error": task_failure_machine_error(event), + "failure_reason": failure_reason, + "locality_failure": locality_failure, + "machine_error": machine_error, "stdout_bytes": event_u64(event, "stdout_bytes").unwrap_or(0), "stderr_bytes": event_u64(event, "stderr_bytes").unwrap_or(0), }) @@ -3416,20 +3467,78 @@ fn task_failure_reason(event: &Value) -> Value { } } -fn task_failure_machine_error(event: &Value) -> Value { +fn task_failure_machine_error_from_reason(event: &Value, reason: &Value) -> Value { match event.get("terminal_state").and_then(Value::as_str) { Some("failed") => { - let reason = task_failure_reason(event) - .as_str() - .unwrap_or("task failed") - .to_owned(); - cli_error_summary_with_default(&reason, "program") + let reason = reason.as_str().unwrap_or("task failed").to_owned(); + let mut summary = cli_error_summary_with_default(&reason, "program"); + if let Some(object) = summary.as_object_mut() { + if message_mentions_locality_failure(&reason.to_ascii_lowercase()) { + object.insert("locality_failure".to_owned(), json!(true)); + object.insert( + "next_actions".to_owned(), + json!(locality_failure_next_actions(&reason)), + ); + } + } + summary } Some("cancelled") => cli_error_summary_for_category("program", "task cancelled"), _ => Value::Null, } } +fn task_locality_failure_from_reason(reason: &Value) -> Value { + let Some(reason) = reason.as_str() else { + return Value::Null; + }; + let lower = reason.to_ascii_lowercase(); + if !message_mentions_locality_failure(&lower) { + return Value::Null; + } + let affected_data = if lower.contains("source snapshot") { + "source_snapshot" + } else if lower.contains("artifact") { + "artifact" + } else { + "direct_transfer" + }; + json!({ + "category": "connectivity", + "affected_data": affected_data, + "reason": reason, + "coordinator_bulk_relay_used": false, + "safe_failure": true, + "safe_next_actions": locality_failure_next_actions(reason), + }) +} + +fn locality_failure_next_actions(reason: &str) -> Vec<&'static str> { + let lower = reason.to_ascii_lowercase(); + if lower.contains("source snapshot") { + return vec![ + "attach or select a node that already has the required source snapshot", + "rerun source preparation on an attached node", + "restore direct node-to-node connectivity and retry", + "do not rely on coordinator bulk source relay", + ]; + } + if lower.contains("artifact") { + return vec![ + "attach or select a node that already has the required artifact", + "explicitly export or download the artifact before retrying", + "restore direct node-to-node connectivity and retry", + "do not rely on coordinator bulk artifact relay", + ]; + } + vec![ + "check node direct connectivity and NAT traversal", + "attach a node with the needed source/artifact locality", + "retry after node connectivity is restored", + "do not rely on coordinator bulk relay", + ] +} + fn process_state_from_tasks(task_events: Option<&Value>) -> &'static str { let events = task_event_values(task_events); if events.is_empty() { @@ -5346,6 +5455,11 @@ mod tests { "capability", 24, ), + ( + "scheduler placement failed: no capable node for placement: source snapshot unavailable and direct connectivity unavailable", + "connectivity", + 25, + ), ( "failed to connect to coordinator: connection refused", "connectivity", @@ -7372,7 +7486,8 @@ mod tests { let response = concat!( r#"{"type":"task_events","events":["#, r#"{"tenant":"tenant","project":"project","process":"vp","node":"node-a","task":"task-a","placement":{"node":"node-a","score":120,"reasons":["warm environment cache","source snapshot already local"]},"terminal_state":"completed","status_code":0,"stdout_bytes":12,"stderr_bytes":0,"stdout_tail":"ok","stderr_tail":"","stdout_truncated":false,"stderr_truncated":false,"artifact_path":"/vfs/artifacts/app.txt","artifact_digest":"sha256:artifact","artifact_size_bytes":12},"#, - r#"{"tenant":"tenant","project":"project","process":"vp","node":"node-b","task":"task-b","terminal_state":"failed","status_code":1,"stdout_bytes":0,"stderr_bytes":7,"stdout_tail":"","stderr_tail":"boom","stdout_truncated":false,"stderr_truncated":false,"artifact_path":null,"artifact_digest":null,"artifact_size_bytes":null}"#, + r#"{"tenant":"tenant","project":"project","process":"vp","node":"node-b","task":"task-b","terminal_state":"failed","status_code":1,"stdout_bytes":0,"stderr_bytes":7,"stdout_tail":"","stderr_tail":"boom","stdout_truncated":false,"stderr_truncated":false,"artifact_path":null,"artifact_digest":null,"artifact_size_bytes":null},"#, + r#"{"tenant":"tenant","project":"project","process":"vp","node":"node-c","task":"task-c","terminal_state":"failed","status_code":1,"stdout_bytes":0,"stderr_bytes":71,"stdout_tail":"","stderr_tail":"source snapshot unavailable and direct connectivity unavailable","stdout_truncated":false,"stderr_truncated":false,"artifact_path":null,"artifact_digest":null,"artifact_size_bytes":null}"#, r#"]}"# ); for _ in 0..4 { @@ -7418,7 +7533,7 @@ mod tests { server.join().unwrap(); assert_eq!(process["state"], "has_failed_tasks"); - assert_eq!(process["current_task_count"], 2); + assert_eq!(process["current_task_count"], 3); assert_eq!( process["current_tasks"][0]["node_placement"]["node"], "node-a" @@ -7434,6 +7549,27 @@ mod tests { assert_eq!(tasks["tasks"][1]["failure_reason"], "boom"); assert_eq!(tasks["tasks"][1]["machine_error"]["category"], "program"); assert_eq!(tasks["tasks"][1]["machine_error"]["stable_exit_code"], 27); + assert_eq!( + tasks["tasks"][2]["locality_failure"]["affected_data"], + "source_snapshot" + ); + assert_eq!( + tasks["tasks"][2]["locality_failure"]["coordinator_bulk_relay_used"], + false + ); + assert_eq!( + tasks["tasks"][2]["machine_error"]["category"], + "connectivity" + ); + assert_eq!(tasks["tasks"][2]["machine_error"]["stable_exit_code"], 25); + assert_eq!(tasks["tasks"][2]["machine_error"]["locality_failure"], true); + assert!(tasks["tasks"][2]["machine_error"]["next_actions"] + .as_array() + .unwrap() + .iter() + .any(|action| action == "rerun source preparation on an attached node")); + assert!(rendered_tasks.contains("locality task-c: source_snapshot")); + assert!(rendered_tasks.contains("do not rely on coordinator bulk source relay")); assert_eq!(logs["log_entries"].as_array().unwrap().len(), 1); assert_eq!(logs["log_entries"][0]["task"], "task-a"); assert_eq!(logs["log_entries"][0]["stdout_tail"], "ok"); diff --git a/scripts/cli-first-contract-smoke.js b/scripts/cli-first-contract-smoke.js index 1cd544e..bcbc4aa 100644 --- a/scripts/cli-first-contract-smoke.js +++ b/scripts/cli-first-contract-smoke.js @@ -87,6 +87,11 @@ expect( "artifact export explicit local byte write", /artifact export --to ` writes bytes[\s\S]*explicit bounded download stream[\s\S]*complete staged content is available/ ); +expect( + criteria, + "locality failure safe guidance", + /If direct transfer or locality assumptions fail[\s\S]*connectivity-category safe failures[\s\S]*coordinator bulk relay was not used/ +); expect( cliFirstAcceptance, "CLI-first acceptance report", @@ -198,11 +203,15 @@ for (const [name, pattern] of [ ["CLI command exit-code coverage", /fn command_report_exit_code_marks_command_failures_only\(\)/], ["CLI top-level logout alias coverage", /fn top_level_logout_alias_removes_only_cli_session_state\(\)/], ["CLI run rejection category coverage", /fn run_rejection_reports_machine_readable_error_category\(\)/], + ["CLI locality failure classifier", /fn classify_cli_error_message\(message: &str\)[\s\S]*message_mentions_locality_failure\(&message\)[\s\S]*return "connectivity"/], + ["CLI locality failure report helper", /fn task_locality_failure_from_reason\(reason: &Value\) -> Value[\s\S]*coordinator_bulk_relay_used[\s\S]*safe_next_actions/], + ["CLI locality failure human output", /fn push_task_locality_failures\(lines: &mut Vec, tasks: &\[Value\]\)[\s\S]*locality \{task_name\}/], ["node attach grant disclosure coverage", /fn node_attach_discloses_dangerous_capability_grants\(\)/], ["node enroll public API grant coverage", /fn node_enroll_reports_short_lived_public_api_grant\(\)/], ["quota local status coverage", /fn quota_status_uses_project_config_and_generic_public_limits\(\)/], ["quota coordinator usage coverage", /fn quota_status_queries_public_coordinator_usage\(\)/], ["task event summary coverage", /fn process_task_log_and_artifact_reports_summarize_task_events\(\)/], + ["task locality failure summary coverage", /fn process_task_log_and_artifact_reports_summarize_task_events\(\)[\s\S]*source snapshot unavailable and direct connectivity unavailable[\s\S]*locality_failure/], ["artifact download/export report coverage", /fn artifact_download_and_export_reports_expose_safe_session_boundaries\(\)/], ["process control report coverage", /fn process_restart_and_cancel_reports_expose_control_boundaries\(\)/], ["task restart report coverage", /fn task_restart_reports_clean_boundary_requirements\(\)/],