| 1 | //! One tally of a workflow run's agents. |
| 2 | //! |
| 3 | //! Every surface that says "N of M finished" should read this function, so |
| 4 | //! there is one count, not one per reducer. It is pure: typed run events in, |
| 5 | //! counts and the agents that stopped out. A failed or cancelled agent is |
| 6 | //! never counted as finished. |
| 7 | |
| 8 | use std::collections::HashMap; |
| 9 | |
| 10 | use serde::{Deserialize, Serialize}; |
| 11 | |
| 12 | use super::{ |
| 13 | IrWorkflowRunStatus, TaskCompletion, WorkflowRunStatus, WorkflowUiEvent, WorkflowUiEventKind, |
| 14 | truncate_chars, |
| 15 | }; |
| 16 | |
| 17 | /// Why a workflow agent stopped without finishing, as a closed set the UI can |
| 18 | /// render without parsing prose. |
| 19 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 20 | #[serde(rename_all = "snake_case")] |
| 21 | pub(super) enum WorkflowTaskStopKind { |
| 22 | /// The agent's wall-time limit ran out. |
| 23 | WallTime, |
| 24 | /// The agent used its last allowed step. |
| 25 | Steps, |
| 26 | /// The agent itself failed (error result, failed delivery check, ...). |
| 27 | Agent, |
| 28 | /// The agent was cancelled. |
| 29 | Cancelled, |
| 30 | /// The agent finished, but its reply did not match the task's schema. |
| 31 | Schema, |
| 32 | /// The task was refused before any agent started. |
| 33 | Dispatch, |
| 34 | } |
| 35 | |
| 36 | impl WorkflowTaskStopKind { |
| 37 | fn phrase(self) -> &'static str { |
| 38 | match self { |
| 39 | Self::WallTime => "stopped at the time limit", |
| 40 | Self::Steps => "stopped at the step limit", |
| 41 | Self::Agent => "failed", |
| 42 | Self::Cancelled => "cancelled", |
| 43 | Self::Schema => "reply did not match the schema", |
| 44 | Self::Dispatch => "could not start", |
| 45 | } |
| 46 | } |
| 47 | } |
| 48 | |
| 49 | const REASON_MAX_CHARS: usize = 240; |
| 50 | |
| 51 | /// Reason and kind for a terminal child. A finished child has neither. |
| 52 | pub(super) fn task_stop( |
| 53 | completion: &TaskCompletion, |
| 54 | ) -> (Option<String>, Option<WorkflowTaskStopKind>) { |
| 55 | match completion { |
| 56 | TaskCompletion::Completed { .. } => (None, None), |
| 57 | TaskCompletion::Failed { message } => ( |
| 58 | Some(bounded_reason(message)), |
| 59 | Some(WorkflowTaskStopKind::Agent), |
| 60 | ), |
| 61 | TaskCompletion::Cancelled => (None, Some(WorkflowTaskStopKind::Cancelled)), |
| 62 | TaskCompletion::BudgetExhausted { message } => { |
| 63 | // The sub-agent loop names its two limits in the checkpoint |
| 64 | // reason; the same substring decides its own hand-back path. |
| 65 | let kind = if message.contains("wall-time") { |
| 66 | WorkflowTaskStopKind::WallTime |
| 67 | } else { |
| 68 | WorkflowTaskStopKind::Steps |
| 69 | }; |
| 70 | (Some(bounded_reason(message)), Some(kind)) |
| 71 | } |
| 72 | } |
| 73 | } |
| 74 | |
| 75 | fn bounded_reason(message: &str) -> String { |
| 76 | let first_line = message.lines().next().unwrap_or_default().trim(); |
| 77 | truncate_chars(first_line, REASON_MAX_CHARS) |
| 78 | } |
| 79 | |
| 80 | /// One agent that did not finish, in the order it stopped. |
| 81 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 82 | pub(super) struct StoppedAgent { |
| 83 | pub(super) label: String, |
| 84 | pub(super) kind: WorkflowTaskStopKind, |
| 85 | pub(super) reason: Option<String>, |
| 86 | } |
| 87 | |
| 88 | #[derive(Debug, Default, Clone, PartialEq, Eq)] |
| 89 | pub(super) struct WorkflowTally { |
| 90 | pub(super) queued: usize, |
| 91 | pub(super) running: usize, |
| 92 | pub(super) finished: usize, |
| 93 | pub(super) failed: usize, |
| 94 | pub(super) cancelled: usize, |
| 95 | pub(super) not_started: usize, |
| 96 | pub(super) stopped: Vec<StoppedAgent>, |
| 97 | } |
| 98 | |
| 99 | #[derive(Clone, Copy, PartialEq, Eq)] |
| 100 | enum RowState { |
| 101 | Running, |
| 102 | Finished, |
| 103 | Failed, |
| 104 | Cancelled, |
| 105 | } |
| 106 | |
| 107 | impl WorkflowTally { |
| 108 | pub(super) fn from_events(events: &[WorkflowUiEvent]) -> Self { |
| 109 | let mut tally = Self::default(); |
| 110 | let mut labels: HashMap<&str, String> = HashMap::new(); |
| 111 | let mut rows: Vec<(&str, RowState)> = Vec::new(); |
| 112 | let mut queued_tickets: Vec<u64> = Vec::new(); |
| 113 | for event in events { |
| 114 | match &event.kind { |
| 115 | WorkflowUiEventKind::TaskQueued { ticket, .. } => queued_tickets.push(*ticket), |
| 116 | WorkflowUiEventKind::TaskStarted(started) => { |
| 117 | if let Some(ticket) = started.queue_ticket { |
| 118 | queued_tickets.retain(|queued| *queued != ticket); |
| 119 | } |
| 120 | let label = started |
| 121 | .workflow_task_label |
| 122 | .clone() |
| 123 | .or_else(|| started.label.clone()) |
| 124 | .unwrap_or_else(|| started.task_id.clone()); |
| 125 | labels.insert(started.task_id.as_str(), label); |
| 126 | rows.push((started.task_id.as_str(), RowState::Running)); |
| 127 | } |
| 128 | WorkflowUiEventKind::TaskCompleted { |
| 129 | task_id, |
| 130 | status, |
| 131 | reason, |
| 132 | kind, |
| 133 | .. |
| 134 | } => { |
| 135 | let state = row_state(*status); |
| 136 | set_row(&mut rows, task_id, state); |
| 137 | if matches!(state, RowState::Failed | RowState::Cancelled) { |
| 138 | tally.stopped.push(StoppedAgent { |
| 139 | label: label_for(&labels, task_id), |
| 140 | kind: kind.unwrap_or(if state == RowState::Cancelled { |
| 141 | WorkflowTaskStopKind::Cancelled |
| 142 | } else { |
| 143 | WorkflowTaskStopKind::Agent |
| 144 | }), |
| 145 | reason: reason.clone(), |
| 146 | }); |
| 147 | } |
| 148 | } |
| 149 | // The reply is decoded after the child settles, so a schema |
| 150 | // failure turns an already-finished row into a failed one. |
| 151 | WorkflowUiEventKind::TaskSchemaValidationFailed { task_id, message } => { |
| 152 | if set_row(&mut rows, task_id, RowState::Failed) { |
| 153 | tally.stopped.push(StoppedAgent { |
| 154 | label: label_for(&labels, task_id), |
| 155 | kind: WorkflowTaskStopKind::Schema, |
| 156 | reason: Some(bounded_reason(message)), |
| 157 | }); |
| 158 | } |
| 159 | } |
| 160 | WorkflowUiEventKind::TaskDispatchFailed { |
| 161 | label, |
| 162 | message, |
| 163 | queue_ticket, |
| 164 | .. |
| 165 | } => { |
| 166 | if let Some(ticket) = queue_ticket { |
| 167 | queued_tickets.retain(|queued| queued != ticket); |
| 168 | } |
| 169 | tally.not_started += 1; |
| 170 | tally.stopped.push(StoppedAgent { |
| 171 | label: label.clone().unwrap_or_else(|| "task".to_string()), |
| 172 | kind: WorkflowTaskStopKind::Dispatch, |
| 173 | reason: Some(bounded_reason(message)), |
| 174 | }); |
| 175 | } |
| 176 | // A settled run has nothing left waiting: a wait it cut short |
| 177 | // never became an agent. |
| 178 | WorkflowUiEventKind::RunCompleted { .. } |
| 179 | | WorkflowUiEventKind::RunCancelled { .. } => queued_tickets.clear(), |
| 180 | _ => {} |
| 181 | } |
| 182 | } |
| 183 | for (_, state) in rows { |
| 184 | match state { |
| 185 | RowState::Running => tally.running += 1, |
| 186 | RowState::Finished => tally.finished += 1, |
| 187 | RowState::Failed => tally.failed += 1, |
| 188 | RowState::Cancelled => tally.cancelled += 1, |
| 189 | } |
| 190 | } |
| 191 | tally.queued = queued_tickets.len(); |
| 192 | tally |
| 193 | } |
| 194 | |
| 195 | /// Recount agents from the driver's own per-task ledger. A run keeps only |
| 196 | /// its newest events, so once older ones were trimmed the events alone |
| 197 | /// under-count; the ledger and the exact dispatch-failure total do not. |
| 198 | /// `stopped` stays the named sample the retained events carry. |
| 199 | pub(super) fn count_from_ledger( |
| 200 | &mut self, |
| 201 | statuses: impl IntoIterator<Item = IrWorkflowRunStatus>, |
| 202 | dispatch_failures: u64, |
| 203 | ) { |
| 204 | let (mut running, mut finished, mut failed, mut cancelled) = (0, 0, 0, 0); |
| 205 | for status in statuses { |
| 206 | match row_state(status) { |
| 207 | RowState::Running => running += 1, |
| 208 | RowState::Finished => finished += 1, |
| 209 | RowState::Failed => failed += 1, |
| 210 | RowState::Cancelled => cancelled += 1, |
| 211 | } |
| 212 | } |
| 213 | if running + finished + failed + cancelled == 0 { |
| 214 | return; |
| 215 | } |
| 216 | self.running = running; |
| 217 | self.finished = finished; |
| 218 | self.failed = failed; |
| 219 | self.cancelled = cancelled; |
| 220 | self.not_started = self |
| 221 | .not_started |
| 222 | .max(usize::try_from(dispatch_failures).unwrap_or(usize::MAX)); |
| 223 | } |
| 224 | |
| 225 | /// Agents that started, plus tasks refused before they could. |
| 226 | pub(super) fn total(&self) -> usize { |
| 227 | self.running + self.finished + self.failed + self.cancelled + self.not_started |
| 228 | } |
| 229 | |
| 230 | /// "3 of 5 agents finished, 2 failed (evals: stopped at the step limit)". |
| 231 | /// Zero counts are left out. |
| 232 | pub(super) fn agents_sentence(&self) -> String { |
| 233 | let total = self.total(); |
| 234 | if total == 0 && self.queued == 0 { |
| 235 | return "no agents ran".to_string(); |
| 236 | } |
| 237 | let noun = if total == 1 { "agent" } else { "agents" }; |
| 238 | let mut sentence = format!("{} of {total} {noun} finished", self.finished); |
| 239 | for (count, word) in [ |
| 240 | (self.failed, "failed"), |
| 241 | (self.cancelled, "cancelled"), |
| 242 | (self.not_started, "could not start"), |
| 243 | (self.running, "still running"), |
| 244 | (self.queued, "still waiting for a slot"), |
| 245 | ] { |
| 246 | if count > 0 { |
| 247 | sentence.push_str(&format!(", {count} {word}")); |
| 248 | } |
| 249 | } |
| 250 | const SHOWN: usize = 4; |
| 251 | if !self.stopped.is_empty() { |
| 252 | let mut details: Vec<String> = self |
| 253 | .stopped |
| 254 | .iter() |
| 255 | .take(SHOWN) |
| 256 | .map(|agent| { |
| 257 | let label = truncate_chars(&agent.label, 48); |
| 258 | match (agent.kind, agent.reason.as_deref()) { |
| 259 | ( |
| 260 | WorkflowTaskStopKind::Agent | WorkflowTaskStopKind::Dispatch, |
| 261 | Some(reason), |
| 262 | ) => { |
| 263 | format!( |
| 264 | "{label}: {}: {}", |
| 265 | agent.kind.phrase(), |
| 266 | truncate_chars(reason, 80) |
| 267 | ) |
| 268 | } |
| 269 | (kind, _) => format!("{label}: {}", kind.phrase()), |
| 270 | } |
| 271 | }) |
| 272 | .collect(); |
| 273 | // The named list can be a sample of a trimmed run; the counts |
| 274 | // are the whole run, so "more" is measured against them. |
| 275 | let stopped_total = self |
| 276 | .stopped |
| 277 | .len() |
| 278 | .max(self.failed + self.cancelled + self.not_started); |
| 279 | if stopped_total > details.len() { |
| 280 | details.push(format!("{} more", stopped_total - details.len())); |
| 281 | } |
| 282 | sentence.push_str(&format!(" ({})", details.join("; "))); |
| 283 | } |
| 284 | sentence |
| 285 | } |
| 286 | } |
| 287 | |
| 288 | fn row_state(status: IrWorkflowRunStatus) -> RowState { |
| 289 | match status { |
| 290 | IrWorkflowRunStatus::Succeeded => RowState::Finished, |
| 291 | IrWorkflowRunStatus::Cancelled => RowState::Cancelled, |
| 292 | IrWorkflowRunStatus::Pending | IrWorkflowRunStatus::Running => RowState::Running, |
| 293 | IrWorkflowRunStatus::Failed | IrWorkflowRunStatus::BudgetExceeded => RowState::Failed, |
| 294 | } |
| 295 | } |
| 296 | |
| 297 | /// Set a row's state, adding the row when its `task_started` event was |
| 298 | /// trimmed from the retained tail. True when the row is new or changed. |
| 299 | fn set_row<'a>(rows: &mut Vec<(&'a str, RowState)>, task_id: &'a str, state: RowState) -> bool { |
| 300 | match rows.iter_mut().find(|(id, _)| *id == task_id) { |
| 301 | Some(row) => { |
| 302 | let changed = row.1 != state; |
| 303 | row.1 = state; |
| 304 | changed |
| 305 | } |
| 306 | None => { |
| 307 | rows.push((task_id, state)); |
| 308 | true |
| 309 | } |
| 310 | } |
| 311 | } |
| 312 | |
| 313 | fn label_for(labels: &HashMap<&str, String>, task_id: &str) -> String { |
| 314 | labels |
| 315 | .get(task_id) |
| 316 | .cloned() |
| 317 | .unwrap_or_else(|| task_id.to_string()) |
| 318 | } |
| 319 | |
| 320 | /// Plain word for a run's end state. Degraded is not a failure: the script |
| 321 | /// returned, with gaps. |
| 322 | pub(super) fn run_outcome_phrase(status: WorkflowRunStatus) -> &'static str { |
| 323 | match status { |
| 324 | WorkflowRunStatus::Running => "is still running", |
| 325 | WorkflowRunStatus::Completed => "finished", |
| 326 | WorkflowRunStatus::Degraded => "finished with gaps", |
| 327 | WorkflowRunStatus::Failed => "failed", |
| 328 | WorkflowRunStatus::Cancelled => "was stopped", |
| 329 | } |
| 330 | } |
| 331 | |
| 332 | #[cfg(test)] |
| 333 | mod tests { |
| 334 | use super::super::WorkflowTaskStartedEvent; |
| 335 | use super::*; |
| 336 | |
| 337 | fn started(task_id: &str, label: &str, phase: &str, ticket: Option<u64>) -> WorkflowUiEvent { |
| 338 | WorkflowUiEvent::at( |
| 339 | 1, |
| 340 | "session", |
| 341 | WorkflowUiEventKind::TaskStarted(Box::new(WorkflowTaskStartedEvent { |
| 342 | task_id: task_id.to_string(), |
| 343 | label: None, |
| 344 | role: None, |
| 345 | profile: None, |
| 346 | model: None, |
| 347 | strength: None, |
| 348 | thinking: None, |
| 349 | requested_reasoning: None, |
| 350 | effective_reasoning: None, |
| 351 | resolved_role: None, |
| 352 | resolved_profile: None, |
| 353 | resolved_provider: "local".to_string(), |
| 354 | resolved_model: "stub".to_string(), |
| 355 | route_source: "session".to_string(), |
| 356 | child_route: None, |
| 357 | worktree: false, |
| 358 | workspace: None, |
| 359 | git_branch: None, |
| 360 | parent_task_id: None, |
| 361 | depth: 1, |
| 362 | workflow_run_id: Some("workflow_view".to_string()), |
| 363 | workflow_phase_id: Some(phase.to_string()), |
| 364 | workflow_task_label: Some(label.to_string()), |
| 365 | workflow_child_index: Some(0), |
| 366 | queue_ticket: ticket, |
| 367 | fleet_receipt: None, |
| 368 | })), |
| 369 | ) |
| 370 | } |
| 371 | |
| 372 | fn completed(task_id: &str, completion: TaskCompletion) -> WorkflowUiEvent { |
| 373 | let status = match completion { |
| 374 | TaskCompletion::Completed { .. } => IrWorkflowRunStatus::Succeeded, |
| 375 | TaskCompletion::Failed { .. } => IrWorkflowRunStatus::Failed, |
| 376 | TaskCompletion::Cancelled => IrWorkflowRunStatus::Cancelled, |
| 377 | TaskCompletion::BudgetExhausted { .. } => IrWorkflowRunStatus::BudgetExceeded, |
| 378 | }; |
| 379 | let (reason, kind) = task_stop(&completion); |
| 380 | WorkflowUiEvent::at( |
| 381 | 2, |
| 382 | "session", |
| 383 | WorkflowUiEventKind::TaskCompleted { |
| 384 | task_id: task_id.to_string(), |
| 385 | status, |
| 386 | reason, |
| 387 | kind, |
| 388 | usage: None, |
| 389 | }, |
| 390 | ) |
| 391 | } |
| 392 | |
| 393 | #[test] |
| 394 | fn stop_kind_names_the_limit_that_ended_the_agent() { |
| 395 | let steps = TaskCompletion::BudgetExhausted { |
| 396 | message: "child step budget exhausted for task execution (limit: 40; used: 40; any remaining turn is reserved for hand-back)".to_string(), |
| 397 | }; |
| 398 | let wall = TaskCompletion::BudgetExhausted { |
| 399 | message: "child wall-time budget exhausted during task execution; remaining time is reserved for hand-back.".to_string(), |
| 400 | }; |
| 401 | assert_eq!(task_stop(&steps).1, Some(WorkflowTaskStopKind::Steps)); |
| 402 | assert!(task_stop(&steps).0.unwrap().contains("limit: 40")); |
| 403 | assert_eq!(task_stop(&wall).1, Some(WorkflowTaskStopKind::WallTime)); |
| 404 | assert_eq!( |
| 405 | task_stop(&TaskCompletion::Completed { |
| 406 | text: "done".to_string() |
| 407 | }), |
| 408 | (None, None) |
| 409 | ); |
| 410 | let long = TaskCompletion::Failed { |
| 411 | message: format!("{}\nstack line", "x".repeat(1_000)), |
| 412 | }; |
| 413 | let (reason, kind) = task_stop(&long); |
| 414 | assert_eq!(kind, Some(WorkflowTaskStopKind::Agent)); |
| 415 | let reason = reason.unwrap(); |
| 416 | assert!(reason.chars().count() <= REASON_MAX_CHARS); |
| 417 | assert!(!reason.contains("stack line"), "one line, never the stack"); |
| 418 | } |
| 419 | |
| 420 | #[test] |
| 421 | fn failed_and_cancelled_agents_are_never_counted_as_finished() { |
| 422 | // Real run 511c2203's shape: every agent died at its limit. |
| 423 | let events: Vec<WorkflowUiEvent> = (0..4) |
| 424 | .flat_map(|index| { |
| 425 | let id = format!("agent_{index}"); |
| 426 | [ |
| 427 | started(&id, &format!("slot-{index}"), "survey", None), |
| 428 | completed( |
| 429 | &id, |
| 430 | TaskCompletion::BudgetExhausted { |
| 431 | message: "child step budget exhausted".to_string(), |
| 432 | }, |
| 433 | ), |
| 434 | ] |
| 435 | }) |
| 436 | .collect(); |
| 437 | let tally = WorkflowTally::from_events(&events); |
| 438 | assert_eq!((tally.finished, tally.failed, tally.total()), (0, 4, 4)); |
| 439 | let sentence = tally.agents_sentence(); |
| 440 | assert!( |
| 441 | sentence.starts_with("0 of 4 agents finished, 4 failed"), |
| 442 | "{sentence}" |
| 443 | ); |
| 444 | assert!(sentence.contains("slot-0: stopped at the step limit")); |
| 445 | assert_eq!(sentence.matches("stopped at the step limit").count(), 4); |
| 446 | assert!(!sentence.contains("more"), "{sentence}"); |
| 447 | } |
| 448 | |
| 449 | #[test] |
| 450 | fn tally_reads_schema_failures_dispatch_refusals_and_queued_tasks() { |
| 451 | let queued = |ticket: u64, label: &str| { |
| 452 | WorkflowUiEvent::at( |
| 453 | 0, |
| 454 | "session", |
| 455 | WorkflowUiEventKind::TaskQueued { |
| 456 | ticket, |
| 457 | label: Some(label.to_string()), |
| 458 | phase: None, |
| 459 | }, |
| 460 | ) |
| 461 | }; |
| 462 | let mut events = vec![ |
| 463 | queued(0, "tools-approval"), |
| 464 | queued(2, "evals"), |
| 465 | started("a", "loop-prompt", "survey", None), |
| 466 | completed( |
| 467 | "a", |
| 468 | TaskCompletion::Completed { |
| 469 | text: "ok".to_string(), |
| 470 | }, |
| 471 | ), |
| 472 | started("b", "tools-approval", "survey", Some(0)), |
| 473 | completed( |
| 474 | "b", |
| 475 | TaskCompletion::Completed { |
| 476 | text: "not json".to_string(), |
| 477 | }, |
| 478 | ), |
| 479 | WorkflowUiEvent::at( |
| 480 | 3, |
| 481 | "session", |
| 482 | WorkflowUiEventKind::TaskSchemaValidationFailed { |
| 483 | task_id: "b".to_string(), |
| 484 | message: "expected object".to_string(), |
| 485 | }, |
| 486 | ), |
| 487 | WorkflowUiEvent::at( |
| 488 | 4, |
| 489 | "session", |
| 490 | WorkflowUiEventKind::TaskDispatchFailed { |
| 491 | label: Some("evals".to_string()), |
| 492 | phase: None, |
| 493 | message: "cwd is outside the workspace".to_string(), |
| 494 | queue_ticket: Some(2), |
| 495 | }, |
| 496 | ), |
| 497 | queued(1, "surfaces"), |
| 498 | started("c", "models-context", "survey", None), |
| 499 | completed("c", TaskCompletion::Cancelled), |
| 500 | ]; |
| 501 | let tally = WorkflowTally::from_events(&events); |
| 502 | assert_eq!(tally.finished, 1); |
| 503 | assert_eq!(tally.failed, 1); |
| 504 | assert_eq!(tally.cancelled, 1); |
| 505 | assert_eq!(tally.not_started, 1); |
| 506 | assert_eq!( |
| 507 | tally.queued, 1, |
| 508 | "ticket 0 started and ticket 2 was refused; ticket 1 still waits" |
| 509 | ); |
| 510 | assert_eq!(tally.total(), 4); |
| 511 | let sentence = tally.agents_sentence(); |
| 512 | assert_eq!( |
| 513 | sentence, |
| 514 | "1 of 4 agents finished, 1 failed, 1 cancelled, 1 could not start, 1 still waiting for a slot \ |
| 515 | (tools-approval: reply did not match the schema; evals: could not start: cwd is outside the workspace; \ |
| 516 | models-context: cancelled)" |
| 517 | ); |
| 518 | |
| 519 | // A settled run has nothing waiting; the cut-short wait is not an agent. |
| 520 | events.push(WorkflowUiEvent::at( |
| 521 | 6, |
| 522 | "session", |
| 523 | WorkflowUiEventKind::RunCancelled { |
| 524 | reason: "stopped".to_string(), |
| 525 | }, |
| 526 | )); |
| 527 | let settled = WorkflowTally::from_events(&events); |
| 528 | assert_eq!((settled.queued, settled.total()), (0, 4)); |
| 529 | } |
| 530 | |
| 531 | #[test] |
| 532 | fn empty_run_says_no_agents_ran() { |
| 533 | assert_eq!( |
| 534 | WorkflowTally::from_events(&[]).agents_sentence(), |
| 535 | "no agents ran" |
| 536 | ); |
| 537 | } |
| 538 | } |
| 539 |