返回 CodeWhale
journal.rs
根目录 / crates / tui / src / tools / workflow / journal.rs
1 use super::{
2 SharedWorkflowControllers, SharedWorkflowLifecycles, SharedWorkflowRuns,
3 WorkflowDispatchFailure, WorkflowRunRecord, WorkflowRunStatus, WorkflowUiEvent,
4 WorkflowUiEventKind, WorkflowWorkLifecycle,
5 };
6 use serde::{Deserialize, Serialize};
7 use std::collections::HashMap;
8 use std::io::{BufRead, Write};
9 use std::path::{Path, PathBuf};
10 use std::sync::{Arc, Mutex, OnceLock};
11 use tracing::warn;
12
13 pub(super) const CODEWHALE_DIR: &str = ".codewhale";
14 pub(super) const WORKFLOW_RUNS_FILE: &str = "workflow-runs.jsonl";
15
16 /// Per-workspace workflow state shared across tool-registry rebuilds.
17 pub(super) struct WorkflowWorkspaceState {
18 pub runs: SharedWorkflowRuns,
19 pub controllers: SharedWorkflowControllers,
20 lifecycles: SharedWorkflowLifecycles,
21 journal: WorkflowRunJournal,
22 }
23
24 impl WorkflowWorkspaceState {
25 pub fn open(workspace: &Path) -> Arc<Self> {
26 Self::open_inner(workspace, true)
27 }
28
29 /// Hydrate the journal without rewriting leftover `running` rows to
30 /// `failed`. Host cancel uses this after a restart so a controller-less
31 /// run can still be marked cancelled instead of looking like a crash.
32 pub fn open_preserving_running(workspace: &Path) -> Arc<Self> {
33 Self::open_inner(workspace, false)
34 }
35
36 fn open_inner(workspace: &Path, recover_orphans: bool) -> Arc<Self> {
37 let journal = WorkflowRunJournal::open(workspace);
38 let runs = Arc::new(Mutex::new(journal.hydrate_runs(recover_orphans)));
39 Arc::new(Self {
40 runs,
41 controllers: Arc::new(Mutex::new(HashMap::new())),
42 lifecycles: Arc::new(Mutex::new(HashMap::new())),
43 journal,
44 })
45 }
46
47 pub fn attach_lifecycle(&self, run_id: &str, lifecycle: WorkflowWorkLifecycle) {
48 self.lifecycles
49 .lock()
50 .unwrap_or_else(|poison| poison.into_inner())
51 .entry(run_id.to_string())
52 .or_insert(lifecycle);
53 }
54
55 pub fn reconcile_snapshot(&self, record: &WorkflowRunRecord) {
56 let lifecycle = self
57 .lifecycles
58 .lock()
59 .unwrap_or_else(|poison| poison.into_inner())
60 .get(&record.run_id)
61 .cloned();
62 if let Some(lifecycle) = lifecycle
63 && let Err(err) = lifecycle.reconcile_record(record)
64 {
65 warn!(
66 run_id = record.run_id,
67 "workflow Work reconciliation failed: {err}"
68 );
69 }
70 }
71
72 pub fn reconcile_cancel(&self, run_id: &str, outcome: super::CancelOutcome) {
73 let lifecycle = self
74 .lifecycles
75 .lock()
76 .unwrap_or_else(|poison| poison.into_inner())
77 .get(run_id)
78 .cloned();
79 if let Some(lifecycle) = lifecycle
80 && let Err(err) = lifecycle.reconcile_cancel(outcome)
81 {
82 warn!(run_id, "workflow cancellation reconciliation failed: {err}");
83 }
84 }
85
86 pub fn mark_owner_missing(&self, run_id: &str) {
87 let lifecycle = self
88 .lifecycles
89 .lock()
90 .unwrap_or_else(|poison| poison.into_inner())
91 .get(run_id)
92 .cloned();
93 if let Some(lifecycle) = lifecycle {
94 lifecycle.reconcile_missing();
95 }
96 }
97
98 pub fn try_record_snapshot(&self, record: &WorkflowRunRecord) -> Result<(), String> {
99 self.journal
100 .append_snapshot(record)
101 .map_err(|err| err.to_string())
102 }
103
104 pub fn record_snapshot(&self, record: &WorkflowRunRecord) {
105 if let Err(err) = self.try_record_snapshot(record) {
106 warn!("workflow journal snapshot failed: {err}");
107 }
108 }
109
110 pub fn record_progress(&self, run_id: &str, message: &str) {
111 if let Err(err) = self.journal.append_progress(run_id, message) {
112 warn!("workflow journal progress failed: {err}");
113 }
114 }
115
116 pub fn record_event(&self, run_id: &str, event: &WorkflowUiEvent) {
117 if let Err(err) = self.journal.append_event(run_id, event) {
118 warn!("workflow journal event failed: {err}");
119 }
120 }
121
122 /// Durable journal location for full-fidelity run detail (#2974).
123 pub fn journal_path(&self) -> &Path {
124 &self.journal.ledger_path
125 }
126 }
127
128 fn workspace_store() -> &'static Mutex<HashMap<PathBuf, Arc<WorkflowWorkspaceState>>> {
129 static STORE: OnceLock<Mutex<HashMap<PathBuf, Arc<WorkflowWorkspaceState>>>> = OnceLock::new();
130 STORE.get_or_init(|| Mutex::new(HashMap::new()))
131 }
132
133 pub(super) fn shared_workflow_state(workspace: &Path) -> Arc<WorkflowWorkspaceState> {
134 let key = workspace
135 .canonicalize()
136 .unwrap_or_else(|_| workspace.to_path_buf());
137 let mut store = workspace_store()
138 .lock()
139 .unwrap_or_else(|poison| poison.into_inner());
140 store
141 .entry(key)
142 .or_insert_with(|| WorkflowWorkspaceState::open(workspace))
143 .clone()
144 }
145
146 /// Read-only lookup that never creates workspace state, a journal
147 /// directory, or a ledger file. Used by the human-only `/structcopy`
148 /// command (#2033), which must stay side-effect free.
149 pub(super) fn peek_shared_workflow_state(workspace: &Path) -> Option<Arc<WorkflowWorkspaceState>> {
150 let key = workspace
151 .canonicalize()
152 .unwrap_or_else(|_| workspace.to_path_buf());
153 workspace_store()
154 .lock()
155 .unwrap_or_else(|poison| poison.into_inner())
156 .get(&key)
157 .cloned()
158 }
159
160 #[derive(Debug, Clone, Serialize, Deserialize)]
161 #[serde(tag = "kind", rename_all = "snake_case")]
162 enum WorkflowJournalRecord {
163 // Boxed: a full run record dwarfs the progress variant
164 // (clippy::large_enum_variant).
165 Snapshot {
166 run: Box<WorkflowRunRecord>,
167 },
168 Progress {
169 run_id: String,
170 message: String,
171 },
172 Event {
173 run_id: String,
174 event: Box<WorkflowUiEvent>,
175 },
176 }
177
178 #[derive(Debug)]
179 struct WorkflowRunJournal {
180 /// Workspace the ledger must stay inside; links below it are refused.
181 root: PathBuf,
182 ledger_path: PathBuf,
183 }
184
185 impl WorkflowRunJournal {
186 /// Opening has no side effects: `.codewhale/` and the ledger appear with
187 /// the first appended record, and only through the confined helpers, so a
188 /// linked `.codewhale` never receives an empty `state/` or ledger.
189 fn open(workspace: &Path) -> Self {
190 Self {
191 root: workspace.to_path_buf(),
192 ledger_path: workspace.join(CODEWHALE_DIR).join(WORKFLOW_RUNS_FILE),
193 }
194 }
195
196 fn hydrate_runs(&self, recover_orphans: bool) -> HashMap<String, WorkflowRunRecord> {
197 let file = match crate::fs_confined::open_read(&self.root, &self.ledger_path) {
198 Ok(file) => file,
199 Err(err) => {
200 if err.kind() != std::io::ErrorKind::NotFound {
201 warn!(
202 "workflow journal unreadable ({}): {err}",
203 self.ledger_path.display()
204 );
205 }
206 return HashMap::new();
207 }
208 };
209 let mut runs = HashMap::new();
210 for line in std::io::BufReader::new(file).lines() {
211 let Ok(line) = line else { continue };
212 let trimmed = line.trim();
213 if trimmed.is_empty() {
214 continue;
215 }
216 let record = match serde_json::from_str::<WorkflowJournalRecord>(trimmed) {
217 Ok(record) => record,
218 Err(err) => {
219 warn!("workflow journal skipped malformed line: {err}");
220 continue;
221 }
222 };
223 match record {
224 WorkflowJournalRecord::Snapshot { run } => {
225 let mut run = *run;
226 run.normalize_bounded_ledgers();
227 runs.insert(run.run_id.clone(), run);
228 }
229 WorkflowJournalRecord::Progress { run_id, message } => {
230 if let Some(run) = runs.get_mut(&run_id) {
231 run.push_progress(message);
232 }
233 }
234 WorkflowJournalRecord::Event { run_id, event } => {
235 if let Some(run) = runs.get_mut(&run_id) {
236 let event = *event;
237 if let WorkflowUiEventKind::TaskDispatchFailed {
238 label,
239 phase,
240 message,
241 ..
242 } = &event.kind
243 {
244 run.push_dispatch_failure(WorkflowDispatchFailure {
245 at_ms: event.at_ms,
246 label: label.clone(),
247 phase: phase.clone(),
248 message: message.clone(),
249 });
250 }
251 run.push_event(event);
252 }
253 }
254 }
255 }
256 // Journals written before #2974 have no counters; rebuild them
257 // from the retained tail so summaries stay truthful.
258 for run in runs.values_mut() {
259 run.normalize_bounded_ledgers();
260 run.events_total = run.events_total.max(run.events.len() as u64);
261 }
262 // A run journaled as Running belongs to a process that is gone;
263 // without this it would show as live forever after a restart.
264 // Host cancel skips this rewrite so it can still mark the line
265 // cancelled with an honest "nothing live to stop" receipt.
266 if recover_orphans {
267 let mut recovered = Vec::new();
268 for run in runs.values_mut() {
269 if run.status == WorkflowRunStatus::Running {
270 run.status = WorkflowRunStatus::Failed;
271 run.lifecycle_seq = run.lifecycle_seq.saturating_add(1);
272 run.completed_at_ms.get_or_insert_with(super::now_ms);
273 run.error = Some(
274 "process exited before the run completed (recovered on startup)"
275 .to_string(),
276 );
277 recovered.push(run.clone());
278 }
279 }
280 // The recovery decision is owner truth, not a presentation-only
281 // repair. Append it so another restart replays the same terminal
282 // sequence instead of rediscovering and incrementing it again.
283 for run in recovered {
284 if let Err(err) = self.append_snapshot(&run) {
285 warn!(
286 run_id = run.run_id,
287 "workflow recovery snapshot append failed: {err}"
288 );
289 }
290 }
291 }
292 runs
293 }
294
295 fn append_record(&self, record: &WorkflowJournalRecord) -> std::io::Result<()> {
296 let mut line =
297 serde_json::to_string(record).map_err(|err| std::io::Error::other(err.to_string()))?;
298 line.push('\n');
299 let mut file = crate::fs_confined::open_append(&self.root, &self.ledger_path)?;
300 file.write_all(line.as_bytes())?;
301 file.flush()?;
302 Ok(())
303 }
304
305 fn append_snapshot(&self, record: &WorkflowRunRecord) -> std::io::Result<()> {
306 self.append_record(&WorkflowJournalRecord::Snapshot {
307 run: Box::new(record.clone()),
308 })
309 }
310
311 fn append_progress(&self, run_id: &str, message: &str) -> std::io::Result<()> {
312 self.append_record(&WorkflowJournalRecord::Progress {
313 run_id: run_id.to_string(),
314 message: message.to_string(),
315 })
316 }
317
318 fn append_event(&self, run_id: &str, event: &WorkflowUiEvent) -> std::io::Result<()> {
319 self.append_record(&WorkflowJournalRecord::Event {
320 run_id: run_id.to_string(),
321 event: Box::new(event.clone()),
322 })
323 }
324 }
325
326 #[cfg(test)]
327 mod tests {
328 use super::super::{WORKFLOW_RUN_DISPATCH_FAILURES_MAX_RETAINED, WorkflowUiEventKind};
329 use super::*;
330
331 /// #5582: a degraded run must never project as an ordinary success
332 /// to owner-level consumers.
333 #[test]
334 fn owner_snapshot_keeps_degraded_distinct_from_completed() {
335 use crate::work_graph::OwnerState;
336 assert_eq!(
337 crate::tools::workflow::owner_state_for_run_status(WorkflowRunStatus::Degraded),
338 OwnerState::Degraded
339 );
340 assert_eq!(
341 crate::tools::workflow::owner_state_for_run_status(WorkflowRunStatus::Completed),
342 OwnerState::Completed
343 );
344 }
345
346 fn sample_record(run_id: &str, status: WorkflowRunStatus) -> WorkflowRunRecord {
347 WorkflowRunRecord {
348 run_id: run_id.to_string(),
349 owner_session_id: Some("session-journal".to_string()),
350 status,
351 lifecycle_seq: 1,
352 started_at_ms: 1,
353 completed_at_ms: None,
354 source_path: None,
355 workflow_id: Some("fixture".to_string()),
356 workflow_goal: Some("journal test".to_string()),
357 token_budget: None,
358 child_ids: Vec::new(),
359 progress_count: 0,
360 progress: Vec::new(),
361 events: Vec::new(),
362 schema_errors: Vec::new(),
363 schema_repairs: Vec::new(),
364 schema_repair_count: 0,
365 dispatch_failure_count: 0,
366 dispatch_failures: Vec::new(),
367 result: None,
368 execution: None,
369 error: None,
370 verify_on_complete: false,
371 verification: None,
372 plan_approval: None,
373 gate_status: Vec::new(),
374 usage: None,
375 events_total: 0,
376 events_dropped: 0,
377 }
378 }
379
380 /// Opening the journal creates nothing; the first record creates the
381 /// ledger under a private `.codewhale/`.
382 #[test]
383 fn opening_the_journal_creates_nothing_until_a_record_is_appended() {
384 let tmp = tempfile::tempdir().expect("tempdir");
385 let state = WorkflowWorkspaceState::open(tmp.path());
386 assert!(!tmp.path().join(CODEWHALE_DIR).exists());
387 state.record_snapshot(&sample_record("workflow_lazy", WorkflowRunStatus::Running));
388 assert!(state.journal_path().is_file());
389 }
390
391 /// A linked `.codewhale` gets no empty `state/` or ledger, and appends and
392 /// reads are refused instead of writing through the link.
393 #[cfg(unix)]
394 #[test]
395 fn journal_never_writes_through_a_linked_codewhale_dir() {
396 let tmp = tempfile::tempdir().expect("tempdir");
397 let workspace = tmp.path().join("workspace");
398 let outside = tmp.path().join("outside");
399 std::fs::create_dir_all(&workspace).expect("mkdir workspace");
400 std::fs::create_dir_all(&outside).expect("mkdir outside");
401 std::os::unix::fs::symlink(&outside, workspace.join(CODEWHALE_DIR)).expect("link");
402
403 let journal = WorkflowRunJournal::open(&workspace);
404 assert!(journal.hydrate_runs(true).is_empty());
405 let err = journal
406 .append_snapshot(&sample_record("workflow_link", WorkflowRunStatus::Running))
407 .expect_err("append through a linked .codewhale must be refused");
408 assert!(!err.to_string().is_empty());
409
410 // The state-level entry points degrade to a warning, not a write.
411 let state = WorkflowWorkspaceState::open(&workspace);
412 state.record_progress("workflow_link", "phase: scan");
413 assert_eq!(
414 std::fs::read_dir(&outside).expect("read outside").count(),
415 0
416 );
417 }
418
419 #[test]
420 fn workflow_journal_hydrates_snapshots_and_progress() {
421 let tmp = tempfile::tempdir().expect("tempdir");
422 let state = WorkflowWorkspaceState::open(tmp.path());
423 let running = sample_record("workflow_abc", WorkflowRunStatus::Running);
424 state.record_snapshot(&running);
425 state.record_progress("workflow_abc", "phase: scan");
426 state.record_event(
427 "workflow_abc",
428 &WorkflowUiEvent::at(
429 5,
430 "session-journal",
431 WorkflowUiEventKind::PhaseStarted {
432 title: "scan".to_string(),
433 },
434 ),
435 );
436
437 let completed = WorkflowRunRecord {
438 status: WorkflowRunStatus::Completed,
439 completed_at_ms: Some(99),
440 progress: vec!["phase: scan".to_string()],
441 events: vec![WorkflowUiEvent::at(
442 5,
443 "session-journal",
444 WorkflowUiEventKind::PhaseStarted {
445 title: "scan".to_string(),
446 },
447 )],
448 ..sample_record("workflow_abc", WorkflowRunStatus::Completed)
449 };
450 state.record_snapshot(&completed);
451 state.record_event(
452 "workflow_abc",
453 &WorkflowUiEvent::at(
454 6,
455 "session-journal",
456 WorkflowUiEventKind::HandoffPromoted {
457 artifact_id: "workflow_abc:scout-1:scout-gate:findings".to_string(),
458 gate_id: "scout-gate".to_string(),
459 kind: "findings".to_string(),
460 from_role: "scout".to_string(),
461 to_role: "implementer".to_string(),
462 producer_task_id: "scout-1".to_string(),
463 },
464 ),
465 );
466 state.record_event(
467 "workflow_abc",
468 &WorkflowUiEvent::at(
469 7,
470 "session-journal",
471 WorkflowUiEventKind::HandoffConsumed {
472 artifact_id: "workflow_abc:scout-1:scout-gate:findings".to_string(),
473 kind: "findings".to_string(),
474 from_role: "scout".to_string(),
475 to_role: "implementer".to_string(),
476 consumer_task_id: "implementer-1".to_string(),
477 },
478 ),
479 );
480
481 let reloaded = WorkflowWorkspaceState::open(tmp.path());
482 let runs = reloaded
483 .runs
484 .lock()
485 .expect("runs lock")
486 .get("workflow_abc")
487 .cloned()
488 .expect("hydrated run");
489 assert_eq!(runs.status, WorkflowRunStatus::Completed);
490 assert_eq!(runs.progress, vec!["phase: scan"]);
491 assert_eq!(runs.events.len(), 3);
492 assert_eq!(runs.events[0].event_type(), "phase_started");
493 let promoted = serde_json::to_value(&runs.events[1]).expect("promoted receipt");
494 assert_eq!(promoted["type"], "handoff_promoted");
495 assert_eq!(
496 promoted["artifact_id"],
497 "workflow_abc:scout-1:scout-gate:findings"
498 );
499 assert_eq!(promoted["gate_id"], "scout-gate");
500 assert_eq!(promoted["producer_task_id"], "scout-1");
501 assert!(promoted.get("payload").is_none(), "{promoted}");
502 let consumed = serde_json::to_value(&runs.events[2]).expect("consumed receipt");
503 assert_eq!(consumed["type"], "handoff_consumed");
504 assert_eq!(consumed["artifact_id"], promoted["artifact_id"]);
505 assert_eq!(consumed["consumer_task_id"], "implementer-1");
506 assert!(consumed.get("payload").is_none(), "{consumed}");
507 assert_eq!(runs.completed_at_ms, Some(99));
508
509 // The event-line replay above must also survive compaction into a
510 // final Snapshot record containing both handoff variants.
511 reloaded.record_snapshot(&runs);
512 let reopened = WorkflowWorkspaceState::open(tmp.path());
513 let compacted = reopened
514 .runs
515 .lock()
516 .expect("runs lock")
517 .get("workflow_abc")
518 .cloned()
519 .expect("snapshot with handoff receipts");
520 assert_eq!(
521 compacted
522 .events
523 .iter()
524 .map(WorkflowUiEvent::event_type)
525 .collect::<Vec<_>>(),
526 vec!["phase_started", "handoff_promoted", "handoff_consumed"]
527 );
528 }
529
530 #[test]
531 fn workflow_journal_rebuilds_a_bounded_exact_rejection_ledger() {
532 let tmp = tempfile::tempdir().expect("tempdir");
533 let state = WorkflowWorkspaceState::open(tmp.path());
534 state.record_snapshot(&sample_record(
535 "workflow_rejections",
536 WorkflowRunStatus::Running,
537 ));
538 let total = WORKFLOW_RUN_DISPATCH_FAILURES_MAX_RETAINED + 5;
539 for index in 0..total {
540 let message = format!("invalid task options {index}");
541 state.record_progress(
542 "workflow_rejections",
543 &format!("dispatch failed for rejected-{index}: {message}"),
544 );
545 state.record_event(
546 "workflow_rejections",
547 &WorkflowUiEvent::at(
548 index as u64,
549 "session-journal",
550 WorkflowUiEventKind::TaskDispatchFailed {
551 label: Some(format!("rejected-{index}")),
552 phase: Some("fan-out".to_string()),
553 message,
554 queue_ticket: None,
555 },
556 ),
557 );
558 }
559 drop(state);
560
561 let reloaded = WorkflowWorkspaceState::open(tmp.path());
562 let run = reloaded
563 .runs
564 .lock()
565 .expect("runs lock")
566 .get("workflow_rejections")
567 .cloned()
568 .expect("hydrated rejection run");
569 assert_eq!(run.progress_count, total as u64);
570 assert_eq!(run.progress.len(), total);
571 assert_eq!(run.dispatch_failure_count, total as u64);
572 assert_eq!(
573 run.dispatch_failures.len(),
574 WORKFLOW_RUN_DISPATCH_FAILURES_MAX_RETAINED
575 );
576 assert_eq!(
577 run.dispatch_failures
578 .first()
579 .and_then(|failure| failure.label.as_deref()),
580 Some("rejected-5")
581 );
582 drop(reloaded);
583
584 // Restart recovery appends a compact snapshot. Replaying the
585 // journal again must not double-count its earlier event lines.
586 let reopened = WorkflowWorkspaceState::open(tmp.path());
587 let run = reopened
588 .runs
589 .lock()
590 .expect("runs lock")
591 .get("workflow_rejections")
592 .cloned()
593 .expect("rehydrated rejection run");
594 assert_eq!(run.dispatch_failure_count, total as u64);
595 assert_eq!(
596 run.dispatch_failures.len(),
597 WORKFLOW_RUN_DISPATCH_FAILURES_MAX_RETAINED
598 );
599 }
600
601 #[test]
602 fn workflow_journal_marks_orphaned_running_runs_failed() {
603 let tmp = tempfile::tempdir().expect("tempdir");
604 let state = WorkflowWorkspaceState::open(tmp.path());
605 state.record_snapshot(&sample_record(
606 "workflow_orphan",
607 WorkflowRunStatus::Running,
608 ));
609
610 let reloaded = WorkflowWorkspaceState::open(tmp.path());
611 let run = reloaded
612 .runs
613 .lock()
614 .expect("runs lock")
615 .get("workflow_orphan")
616 .cloned()
617 .expect("hydrated run");
618 assert_eq!(run.status, WorkflowRunStatus::Failed);
619 assert_eq!(
620 run.lifecycle_seq, 2,
621 "restart recovery is a durable owner lifecycle transition"
622 );
623 assert!(
624 run.completed_at_ms.is_some(),
625 "restart recovery must terminalize the durable owner record"
626 );
627 assert!(
628 run.error
629 .as_deref()
630 .is_some_and(|error| error.contains("process exited")),
631 "expected orphan recovery error, got {:?}",
632 run.error
633 );
634
635 let reopened = WorkflowWorkspaceState::open(tmp.path());
636 let replayed = reopened
637 .runs
638 .lock()
639 .expect("runs lock")
640 .get("workflow_orphan")
641 .cloned()
642 .expect("durably recovered run");
643 assert_eq!(replayed.status, WorkflowRunStatus::Failed);
644 assert_eq!(
645 replayed.lifecycle_seq, 2,
646 "reopening must replay the recovery snapshot without another transition"
647 );
648 }
649
650 #[test]
651 fn host_cancel_hydrates_a_journal_without_live_process_state() {
652 let tmp = tempfile::tempdir().expect("tempdir");
653 let state = WorkflowWorkspaceState::open(tmp.path());
654 let record = sample_record("workflow_prior", WorkflowRunStatus::Running);
655 state.record_snapshot(&record);
656 drop(state);
657
658 assert!(
659 peek_shared_workflow_state(tmp.path()).is_none(),
660 "writing the journal must not insert process-wide live state"
661 );
662
663 let line = super::super::host_cancel_workflow(
664 tmp.path(),
665 "workflow_prior",
666 Some("session-journal"),
667 )
668 .expect("a journaled run must be visible to host cancel after restart");
669 assert_eq!(line.run_id, "workflow_prior");
670 assert_eq!(line.status, "cancelled");
671 assert!(
672 line.error
673 .as_deref()
674 .is_some_and(|error| error.contains("no live process")),
675 "controller-less cancel must leave an honest receipt, got {:?}",
676 line.error
677 );
678
679 let reopened = WorkflowWorkspaceState::open(tmp.path());
680 let replayed = reopened
681 .runs
682 .lock()
683 .expect("runs lock")
684 .get("workflow_prior")
685 .cloned()
686 .expect("cancelled journal line");
687 assert_eq!(replayed.status, WorkflowRunStatus::Cancelled);
688 }
689
690 #[test]
691 fn host_stage_is_derived_from_typed_owner_events() {
692 let mut record = sample_record("workflow_stage", WorkflowRunStatus::Running);
693 record.push_event(WorkflowUiEvent::at(
694 1,
695 "session-journal",
696 WorkflowUiEventKind::RunStarted {
697 workflow_id: Some("fixture".to_string()),
698 workflow_goal: Some("review release".to_string()),
699 source_path: None,
700 token_budget: None,
701 max_concurrent: None,
702 },
703 ));
704 assert_eq!(super::super::host_workflow_stage(&record), "queued");
705
706 record.push_event(WorkflowUiEvent::at(
707 2,
708 "session-journal",
709 WorkflowUiEventKind::PhaseStarted {
710 title: "review".to_string(),
711 },
712 ));
713 assert_eq!(super::super::host_workflow_stage(&record), "running");
714
715 record.push_event(WorkflowUiEvent::at(
716 3,
717 "session-journal",
718 WorkflowUiEventKind::TaskStarted(Box::new(super::super::WorkflowTaskStartedEvent {
719 task_id: "reviewer-1".to_string(),
720 label: Some("reviewer".to_string()),
721 role: None,
722 profile: None,
723 model: None,
724 strength: None,
725 thinking: None,
726 requested_reasoning: None,
727 effective_reasoning: None,
728 resolved_role: Some("reviewer".to_string()),
729 resolved_profile: None,
730 resolved_provider: "local".to_string(),
731 resolved_model: "stub".to_string(),
732 route_source: "session".to_string(),
733 child_route: None,
734 worktree: false,
735 workspace: None,
736 git_branch: None,
737 parent_task_id: None,
738 depth: 0,
739 workflow_run_id: Some("workflow_stage".to_string()),
740 workflow_phase_id: Some("review".to_string()),
741 workflow_task_label: Some("reviewer".to_string()),
742 workflow_child_index: Some(0),
743 queue_ticket: None,
744 fleet_receipt: None,
745 })),
746 ));
747 assert_eq!(super::super::host_workflow_stage(&record), "waiting");
748
749 record.push_event(WorkflowUiEvent::at(
750 4,
751 "session-journal",
752 WorkflowUiEventKind::TaskCompleted {
753 task_id: "reviewer-1".to_string(),
754 status: super::super::IrWorkflowRunStatus::Succeeded,
755 reason: None,
756 kind: None,
757 usage: None,
758 },
759 ));
760 assert_eq!(super::super::host_workflow_stage(&record), "running");
761
762 record.status = WorkflowRunStatus::Completed;
763 assert_eq!(super::super::host_workflow_stage(&record), "completed");
764 record.status = WorkflowRunStatus::Failed;
765 assert_eq!(super::super::host_workflow_stage(&record), "failed");
766 record.status = WorkflowRunStatus::Cancelled;
767 assert_eq!(super::super::host_workflow_stage(&record), "cancelled");
768 }
769
770 #[test]
771 fn host_run_details_derive_phases_and_child_states_from_the_journal() {
772 let tmp = tempfile::tempdir().expect("tempdir");
773 let state = WorkflowWorkspaceState::open(tmp.path());
774 let mut record = sample_record("workflow_detail", WorkflowRunStatus::Running);
775 record.workflow_goal = Some("audit provider errors".to_string());
776 for message in ["phase: scan", "child slow-1 done", "child slow-2 failed"] {
777 record.push_progress(message.to_string());
778 }
779 state.record_snapshot(&record);
780 drop(state);
781
782 let phase: WorkflowUiEvent = serde_json::from_value(serde_json::json!({
783 "at_ms": 1,
784 "owner_session_id": "session-journal",
785 "type": "phase_started",
786 "title": "scan"
787 }))
788 .expect("phase_started event");
789 let started: WorkflowUiEvent = WorkflowUiEvent::at(
790 2,
791 "session-journal",
792 WorkflowUiEventKind::TaskStarted(Box::new(super::super::WorkflowTaskStartedEvent {
793 task_id: "child-1".to_string(),
794 label: Some("slow-1".to_string()),
795 role: None,
796 profile: None,
797 model: None,
798 strength: None,
799 thinking: None,
800 requested_reasoning: None,
801 effective_reasoning: None,
802 resolved_role: Some("explore".to_string()),
803 resolved_profile: None,
804 resolved_provider: "deepseek".to_string(),
805 resolved_model: "deepseek-v4-flash".to_string(),
806 route_source: "session".to_string(),
807 child_route: None,
808 worktree: false,
809 workspace: None,
810 git_branch: None,
811 parent_task_id: None,
812 depth: 0,
813 workflow_run_id: Some("workflow_detail".to_string()),
814 workflow_phase_id: Some("scan".to_string()),
815 workflow_task_label: None,
816 workflow_child_index: Some(0),
817 queue_ticket: None,
818 fleet_receipt: None,
819 })),
820 );
821 let completed: WorkflowUiEvent = serde_json::from_value(serde_json::json!({
822 "at_ms": 3,
823 "owner_session_id": "session-journal",
824 "type": "task_completed",
825 "task_id": "child-1",
826 "status": "failed"
827 }))
828 .expect("task_completed event");
829 let replay = WorkflowWorkspaceState::open(tmp.path());
830 replay.record_event("workflow_detail", &phase);
831 replay.record_event("workflow_detail", &started);
832 replay.record_event("workflow_detail", &completed);
833 drop(replay);
834
835 let details = super::super::host_workflow_run_details(tmp.path(), Some("session-journal"));
836 assert_eq!(details.len(), 1, "one journaled run");
837 let detail = &details[0];
838 assert_eq!(detail.line.run_id, "workflow_detail");
839 // Journal-only `running` rows hydrate through restart-orphan
840 // recovery (the same rewrite `WorkflowWorkspaceState::open`
841 // applies), so the host projection reports the run as failed —
842 // live in-process runs keep `running` via the shared state.
843 assert_eq!(detail.line.status, "failed");
844 assert_eq!(detail.line.label, "audit provider errors");
845 assert_eq!(detail.phases, vec!["scan".to_string()]);
846 assert_eq!(detail.children.len(), 1);
847 let child = &detail.children[0];
848 assert_eq!(child.task_id, "child-1");
849 assert_eq!(child.label.as_deref(), Some("slow-1"));
850 assert_eq!(child.role.as_deref(), Some("explore"));
851 assert_eq!(child.model.as_deref(), Some("deepseek-v4-flash"));
852 assert_eq!(child.phase.as_deref(), Some("scan"));
853 assert_eq!(
854 child.state, "failed",
855 "terminal event must win over running"
856 );
857 assert_eq!(detail.progress_tail.len(), 3);
858 assert!(!detail.has_result);
859
860 // Session ownership fences the projection: a foreign session
861 // sees nothing, exactly like every other host control.
862 assert!(
863 super::super::host_workflow_run_details(tmp.path(), Some("session-other")).is_empty()
864 );
865 }
866 }
867
867 lines RUST