Expose task placement reasons in CLI summaries

This commit is contained in:
Michel Paulissen 2026-07-03 22:15:11 +02:00
parent fe6c2d5042
commit 218ab65e47
4 changed files with 161 additions and 33 deletions

View file

@ -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<Placement>,
pub terminal_state: TaskTerminalState,
pub status_code: Option<i32>,
pub stdout_bytes: u64,
@ -717,6 +719,7 @@ pub struct CoordinatorService {
task_events: Vec<TaskCompletionEvent>,
debug_audit_events: Vec<DebugAuditEvent>,
task_assignments: BTreeMap<TaskAssignmentKey, VecDeque<TaskAssignment>>,
task_placements: BTreeMap<TaskControlKey, Placement>,
active_tasks: BTreeSet<TaskControlKey>,
task_cancellations: BTreeSet<TaskControlKey>,
process_cancellations: BTreeSet<ProcessControlKey>,
@ -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]