diff --git a/DISASMER_PUBLIC_TREE.json b/DISASMER_PUBLIC_TREE.json index ff503bb..a254a42 100644 --- a/DISASMER_PUBLIC_TREE.json +++ b/DISASMER_PUBLIC_TREE.json @@ -1,7 +1,7 @@ { "kind": "disasmer-filtered-public-tree", - "source_commit": "96f60dfba4756607edf60c08519050e2e4c2bf37", - "release_name": "dryrun-96f60dfba475", + "source_commit": "9b295127cfd7e4e87d4518a817563ca8e695ed0b", + "release_name": "dryrun-9b295127cfd7", "filtered_out": [ "private/**", "experiments/**", diff --git a/crates/disasmer-cli/src/main.rs b/crates/disasmer-cli/src/main.rs index a840bb4..641ae1e 100644 --- a/crates/disasmer-cli/src/main.rs +++ b/crates/disasmer-cli/src/main.rs @@ -1783,8 +1783,12 @@ fn human_report(value: &Value) -> String { if let Some(current_task_count) = value.get("current_task_count").and_then(Value::as_u64) { lines.push(format!("tasks: {current_task_count}")); } + if let Some(tasks) = value.get("current_tasks").and_then(Value::as_array) { + push_task_placement_reasons(&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); } if let Some(log_entries) = value.get("log_entries").and_then(Value::as_array) { lines.push(format!("log entries: {}", log_entries.len())); @@ -1976,6 +1980,33 @@ fn push_string_field(lines: &mut Vec, value: &Value, key: &str, label: & } } +fn push_task_placement_reasons(lines: &mut Vec, tasks: &[Value]) { + for task in tasks { + let Some(placement) = task.get("node_placement") else { + continue; + }; + let Some(reasons) = placement.get("reasons").and_then(Value::as_array) else { + continue; + }; + let reasons = reasons.iter().filter_map(Value::as_str).collect::>(); + if reasons.is_empty() { + continue; + } + let task_name = task + .get("task") + .and_then(Value::as_str) + .unwrap_or("unknown"); + let node = placement + .get("node") + .and_then(Value::as_str) + .unwrap_or("unknown"); + lines.push(format!( + "placement {task_name}: {node} ({})", + reasons.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}")); @@ -2514,6 +2545,15 @@ fn task_summaries(task_events: Option<&Value>) -> Value { let terminal_state = event_string(event, "terminal_state").unwrap_or_else(|| "unknown".to_owned()); let node = event_string(event, "node"); + let placement = event.get("placement"); + let placement_reasons = placement + .and_then(|placement| placement.get("reasons")) + .cloned() + .unwrap_or_else(|| json!([])); + let placement_score = placement + .and_then(|placement| placement.get("score")) + .cloned() + .unwrap_or(Value::Null); json!({ "process": event_string(event, "process"), "task": task, @@ -2523,6 +2563,9 @@ fn task_summaries(task_events: Option<&Value>) -> Value { "node_placement": { "node": node, "source": "coordinator_task_event", + "score": placement_score, + "reasons": placement_reasons, + "explanation_available": placement.is_some(), }, "failure_reason": task_failure_reason(event), "stdout_bytes": event_u64(event, "stdout_bytes").unwrap_or(0), @@ -5345,7 +5388,7 @@ mod tests { let server = std::thread::spawn(move || { let response = concat!( r#"{"type":"task_events","events":["#, - r#"{"tenant":"tenant","project":"project","process":"vp","node":"node-a","task":"task-a","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-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#"]}"# ); @@ -5397,6 +5440,14 @@ mod tests { process["current_tasks"][0]["node_placement"]["node"], "node-a" ); + assert_eq!( + process["current_tasks"][0]["node_placement"]["reasons"][0], + "warm environment cache" + ); + assert_eq!(tasks["tasks"][0]["node_placement"]["score"], 120); + let rendered_tasks = human_report(&tasks); + assert!(rendered_tasks.contains("placement task-a: node-a")); + assert!(rendered_tasks.contains("source snapshot already local")); assert_eq!(tasks["tasks"][1]["failure_reason"], "boom"); assert_eq!(logs["log_entries"].as_array().unwrap().len(), 1); assert_eq!(logs["log_entries"][0]["task"], "task-a"); diff --git a/crates/disasmer-coordinator/src/service.rs b/crates/disasmer-coordinator/src/service.rs index 4499b82..e403597 100644 --- a/crates/disasmer-coordinator/src/service.rs +++ b/crates/disasmer-coordinator/src/service.rs @@ -410,6 +410,8 @@ pub struct TaskCompletionEvent { pub process: ProcessId, pub node: NodeId, pub task: TaskId, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub placement: Option, pub terminal_state: TaskTerminalState, pub status_code: Option, pub stdout_bytes: u64, @@ -717,6 +719,7 @@ pub struct CoordinatorService { task_events: Vec, debug_audit_events: Vec, task_assignments: BTreeMap>, + task_placements: BTreeMap, active_tasks: BTreeSet, task_cancellations: BTreeSet, process_cancellations: BTreeSet, @@ -745,6 +748,7 @@ impl CoordinatorService { task_events: Vec::new(), debug_audit_events: Vec::new(), task_assignments: BTreeMap::new(), + task_placements: BTreeMap::new(), active_tasks: BTreeSet::new(), task_cancellations: BTreeSet::new(), process_cancellations: BTreeSet::new(), @@ -1199,6 +1203,10 @@ impl CoordinatorService { .retain(|(task_tenant, task_project, _, task_node, _)| { task_tenant != &tenant || task_project != &project || task_node != &node }); + self.task_placements + .retain(|(task_tenant, task_project, _, task_node, _), _| { + task_tenant != &tenant || task_project != &project || task_node != &node + }); self.coordinator.persist(&mut self.store); Ok(CoordinatorResponse::NodeCredentialRevoked { node, @@ -1325,13 +1333,10 @@ impl CoordinatorService { command_args, artifact_path, }; - self.active_tasks.insert(task_control_key( - &tenant, - &project, - &process, - &placement.node, - &task, - )); + let task_key = + task_control_key(&tenant, &project, &process, &placement.node, &task); + self.active_tasks.insert(task_key.clone()); + self.task_placements.insert(task_key, placement.clone()); self.task_assignments .entry((tenant, project, placement.node.clone())) .or_default() @@ -1520,6 +1525,13 @@ impl CoordinatorService { || task_project != &project || task_process != &process }); + self.task_placements.retain( + |(task_tenant, task_project, task_process, _, _), _| { + task_tenant != &tenant + || task_project != &project + || task_process != &process + }, + ); self.task_assignments.retain(|_, assignments| { assignments.retain(|assignment| { assignment.tenant != tenant @@ -1913,12 +1925,13 @@ impl CoordinatorService { .map(VfsPath::new) .transpose() .map_err(|err| CoordinatorServiceError::InvalidArtifactPath(err.to_string()))?; - let event = TaskCompletionEvent { + let mut event = TaskCompletionEvent { tenant: TenantId::new(tenant), project: ProjectId::new(project), process: ProcessId::new(process), node: NodeId::new(node), task: TaskId::new(task), + placement: None, terminal_state: terminal_state .unwrap_or_else(|| TaskTerminalState::from_status_code(status_code)), status_code, @@ -1938,6 +1951,14 @@ impl CoordinatorService { &event.project, &event.process, )?; + let task_key = task_control_key( + &event.tenant, + &event.project, + &event.process, + &event.node, + &event.task, + ); + event.placement = self.task_placements.remove(&task_key); if let (Some(path), Some(digest)) = (&event.artifact_path, artifact_digest) { self.artifact_registry.flush_metadata( artifact_id_from_path(path), @@ -1950,20 +1971,8 @@ impl CoordinatorService { artifact_size_bytes.unwrap_or(stdout_bytes), ); } - self.task_cancellations.remove(&task_control_key( - &event.tenant, - &event.project, - &event.process, - &event.node, - &event.task, - )); - self.active_tasks.remove(&task_control_key( - &event.tenant, - &event.project, - &event.process, - &event.node, - &event.task, - )); + self.task_cancellations.remove(&task_key); + self.active_tasks.remove(&task_key); self.task_events.push(event.clone()); Ok(CoordinatorResponse::TaskRecorded { process: event.process, @@ -4862,10 +4871,10 @@ mod tests { project: "project".to_owned(), node: "worker-linux".to_owned(), capabilities: linux_capabilities(), - cached_environment_digests: Vec::new(), - dependency_cache_digests: Vec::new(), - source_snapshots: Vec::new(), - artifact_locations: Vec::new(), + cached_environment_digests: vec![Digest::sha256("env")], + dependency_cache_digests: vec![Digest::sha256("deps")], + source_snapshots: vec![Digest::sha256("source")], + artifact_locations: vec!["bootstrap-artifact".to_owned()], direct_connectivity: true, online: true, }) @@ -4884,6 +4893,13 @@ mod tests { else { panic!("expected coordinator-side process start"); }; + service + .handle_request(CoordinatorRequest::ReconnectNode { + node: "worker-linux".to_owned(), + process: "vp-control".to_owned(), + epoch, + }) + .unwrap(); let CoordinatorResponse::TaskLaunched { process, @@ -4901,11 +4917,11 @@ mod tests { process: "vp-control".to_owned(), task: "compile-linux".to_owned(), environment: None, - environment_digest: None, + environment_digest: Some(Digest::sha256("env")), required_capabilities: vec![Capability::Command], - dependency_cache: None, - source_snapshot: None, - required_artifacts: Vec::new(), + dependency_cache: Some(Digest::sha256("deps")), + source_snapshot: Some(Digest::sha256("source")), + required_artifacts: vec!["bootstrap-artifact".to_owned()], quota_available: true, policy_allowed: true, command: "cargo".to_owned(), @@ -4919,6 +4935,14 @@ mod tests { assert_eq!(process, ProcessId::from("vp-control")); assert_eq!(task, TaskId::from("compile-linux")); assert_eq!(placement.node, NodeId::from("worker-linux")); + assert!(placement + .reasons + .iter() + .any(|reason| reason.contains("environment"))); + assert!(placement + .reasons + .iter() + .any(|reason| reason.contains("source"))); assert_eq!(assignment.node, NodeId::from("worker-linux")); assert_eq!(assignment.epoch, epoch); @@ -4949,6 +4973,44 @@ mod tests { panic!("expected empty worker assignment poll"); }; assert!(assignment.is_none()); + + service + .handle_request(CoordinatorRequest::TaskCompleted { + tenant: "tenant".to_owned(), + project: "project".to_owned(), + process: "vp-control".to_owned(), + node: "worker-linux".to_owned(), + task: "compile-linux".to_owned(), + terminal_state: Some(TaskTerminalState::Completed), + status_code: Some(0), + stdout_bytes: 12, + stderr_bytes: 0, + stdout_tail: "ok".to_owned(), + stderr_tail: String::new(), + stdout_truncated: false, + stderr_truncated: false, + artifact_path: Some("/vfs/artifacts/dap-output.txt".to_owned()), + artifact_digest: Some(Digest::sha256("artifact")), + artifact_size_bytes: Some(12), + }) + .unwrap(); + let CoordinatorResponse::TaskEvents { events } = service + .handle_request(CoordinatorRequest::ListTaskEvents { + tenant: "tenant".to_owned(), + project: "project".to_owned(), + actor_user: "user".to_owned(), + process: Some("vp-control".to_owned()), + }) + .unwrap() + else { + panic!("expected task events"); + }; + let event_placement = events[0] + .placement + .as_ref() + .expect("task event should retain launch placement explanation"); + assert_eq!(event_placement.node, NodeId::from("worker-linux")); + assert_eq!(event_placement.reasons, placement.reasons); } #[test] diff --git a/scripts/cli-first-contract-smoke.js b/scripts/cli-first-contract-smoke.js index c10c6fc..abe76bc 100644 --- a/scripts/cli-first-contract-smoke.js +++ b/scripts/cli-first-contract-smoke.js @@ -178,6 +178,11 @@ expect( "coordinator task restart boundary coverage", /fn service_reports_task_restart_boundary_through_public_api\(\)/ ); +expect( + coordinator, + "coordinator task events retain placement reasons", + /pub struct TaskCompletionEvent[\s\S]*placement: Option[\s\S]*task_placements[\s\S]*event\.placement = self\.task_placements\.remove/ +); expect( coordinator, "public debug operation audit event", @@ -203,6 +208,16 @@ expect( "CLI exposes node attach grant disclosures", /struct CapabilityGrantDisclosure[\s\S]*coordinator_policy_limited[\s\S]*fn capability_grant_disclosures/ ); +expect( + cli, + "CLI exposes task placement reasons", + /fn task_summaries[\s\S]*node_placement[\s\S]*reasons[\s\S]*explanation_available/ +); +expect( + cli, + "CLI renders task placement reasons", + /fn push_task_placement_reasons[\s\S]*placement \{task_name\}: \{node\}/ +); expect( cli, "CLI parses dangerous capability overrides",