1use std::path::{Path, PathBuf};
49
50use chrono::{DateTime, Utc};
51use serde_json::Value;
52
53use crate::error::{Error, Result};
54use crate::paths::RunPaths;
55use crate::projections::{read_manifest_opt, read_node_opt, write_manifest, write_node};
56use crate::report::ReportOrigin;
57use crate::schema::{
58 ChildRef, Event, EvidenceStatus, IdValidationError, Kind, Lifecycle, Manifest, MergeTxn, Node,
59 NodeId, RunId, Status, TmuxIdentity, WorkerEvidence, WorkerExit, STATE_SCHEMA_VERSION,
60};
61
62fn corrupt_id(events_path: &Path, ev: &Event, e: &IdValidationError) -> Error {
68 Error::CorruptEventLog {
69 path: events_path.to_path_buf(),
70 reason: format!("event seq={} kind={}: {e}", ev.seq, ev.kind),
71 }
72}
73
74fn opt_run_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<RunId>> {
78 match d.get(field) {
79 None | Some(Value::Null) => Ok(None),
80 Some(Value::String(s)) => RunId::parse_str(s)
81 .map(Some)
82 .map_err(|e| corrupt_id(events_path, ev, &e)),
83 Some(_) => Err(Error::CorruptEventLog {
84 path: events_path.to_path_buf(),
85 reason: format!(
86 "event seq={} kind={} `{field}` must be a JSON string or null",
87 ev.seq, ev.kind
88 ),
89 }),
90 }
91}
92
93fn opt_node_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<NodeId>> {
95 match d.get(field) {
96 None | Some(Value::Null) => Ok(None),
97 Some(Value::String(s)) => NodeId::parse_str(s)
98 .map(Some)
99 .map_err(|e| corrupt_id(events_path, ev, &e)),
100 Some(_) => Err(Error::CorruptEventLog {
101 path: events_path.to_path_buf(),
102 reason: format!(
103 "event seq={} kind={} `{field}` must be a JSON string or null",
104 ev.seq, ev.kind
105 ),
106 }),
107 }
108}
109
110fn data_kind(v: &Value) -> Option<Kind> {
124 match serde_json::from_value::<Kind>(v.clone()) {
125 Ok(Kind::Unknown) | Err(_) => None,
126 Ok(k) => Some(k),
127 }
128}
129
130fn data_status(v: &Value) -> Option<Status> {
131 serde_json::from_value(v.clone()).ok()
132}
133
134fn require_status(ev: &Event, path: PathBuf) -> Result<Status> {
135 data_status(ev.data.get("status").unwrap_or(&Value::Null)).ok_or_else(|| {
136 Error::CorruptEventLog {
137 path,
138 reason: format!("{} missing/invalid `status`", ev.kind),
139 }
140 })
141}
142
143fn want_str<'a>(events_path: &Path, ev: &Event, d: &'a Value, field: &str) -> Result<&'a str> {
144 d.get(field)
145 .and_then(Value::as_str)
146 .ok_or_else(|| Error::CorruptEventLog {
147 path: events_path.to_path_buf(),
148 reason: format!(
149 "event seq={} kind={} missing `{field}` string field",
150 ev.seq, ev.kind
151 ),
152 })
153}
154
155fn optional_bool(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<bool>> {
161 match d.get(field) {
162 None | Some(Value::Null) => Ok(None),
163 Some(Value::Bool(b)) => Ok(Some(*b)),
164 Some(_) => Err(Error::CorruptEventLog {
165 path: events_path.to_path_buf(),
166 reason: format!(
167 "event seq={} kind={} `{field}` must be a JSON boolean or null",
168 ev.seq, ev.kind
169 ),
170 }),
171 }
172}
173
174fn optional_i32(d: &Value, field: &str, events_path: &Path, ev: &Event) -> Result<Option<i32>> {
175 match d.get(field) {
176 None | Some(Value::Null) => Ok(None),
177 Some(v) => {
178 let raw = v.as_i64().ok_or_else(|| Error::CorruptEventLog {
179 path: events_path.to_path_buf(),
180 reason: format!(
181 "event seq={} kind={} `{field}` must be integer",
182 ev.seq, ev.kind
183 ),
184 })?;
185 i32::try_from(raw)
186 .map(Some)
187 .map_err(|_| Error::CorruptEventLog {
188 path: events_path.to_path_buf(),
189 reason: format!(
190 "event seq={} kind={} `{field}` out of i32 range: {raw}",
191 ev.seq, ev.kind
192 ),
193 })
194 }
195 }
196}
197
198fn optional_ts(
199 d: &Value,
200 field: &str,
201 events_path: &Path,
202 ev: &Event,
203) -> Result<Option<DateTime<Utc>>> {
204 match d.get(field) {
205 None | Some(Value::Null) => Ok(None),
206 Some(Value::String(s)) => DateTime::parse_from_rfc3339(s)
207 .map(|dt| Some(dt.with_timezone(&Utc)))
208 .map_err(|_| Error::CorruptEventLog {
209 path: events_path.to_path_buf(),
210 reason: format!(
211 "event seq={} kind={} `{field}` not RFC3339",
212 ev.seq, ev.kind
213 ),
214 }),
215 Some(_) => Err(Error::CorruptEventLog {
216 path: events_path.to_path_buf(),
217 reason: format!(
218 "event seq={} kind={} `{field}` must be RFC3339 string or null",
219 ev.seq, ev.kind
220 ),
221 }),
222 }
223}
224
225pub(crate) enum ProjectionOp {
235 Manifest(Manifest),
237 Node(Node),
239}
240
241pub(crate) fn commit_ops(paths: &RunPaths, ops: Vec<ProjectionOp>) -> Result<()> {
249 for op in ops {
250 match op {
251 ProjectionOp::Manifest(m) => write_manifest(paths, &m)?,
252 ProjectionOp::Node(n) => write_node(paths, &n)?,
253 }
254 }
255 Ok(())
256}
257
258pub(crate) fn reduce_event_to_ops(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
274 if ev.run_id != paths.run_id {
278 return Err(Error::CorruptEventLog {
279 path: paths.events(),
280 reason: format!(
281 "event seq={} envelope run_id {:?} does not match run {:?}",
282 ev.seq,
283 ev.run_id.as_str(),
284 paths.run_id.as_str()
285 ),
286 });
287 }
288 #[allow(clippy::match_same_arms)]
291 match ev.kind.as_str() {
292 "run.created" => reduce_run_created(paths, ev),
293 "run.status" => reduce_run_status(paths, ev),
294 "node.created" => reduce_node_created(paths, ev),
295 "node.status" => reduce_node_status(paths, ev),
296 "node.report" => reduce_node_report(paths, ev),
297 "node.retry" => reduce_node_retry(paths, ev),
298 "worker.exited" => reduce_worker_exited(paths, ev),
299 "worker.evidence.archived" => reduce_worker_evidence_archived(paths, ev),
300 "worker.evidence.failed" => reduce_worker_evidence_failed(paths, ev),
301 "node.death_observed" => reduce_node_death_observed(paths, ev),
302 "node.awaiting_input" => reduce_node_awaiting_input(paths, ev),
303 "node.input_resolved" => reduce_node_input_resolved(paths, ev),
304 KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
305 KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
306 "child.spawned" => reduce_child_spawned(paths, ev),
307 "supervisor.attached" => reduce_supervisor_attached(paths, ev),
308 "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
309 "supervisor.exited" => Ok(vec![]),
310 "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
318 "run.notified" | "run.awaiting_input_notified" => Ok(vec![]),
325 "cleanup.window_missing"
350 | "cleanup.worktree_missing"
351 | "cleanup.branch_remove_failed"
352 | "cleanup.branch_preserved"
353 | "cleanup.discard_authorized"
354 | "cleanup.session_killed"
355 | "cleanup.session_retained" => Ok(vec![]),
356 "supervisor.child_id_quarantined" => Ok(vec![]),
366 _ => Ok(vec![]),
367 }
368}
369
370fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
376 match op {
377 ProjectionOp::Manifest(_) => paths.manifest(),
378 ProjectionOp::Node(n) => paths.node(&n.node_id),
379 }
380}
381
382pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
404 let ops = reduce_event_to_ops(paths, event)?;
405 Ok(ops.iter().map(|op| op_path(paths, op)).collect())
406}
407
408pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
424 let ops = reduce_event_to_ops(paths, ev)?;
425 commit_ops(paths, ops)
426}
427
428pub fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
434 reduce_event_to_ops(paths, ev).map(|_| ())
435}
436
437fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
441 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
442 path: events_path.to_path_buf(),
443 reason: format!(
444 "event seq={} kind={} missing top-level `node_id`",
445 ev.seq, ev.kind
446 ),
447 })
448}
449
450fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
451 if let Some(existing) = read_manifest_opt(paths)? {
455 if existing.run_id != ev.run_id {
456 return Err(Error::CorruptEventLog {
457 path: paths.manifest(),
458 reason: format!(
459 "run.created run_id={} conflicts with existing manifest run_id={}",
460 ev.run_id, existing.run_id
461 ),
462 });
463 }
464 return Ok(vec![]);
465 }
466 let events_path = paths.events();
467 let d = &ev.data;
468 let kind =
469 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
470 path: events_path.clone(),
471 reason: "run.created missing/invalid `kind`".into(),
472 })?;
473 let lifecycle: Lifecycle = serde_json::from_value(
474 d.get("lifecycle").cloned().unwrap_or(Value::Null),
475 )
476 .map_err(|_| Error::CorruptEventLog {
477 path: events_path.clone(),
478 reason: "run.created missing/invalid `lifecycle`".into(),
479 })?;
480 let title = want_str(&events_path, ev, d, "title")?.to_string();
481 let agent_selection: Option<crate::schema::AgentSelection> = d
482 .get("agent_selection")
483 .cloned()
484 .map(serde_json::from_value)
485 .transpose()
486 .map_err(|e| Error::CorruptEventLog {
487 path: events_path.clone(),
488 reason: format!("run.created invalid `agent_selection`: {e}"),
489 })?;
490 if let Some(selection) = &agent_selection {
491 selection
492 .validate()
493 .map_err(|reason| Error::CorruptEventLog {
494 path: events_path.clone(),
495 reason: format!("run.created invalid `agent_selection`: {reason}"),
496 })?;
497 }
498 let m = Manifest {
499 schema_version: STATE_SCHEMA_VERSION,
500 applied_seq: 0,
503 run_id: paths.run_id.clone(),
505 kind,
506 lifecycle,
507 title,
508 status: Status::Pending,
509 created_at: ev.ts,
510 updated_at: ev.ts,
511 source_repo: d
512 .get("source_repo")
513 .and_then(Value::as_str)
514 .map(str::to_string),
515 source_branch: d
516 .get("source_branch")
517 .and_then(Value::as_str)
518 .map(str::to_string),
519 worktree_root: d
520 .get("worktree_root")
521 .and_then(Value::as_str)
522 .map(str::to_string),
523 managed_tmux_session: d
524 .get("managed_tmux_session")
525 .and_then(Value::as_str)
526 .map(str::to_string),
527 notify_cmd: d
528 .get("notify_cmd")
529 .and_then(Value::as_str)
530 .map(str::to_string),
531 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
532 agent_selection,
533 node_count: 0,
534 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
535 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
536 };
537 Ok(vec![ProjectionOp::Manifest(m)])
538}
539
540fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
541 let mut m = match read_manifest_opt(paths)? {
542 Some(m) => m,
543 None => return Ok(vec![]),
544 };
545 let new_status = require_status(ev, paths.events())?;
546 if m.status.is_terminal() {
552 let recovery_seq = ev
553 .data
554 .get("recovery_merge_report_seq")
555 .and_then(Value::as_u64);
556 let children_successful =
557 ev.data.get("children_successful").and_then(Value::as_bool) == Some(true);
558 let recovery_allowed = m.status == Status::Failed
559 && new_status == Status::Done
560 && children_successful
561 && recovery_seq.is_some_and(|wanted| {
562 crate::cancel::read_node_status_facts(paths, Some(ev.seq))
563 .ok()
564 .filter(|facts| {
565 crate::aggregate_terminal_status(facts.iter().map(|fact| fact.status))
566 == Some(Status::Done)
567 })
568 .is_some_and(|facts| {
569 facts
570 .iter()
571 .any(|fact| fact.confirmed_merge_seq == Some(wanted))
572 })
573 });
574 if !recovery_allowed {
575 trace_terminal_noop(ev, m.status, new_status);
576 return Ok(vec![]);
577 }
578 }
579 if m.status == new_status {
580 return Ok(vec![]);
581 }
582 m.status = new_status;
583 m.updated_at = ev.ts;
584 Ok(vec![ProjectionOp::Manifest(m)])
585}
586
587fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
598 let nonempty = |key| {
599 d.get(key)
600 .and_then(Value::as_str)
601 .map(str::trim)
602 .filter(|s| !s.is_empty())
603 .map(str::to_string)
604 };
605 let session = nonempty("tmux_session")?;
606 let window_id = nonempty("tmux_window_id")?;
607 Some(TmuxIdentity {
608 socket: nonempty("tmux_socket"),
609 session,
610 window_id,
611 pane_id: nonempty("tmux_pane_id"),
614 })
615}
616
617fn worker_evidence_from_spawn_data(
618 events_path: &Path,
619 ev: &Event,
620 data: &Value,
621) -> Result<Option<WorkerEvidence>> {
622 let Some(session_id) = data.get("pi_session_id") else {
623 if data.get("pi_session_path").is_some() || data.get("pi_session_cwd").is_some() {
624 return Err(Error::CorruptEventLog {
625 path: events_path.to_path_buf(),
626 reason: format!(
627 "event seq={} Pi evidence fields must be all-or-none",
628 ev.seq
629 ),
630 });
631 }
632 return Ok(None);
633 };
634 if session_id.is_null() {
635 if data
636 .get("pi_session_path")
637 .is_some_and(|value| !value.is_null())
638 || data
639 .get("pi_session_cwd")
640 .is_some_and(|value| !value.is_null())
641 {
642 return Err(Error::CorruptEventLog {
643 path: events_path.to_path_buf(),
644 reason: format!(
645 "event seq={} Pi evidence fields must be all-or-none",
646 ev.seq
647 ),
648 });
649 }
650 return Ok(None);
651 }
652 let session_id = session_id
653 .as_str()
654 .filter(|v| !v.is_empty())
655 .ok_or_else(|| Error::CorruptEventLog {
656 path: events_path.to_path_buf(),
657 reason: format!(
658 "event seq={} pi_session_id must be a non-empty string",
659 ev.seq
660 ),
661 })?;
662 if session_id.len() != 36
663 || !session_id.bytes().enumerate().all(|(index, byte)| {
664 if matches!(index, 8 | 13 | 18 | 23) {
665 byte == b'-'
666 } else {
667 byte.is_ascii_hexdigit()
668 }
669 })
670 {
671 return Err(Error::CorruptEventLog {
672 path: events_path.to_path_buf(),
673 reason: format!("event seq={} pi_session_id must be a UUID", ev.seq),
674 });
675 }
676 let original_cwd = want_str(events_path, ev, data, "pi_session_cwd")?;
677 let live_session_path = want_str(events_path, ev, data, "pi_session_path")?;
678 let expected_live_path = format!(
679 ".creating/pi-sessions/{}/pi-session-{session_id}.jsonl",
680 ev.run_id.as_str()
681 );
682 if live_session_path != expected_live_path {
683 return Err(Error::CorruptEventLog {
684 path: events_path.to_path_buf(),
685 reason: format!(
686 "event seq={} pi_session_path is not the canonical state-relative path",
687 ev.seq
688 ),
689 });
690 }
691 let attempt = match data.get("attempt") {
692 None | Some(Value::Null) => 0,
693 Some(value) => value
694 .as_u64()
695 .and_then(|raw| u32::try_from(raw).ok())
696 .ok_or_else(|| Error::CorruptEventLog {
697 path: events_path.to_path_buf(),
698 reason: format!("event seq={} attempt must be a u32", ev.seq),
699 })?,
700 };
701 Ok(Some(WorkerEvidence {
702 attempt,
703 session_id: session_id.to_string(),
704 original_cwd: original_cwd.to_string(),
705 live_session_path: live_session_path.to_string(),
706 status: EvidenceStatus::Pending,
707 transcript_path: None,
708 resume_path: None,
709 pane_path: None,
710 report_path: None,
711 transcript_sha256: None,
712 error: None,
713 }))
714}
715
716fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
717 let events_path = paths.events();
718 let node_id = require_envelope_node_id(&events_path, ev)?;
721 let is_default_node = node_id.as_str() == "n-0001";
722 if read_node_opt(paths, &node_id)?.is_some() {
724 return Ok(vec![]);
725 }
726 let d = &ev.data;
727 let kind =
728 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
729 path: events_path.clone(),
730 reason: format!(
731 "event seq={} kind=node.created missing/invalid `kind`",
732 ev.seq
733 ),
734 })?;
735 let n = Node {
736 schema_version: STATE_SCHEMA_VERSION,
737 node_id,
738 run_id: paths.run_id.clone(),
740 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
741 kind,
742 status: Status::Pending,
743 task: d.get("task").and_then(Value::as_str).map(str::to_string),
744 worktree_path: d
745 .get("worktree_path")
746 .and_then(Value::as_str)
747 .map(str::to_string),
748 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
749 base_sha: d
750 .get("base_sha")
751 .and_then(Value::as_str)
752 .filter(|s| !s.is_empty())
753 .map(str::to_string),
754 tmux_window: d
755 .get("tmux_window")
756 .and_then(Value::as_str)
757 .map(str::to_string),
758 tmux_identity: tmux_identity_from_data(d),
759 evidence: worker_evidence_from_spawn_data(&events_path, ev, d)?,
760 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
761 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
762 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
763 children: Vec::new(),
764 started_at: Some(ev.ts),
765 updated_at: ev.ts,
766 last_report: None,
767 last_processed_report_seq_by_child: serde_json::Map::default(),
768 retry_attempts: 0,
769 worker_exit: None,
770 pending_merge: None,
771 first_death_at: None,
772 awaiting_input: None,
773 };
774 let mut ops = vec![ProjectionOp::Node(n)];
775 if let Some(mut m) = read_manifest_opt(paths)? {
776 if is_default_node && m.source_branch.is_none() {
781 m.source_branch = d
782 .get("source_branch")
783 .and_then(Value::as_str)
784 .filter(|branch| !branch.is_empty())
785 .map(str::to_string);
786 }
787 m.updated_at = ev.ts;
792 ops.push(ProjectionOp::Manifest(m));
793 }
794 Ok(ops)
795}
796
797fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
816 let events_path = paths.events();
817 let node_id = require_envelope_node_id(&events_path, ev)?;
818 let mut n = match read_node_opt(paths, &node_id)? {
819 Some(n) => n,
820 None => return Ok(vec![]),
821 };
822 if n.status.is_terminal() {
825 tracing::debug!(
826 target: "taskfleet_core::reducer",
827 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
828 "no-op: node.retry against terminal node"
829 );
830 return Ok(vec![]);
831 }
832 let d = &ev.data;
833 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
836 n.base_sha = d
837 .get("base_sha")
838 .and_then(Value::as_str)
839 .filter(|s| !s.is_empty())
840 .map(str::to_string);
841 n.worktree_path = d
842 .get("worktree_path")
843 .and_then(Value::as_str)
844 .map(str::to_string);
845 n.tmux_window = d
846 .get("tmux_window")
847 .and_then(Value::as_str)
848 .map(str::to_string);
849 n.tmux_identity = tmux_identity_from_data(d);
850 n.evidence = worker_evidence_from_spawn_data(&events_path, ev, d)?;
851 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
852 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
853 n.status = Status::Pending;
854 n.started_at = Some(ev.ts);
855 n.updated_at = ev.ts;
856 n.last_report = None;
857 n.pending_merge = None;
863 n.worker_exit = None;
868 n.first_death_at = None;
873 n.awaiting_input = None;
876 n.retry_attempts = d
884 .get("attempt")
885 .and_then(Value::as_u64)
886 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
887 Ok(vec![ProjectionOp::Node(n)])
888}
889
890fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
891 let events_path = paths.events();
892 let node_id = require_envelope_node_id(&events_path, ev)?;
893 let mut n = match read_node_opt(paths, &node_id)? {
894 Some(n) => n,
895 None => return Ok(vec![]),
896 };
897 let new_status = require_status(ev, events_path)?;
898 if n.status.is_terminal() {
901 trace_terminal_noop(ev, n.status, new_status);
902 return Ok(vec![]);
903 }
904 if n.status == new_status {
905 return Ok(vec![]);
906 }
907 n.status = new_status;
908 if new_status.is_terminal() {
914 n.pending_merge = None;
915 n.awaiting_input = None;
916 }
917 n.updated_at = ev.ts;
918 Ok(vec![ProjectionOp::Node(n)])
919}
920
921fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
922 let events_path = paths.events();
923 let node_id = require_envelope_node_id(&events_path, ev)?;
924 let mut n = match read_node_opt(paths, &node_id)? {
925 Some(n) => n,
926 None => return Ok(vec![]),
927 };
928 if n.status.is_terminal() {
940 if ReportOrigin::permits_terminal_merge_recovery(n.status, &ev.data) {
973 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
974 return Ok(vec![]);
975 }
976 tracing::info!(
977 target: "taskfleet_core::reducer",
978 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
979 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
980 );
981 n.last_report = Some(ev.data.clone());
982 n.status = Status::Done;
986 n.pending_merge = None;
990 n.awaiting_input = None;
991 n.updated_at = ev.ts;
992 return Ok(vec![ProjectionOp::Node(n)]);
993 }
994 tracing::debug!(
995 target: "taskfleet_core::reducer",
996 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
997 "no-op: node.report against terminal node"
998 );
999 return Ok(vec![]);
1000 }
1001 let new_status = report_terminal_status(&events_path, ev)?;
1009 n.last_report = Some(ev.data.clone());
1010 n.status = new_status;
1011 n.awaiting_input = None;
1015 n.pending_merge = None;
1020 n.updated_at = ev.ts;
1021 Ok(vec![ProjectionOp::Node(n)])
1022}
1023
1024fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1039 let events_path = paths.events();
1040 let node_id = require_envelope_node_id(&events_path, ev)?;
1041 let attempt = match ev.data.get("attempt") {
1042 None | Some(Value::Null) => 0,
1043 Some(_) => evidence_attempt(&events_path, ev)?,
1044 };
1045 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
1046 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
1047 match (code, signal) {
1052 (Some(_), None) | (None, Some(_)) => {}
1053 _ => {
1054 return Err(Error::CorruptEventLog {
1055 path: events_path,
1056 reason: format!(
1057 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
1058 ev.seq
1059 ),
1060 });
1061 }
1062 }
1063 let mut n = match read_node_opt(paths, &node_id)? {
1064 Some(n) => n,
1065 None => return Ok(vec![]),
1072 };
1073 if attempt != n.retry_attempts || n.worker_exit.is_some() {
1077 return Ok(vec![]);
1078 }
1079 n.worker_exit = Some(WorkerExit {
1080 code,
1081 signal,
1082 at: ev.ts,
1083 });
1084 n.awaiting_input = None;
1088 n.updated_at = ev.ts;
1089 Ok(vec![ProjectionOp::Node(n)])
1090}
1091
1092fn evidence_attempt(events_path: &Path, ev: &Event) -> Result<u32> {
1093 ev.data
1094 .get("attempt")
1095 .and_then(Value::as_u64)
1096 .and_then(|raw| u32::try_from(raw).ok())
1097 .ok_or_else(|| Error::CorruptEventLog {
1098 path: events_path.to_path_buf(),
1099 reason: format!("event seq={} attempt must be a u32", ev.seq),
1100 })
1101}
1102
1103fn reduce_worker_evidence_archived(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1104 let events_path = paths.events();
1105 let node_id = require_envelope_node_id(&events_path, ev)?;
1106 let attempt = evidence_attempt(&events_path, ev)?;
1107 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1108 let mut values = Vec::new();
1109 for field in [
1110 "transcript_path",
1111 "resume_path",
1112 "pane_path",
1113 "report_path",
1114 "transcript_sha256",
1115 ] {
1116 values.push(want_str(&events_path, ev, &ev.data, field)?.to_string());
1117 }
1118 let expected_prefix = format!("evidence/{}/", node_id.as_str());
1119 for (index, suffix) in [
1120 "pi-session.original.jsonl",
1121 "pi-session.resume.jsonl",
1122 "final-pane.log",
1123 "terminal-report.json",
1124 ]
1125 .iter()
1126 .enumerate()
1127 {
1128 if values[index] != format!("{expected_prefix}{suffix}") {
1129 return Err(Error::CorruptEventLog {
1130 path: events_path,
1131 reason: format!(
1132 "event seq={} evidence artifact path is not canonical",
1133 ev.seq
1134 ),
1135 });
1136 }
1137 }
1138 if values[4].len() != 64
1139 || !values[4]
1140 .bytes()
1141 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1142 {
1143 return Err(Error::CorruptEventLog {
1144 path: events_path,
1145 reason: format!(
1146 "event seq={} transcript_sha256 is not lowercase SHA-256",
1147 ev.seq
1148 ),
1149 });
1150 }
1151 let mut node = match read_node_opt(paths, &node_id)? {
1152 Some(node) => node,
1153 None => return Ok(vec![]),
1154 };
1155 let Some(evidence) = node.evidence.as_mut() else {
1156 return Err(Error::CorruptEventLog {
1157 path: events_path,
1158 reason: format!(
1159 "event seq={} evidence archive has no recorded Pi session",
1160 ev.seq
1161 ),
1162 });
1163 };
1164 if evidence.attempt != attempt || evidence.session_id != session_id {
1165 return Ok(vec![]);
1166 }
1167 if evidence.status == EvidenceStatus::Complete {
1168 return Ok(vec![]);
1169 }
1170 evidence.transcript_path = Some(values.remove(0));
1171 evidence.resume_path = Some(values.remove(0));
1172 evidence.pane_path = Some(values.remove(0));
1173 evidence.report_path = Some(values.remove(0));
1174 evidence.transcript_sha256 = Some(values.remove(0));
1175 evidence.status = EvidenceStatus::Complete;
1176 evidence.error = None;
1177 node.updated_at = ev.ts;
1178 Ok(vec![ProjectionOp::Node(node)])
1179}
1180
1181fn reduce_worker_evidence_failed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1182 let events_path = paths.events();
1183 let node_id = require_envelope_node_id(&events_path, ev)?;
1184 let attempt = evidence_attempt(&events_path, ev)?;
1185 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1186 let detail = want_str(&events_path, ev, &ev.data, "error")?.to_string();
1187 let mut node = match read_node_opt(paths, &node_id)? {
1188 Some(node) => node,
1189 None => return Ok(vec![]),
1190 };
1191 let Some(evidence) = node.evidence.as_mut() else {
1192 return Err(Error::CorruptEventLog {
1193 path: events_path,
1194 reason: format!(
1195 "event seq={} evidence failure has no recorded Pi session",
1196 ev.seq
1197 ),
1198 });
1199 };
1200 if evidence.attempt != attempt || evidence.session_id != session_id {
1201 return Ok(vec![]);
1202 }
1203 if evidence.status == EvidenceStatus::Complete {
1204 return Ok(vec![]);
1205 }
1206 evidence.status = EvidenceStatus::Failed;
1207 evidence.error = Some(detail);
1208 node.updated_at = ev.ts;
1209 Ok(vec![ProjectionOp::Node(node)])
1210}
1211
1212fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1225 let events_path = paths.events();
1226 let node_id = require_envelope_node_id(&events_path, ev)?;
1227 let mut n = match read_node_opt(paths, &node_id)? {
1228 Some(n) => n,
1229 None => return Ok(vec![]),
1230 };
1231 if n.first_death_at.is_some()
1237 || n.status.is_terminal()
1238 || n.worker_exit.is_some()
1239 || n.last_report.is_some()
1240 || n.pending_merge.is_some()
1241 {
1242 return Ok(vec![]);
1243 }
1244 n.first_death_at = Some(ev.ts);
1245 n.updated_at = ev.ts;
1246 Ok(vec![ProjectionOp::Node(n)])
1247}
1248
1249fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1258 const MAX_ITEMS: usize = 8;
1259 const MAX_TOPIC_CHARS: usize = 512;
1260 const MAX_OPTIONS: usize = 16;
1261 const MAX_OPTION_CHARS: usize = 256;
1262
1263 let events_path = paths.events();
1264 let node_id = require_envelope_node_id(&events_path, ev)?;
1265 let items = ev
1268 .data
1269 .get("discussion_items")
1270 .and_then(Value::as_array)
1271 .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1272 .ok_or_else(|| Error::CorruptEventLog {
1273 path: events_path.clone(),
1274 reason: format!(
1275 "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1276 ev.seq
1277 ),
1278 })?;
1279 for (index, item) in items.iter().enumerate() {
1280 let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1281 path: events_path.clone(),
1282 reason: format!(
1283 "event seq={} discussion_items[{index}] must be an object",
1284 ev.seq
1285 ),
1286 })?;
1287 let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1288 let default = obj
1289 .get("recommended_default")
1290 .and_then(Value::as_str)
1291 .unwrap_or("");
1292 let options = obj.get("options").and_then(Value::as_array);
1293 let options_valid = options.is_some_and(|values| {
1294 !values.is_empty()
1295 && values.len() <= MAX_OPTIONS
1296 && values.iter().all(|v| {
1297 v.as_str().is_some_and(|s| {
1298 !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1299 })
1300 })
1301 && values.iter().any(|v| v.as_str() == Some(default))
1302 });
1303 if topic.trim().is_empty()
1304 || topic.chars().count() > MAX_TOPIC_CHARS
1305 || default.trim().is_empty()
1306 || !options_valid
1307 {
1308 return Err(Error::CorruptEventLog {
1309 path: events_path.clone(),
1310 reason: format!(
1311 "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1312 ev.seq
1313 ),
1314 });
1315 }
1316 }
1317
1318 let mut n = match read_node_opt(paths, &node_id)? {
1319 Some(n) => n,
1320 None => return Ok(vec![]),
1321 };
1322 if n.status.is_terminal() || n.worker_exit.is_some() {
1323 return Ok(vec![]);
1324 }
1325 if let Some(open) = n.awaiting_input.as_mut() {
1326 if open.discussion_items.len() + items.len() > MAX_ITEMS {
1329 return Err(Error::CorruptEventLog {
1330 path: events_path,
1331 reason: format!(
1332 "event seq={} would exceed {MAX_ITEMS} open discussion items",
1333 ev.seq
1334 ),
1335 });
1336 }
1337 open.discussion_items.extend(items.iter().cloned());
1338 } else {
1339 n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1340 opened_at: ev.ts,
1341 event_seq: ev.seq,
1342 discussion_items: items.clone(),
1343 }));
1344 }
1345 n.updated_at = ev.ts;
1346 Ok(vec![ProjectionOp::Node(n)])
1347}
1348
1349fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1353 let events_path = paths.events();
1354 let node_id = require_envelope_node_id(&events_path, ev)?;
1355 let seq = ev
1358 .data
1359 .get("event_seq")
1360 .and_then(Value::as_u64)
1361 .ok_or_else(|| Error::CorruptEventLog {
1362 path: events_path.clone(),
1363 reason: format!(
1364 "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1365 ev.seq
1366 ),
1367 })?;
1368 let mut n = match read_node_opt(paths, &node_id)? {
1369 Some(n) => n,
1370 None => return Ok(vec![]),
1371 };
1372 let Some(open) = n.awaiting_input.as_ref() else {
1373 return Ok(vec![]);
1374 };
1375 if seq != open.event_seq {
1376 return Ok(vec![]);
1377 }
1378 n.awaiting_input = None;
1379 n.updated_at = ev.ts;
1380 Ok(vec![ProjectionOp::Node(n)])
1381}
1382
1383pub const KIND_MERGE_STARTED: &str = "merge.started";
1387
1388pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1393
1394fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1408 let events_path = paths.events();
1409 let node_id = require_envelope_node_id(&events_path, ev)?;
1410 let txn: MergeTxn =
1411 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1412 path: events_path.clone(),
1413 reason: format!(
1414 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1415 ev.seq
1416 ),
1417 })?;
1418 let mut n = match read_node_opt(paths, &node_id)? {
1419 Some(n) => n,
1420 None => return Ok(vec![]),
1421 };
1422 if n.status.is_terminal() {
1426 return Ok(vec![]);
1427 }
1428 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1431 return Ok(vec![]);
1432 }
1433 n.pending_merge = Some(Box::new(txn));
1434 n.updated_at = ev.ts;
1435 Ok(vec![ProjectionOp::Node(n)])
1436}
1437
1438fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1450 let events_path = paths.events();
1451 let node_id = require_envelope_node_id(&events_path, ev)?;
1452 let op_id = ev
1453 .data
1454 .get("op_id")
1455 .and_then(Value::as_str)
1456 .ok_or_else(|| Error::CorruptEventLog {
1457 path: events_path.clone(),
1458 reason: format!(
1459 "event seq={} kind=merge.aborted is missing string `op_id`",
1460 ev.seq
1461 ),
1462 })?;
1463 let mut n = match read_node_opt(paths, &node_id)? {
1464 Some(n) => n,
1465 None => return Ok(vec![]),
1466 };
1467 match n.pending_merge.as_ref() {
1470 Some(t) if t.op_id == op_id => {}
1471 _ => return Ok(vec![]),
1472 }
1473 n.pending_merge = None;
1474 n.updated_at = ev.ts;
1475 Ok(vec![ProjectionOp::Node(n)])
1476}
1477
1478fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1486 if current == incoming {
1487 tracing::debug!(
1488 target: "taskfleet_core::reducer",
1489 seq = ev.seq, kind = %ev.kind, status = ?current,
1490 "no-op: status re-applied to terminal target"
1491 );
1492 } else {
1493 tracing::warn!(
1494 target: "taskfleet_core::reducer",
1495 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1496 "no-op: ignored conflicting transition from terminal target"
1497 );
1498 }
1499}
1500
1501fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1510 let corrupt = |reason: String| Error::CorruptEventLog {
1511 path: events_path.to_path_buf(),
1512 reason,
1513 };
1514 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1515 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1516 if cancelled {
1517 if success == Some(true) {
1518 return Err(corrupt(format!(
1519 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1520 ev.seq
1521 )));
1522 }
1523 Ok(Status::Cancelled)
1524 } else {
1525 match success {
1526 Some(true) => Ok(Status::Done),
1527 Some(false) => Ok(Status::Failed),
1528 None => Err(corrupt(format!(
1529 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1530 ev.seq
1531 ))),
1532 }
1533 }
1534}
1535
1536fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1537 let events_path = paths.events();
1540 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1541 path: events_path.clone(),
1542 reason: format!(
1543 "event seq={} kind=child.spawned missing parent `node_id`",
1544 ev.seq
1545 ),
1546 })?;
1547 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1548 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1549 let child_node_id = NodeId::parse_str(
1550 ev.data
1551 .get("child_node_id")
1552 .and_then(Value::as_str)
1553 .unwrap_or("n-0001"),
1554 )
1555 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1556 let mut n = match read_node_opt(paths, &parent_node_id)? {
1557 Some(n) => n,
1558 None => return Ok(vec![]),
1559 };
1560 let new_ref = ChildRef {
1561 run_id: child_run_id,
1562 node_id: child_node_id,
1563 };
1564 if n.children.iter().any(|c| c == &new_ref) {
1565 return Ok(vec![]);
1568 }
1569 n.children.push(new_ref);
1570 n.updated_at = ev.ts;
1571 Ok(vec![ProjectionOp::Node(n)])
1572}
1573
1574fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1585 let events_path = paths.events();
1586 let node_id = require_envelope_node_id(&events_path, ev)?;
1587 let raw = ev
1588 .data
1589 .get("pid")
1590 .and_then(Value::as_i64)
1591 .ok_or_else(|| Error::CorruptEventLog {
1592 path: events_path.clone(),
1593 reason: format!(
1594 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1595 ev.seq
1596 ),
1597 })?;
1598 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1599 path: events_path.clone(),
1600 reason: format!(
1601 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1602 ev.seq
1603 ),
1604 })?;
1605 let mut n = match read_node_opt(paths, &node_id)? {
1606 Some(n) => n,
1607 None => return Ok(vec![]),
1608 };
1609 if n.supervisor_pid == Some(pid) {
1610 return Ok(vec![]);
1611 }
1612 n.supervisor_pid = Some(pid);
1613 n.updated_at = ev.ts;
1614 Ok(vec![ProjectionOp::Node(n)])
1615}
1616
1617fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1628 let events_path = paths.events();
1629 let node_id = require_envelope_node_id(&events_path, ev)?;
1630 let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1631 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1635 let report_seq = ev
1636 .data
1637 .get("report_seq")
1638 .and_then(Value::as_u64)
1639 .ok_or_else(|| Error::CorruptEventLog {
1640 path: events_path.clone(),
1641 reason: format!(
1642 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1643 ev.seq
1644 ),
1645 })?;
1646 let mut n = match read_node_opt(paths, &node_id)? {
1647 Some(n) => n,
1648 None => return Ok(vec![]),
1649 };
1650 if let Some(prev) = n
1651 .last_processed_report_seq_by_child
1652 .get(child_run_id)
1653 .and_then(Value::as_u64)
1654 {
1655 if report_seq <= prev {
1656 return Ok(vec![]);
1657 }
1658 }
1659 n.last_processed_report_seq_by_child
1660 .insert(child_run_id.to_string(), Value::from(report_seq));
1661 n.updated_at = ev.ts;
1662 Ok(vec![ProjectionOp::Node(n)])
1663}
1664
1665#[cfg(test)]
1666mod tests {
1667 use super::*;
1668 use crate::schema::Event;
1669 use chrono::Utc;
1670 use tempfile::TempDir;
1671
1672 fn event(run_id: &str) -> Event {
1673 Event {
1674 ts: Utc::now(),
1675 seq: 1,
1676 kind: "run.status".into(),
1677 run_id: RunId::parse_str(run_id).unwrap(),
1678 node_id: None,
1679 idempotency_key: None,
1680 data: serde_json::json!({ "status": "running" }),
1681 }
1682 }
1683
1684 #[test]
1685 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1686 let tmp = TempDir::new().unwrap();
1690 let run_id = "01jxsnap000000000000000000";
1691 let rid = RunId::parse_str(run_id).unwrap();
1692 let dir = crate::run_dir(tmp.path(), &rid);
1693 std::fs::create_dir_all(&dir).unwrap();
1694 let paths = RunPaths::new(dir, run_id).unwrap();
1695
1696 let mut created = event(run_id);
1699 created.kind = "run.created".into();
1700 created.data = serde_json::json!({
1701 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1702 });
1703 apply_event(&paths, &created).expect("run.created applies");
1704 let manifest_before = std::fs::read(paths.manifest()).unwrap();
1705
1706 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1707 let mut ev = event(run_id);
1708 ev.seq = seq;
1709 ev.kind = kind.into();
1710 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1712 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1713 assert!(ops.is_empty(), "{kind} must plan no projection ops");
1714 apply_event(&paths, &ev).expect("audit kind applies as no-op");
1716 }
1717
1718 assert_eq!(
1720 std::fs::read(paths.manifest()).unwrap(),
1721 manifest_before,
1722 "audit events must not mutate the manifest"
1723 );
1724 assert!(!paths.nodes_dir().exists(), "no node projection created");
1725 }
1726
1727 #[test]
1728 fn run_created_folds_harness_when_present_and_defaults_none() {
1729 let tmp = TempDir::new().unwrap();
1730
1731 let run_id = "01jxhrnsaa0000000000000001";
1733 let rid = RunId::parse_str(run_id).unwrap();
1734 let dir = crate::run_dir(tmp.path(), &rid);
1735 std::fs::create_dir_all(&dir).unwrap();
1736 let paths = RunPaths::new(dir, run_id).unwrap();
1737 let mut created = event(run_id);
1738 created.kind = "run.created".into();
1739 created.data = serde_json::json!({
1740 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1741 "harness": "pi", "harness_source": "flag",
1742 });
1743 apply_event(&paths, &created).expect("run.created applies");
1744 let m = read_manifest_opt(&paths).unwrap().unwrap();
1745 assert_eq!(m.harness.as_deref(), Some("pi"));
1746
1747 let run_id2 = "01jxhrnsaa0000000000000002";
1749 let rid2 = RunId::parse_str(run_id2).unwrap();
1750 let dir2 = crate::run_dir(tmp.path(), &rid2);
1751 std::fs::create_dir_all(&dir2).unwrap();
1752 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1753 let mut created2 = event(run_id2);
1754 created2.kind = "run.created".into();
1755 created2.data =
1756 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1757 apply_event(&paths2, &created2).expect("run.created applies");
1758 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1759 assert_eq!(m2.harness, None);
1760 }
1761
1762 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1765 let rid = RunId::parse_str(run_id).unwrap();
1766 let dir = crate::run_dir(tmp.path(), &rid);
1767 std::fs::create_dir_all(&dir).unwrap();
1768 let paths = RunPaths::new(dir, run_id).unwrap();
1769 let mut created = event(run_id);
1770 created.kind = "run.created".into();
1771 created.data =
1772 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1773 apply_event(&paths, &created).expect("run.created applies");
1774 let mut node = event(run_id);
1775 node.seq = 2;
1776 node.kind = "node.created".into();
1777 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1778 node.data = serde_json::json!({
1779 "kind": "spinoff",
1780 "branch": "wt/foo",
1781 "worktree_path": "/tmp/old-wt",
1782 "agent_pid": 111,
1783 });
1784 apply_event(&paths, &node).expect("node.created applies");
1785 paths
1786 }
1787
1788 #[test]
1792 fn node_retry_rewires_node_and_increments_attempts() {
1793 let tmp = TempDir::new().unwrap();
1794 let run_id = "01jxsnap000000000000000000";
1795 let paths = bootstrap_retry_node(&tmp, run_id);
1796
1797 let mut retry = event(run_id);
1798 retry.seq = 3;
1799 retry.kind = "node.retry".into();
1800 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1801 retry.data = serde_json::json!({
1802 "attempt": 1,
1803 "reason": "agent-died",
1804 "branch": "wt/foo-r1",
1805 "base_sha": "a".repeat(40),
1806 "worktree_path": "/tmp/new-wt",
1807 "agent_pid": 222,
1808 "tmux_session": "s",
1809 "tmux_window_id": "@9",
1810 });
1811 apply_event(&paths, &retry).expect("node.retry applies");
1812
1813 let n = read_n0001(&paths);
1814 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1815 assert_eq!(
1816 n.branch.as_deref(),
1817 Some("wt/foo-r1"),
1818 "rewired to new branch"
1819 );
1820 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1821 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1822 assert_eq!(n.status, Status::Pending, "node returns to pending");
1823 assert!(n.last_report.is_none());
1824 assert_eq!(
1825 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1826 Some("@9"),
1827 "rewired tmux identity"
1828 );
1829
1830 let mut retry2 = event(run_id);
1832 retry2.seq = 4;
1833 retry2.kind = "node.retry".into();
1834 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1835 retry2.data = serde_json::json!({
1836 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1837 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1838 });
1839 apply_event(&paths, &retry2).expect("node.retry applies");
1840 assert_eq!(read_n0001(&paths).retry_attempts, 2);
1841 }
1842
1843 #[test]
1844 fn stale_worker_evidence_and_exit_do_not_cross_retry_generation() {
1845 let tmp = TempDir::new().unwrap();
1846 let run_id = "01jxsnap000000000000000000";
1847 let paths = bootstrap_retry_node(&tmp, run_id);
1848 let nid = NodeId::parse_str("n-0001").unwrap();
1849 let current_session = "018f5f64-b137-7d44-b2b4-4f02c3f646e8";
1850
1851 let mut retry = event(run_id);
1852 retry.seq = 3;
1853 retry.kind = "node.retry".into();
1854 retry.node_id = Some(nid.clone());
1855 retry.data = serde_json::json!({
1856 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1857 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1858 "pi_session_id": current_session,
1859 "pi_session_path": format!(".creating/pi-sessions/{run_id}/pi-session-{current_session}.jsonl"),
1860 "pi_session_cwd": "/tmp/new-wt"
1861 });
1862 apply_event(&paths, &retry).unwrap();
1863
1864 let mut stale_failure = event(run_id);
1865 stale_failure.seq = 4;
1866 stale_failure.kind = "worker.evidence.failed".into();
1867 stale_failure.node_id = Some(nid.clone());
1868 stale_failure.data = serde_json::json!({
1869 "attempt": 0,
1870 "session_id": "118f5f64-b137-7d44-b2b4-4f02c3f646e8",
1871 "error": "old attempt"
1872 });
1873 apply_event(&paths, &stale_failure).unwrap();
1874 assert_eq!(
1875 read_n0001(&paths).evidence.unwrap().status,
1876 EvidenceStatus::Pending
1877 );
1878
1879 let mut stale_exit = event(run_id);
1880 stale_exit.seq = 5;
1881 stale_exit.kind = "worker.exited".into();
1882 stale_exit.node_id = Some(nid.clone());
1883 stale_exit.data = serde_json::json!({"attempt":0,"exit_code":9});
1884 apply_event(&paths, &stale_exit).unwrap();
1885 assert!(read_n0001(&paths).worker_exit.is_none());
1886
1887 let mut archived = event(run_id);
1888 archived.seq = 6;
1889 archived.kind = "worker.evidence.archived".into();
1890 archived.node_id = Some(nid);
1891 archived.data = serde_json::json!({
1892 "attempt":1, "session_id":current_session,
1893 "transcript_path":"evidence/n-0001/pi-session.original.jsonl",
1894 "resume_path":"evidence/n-0001/pi-session.resume.jsonl",
1895 "pane_path":"evidence/n-0001/final-pane.log",
1896 "report_path":"evidence/n-0001/terminal-report.json",
1897 "transcript_sha256":"0".repeat(64)
1898 });
1899 apply_event(&paths, &archived).unwrap();
1900 assert_eq!(
1901 read_n0001(&paths).evidence.unwrap().status,
1902 EvidenceStatus::Complete
1903 );
1904 }
1905
1906 #[test]
1910 fn node_retry_against_terminal_node_is_noop() {
1911 let tmp = TempDir::new().unwrap();
1912 let run_id = "01jxsnap000000000000000000";
1913 let paths = bootstrap_retry_node(&tmp, run_id);
1914
1915 let mut report = event(run_id);
1917 report.seq = 3;
1918 report.kind = "node.report".into();
1919 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1920 report.data = serde_json::json!({ "success": true });
1921 apply_event(&paths, &report).expect("node.report applies");
1922 assert_eq!(read_n0001(&paths).status, Status::Done);
1923
1924 let mut retry = event(run_id);
1925 retry.seq = 4;
1926 retry.kind = "node.retry".into();
1927 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1928 retry.data = serde_json::json!({
1929 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1930 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1931 });
1932 apply_event(&paths, &retry).expect("node.retry applies as no-op");
1933
1934 let n = read_n0001(&paths);
1935 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
1936 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
1937 assert_eq!(n.agent_pid, Some(111), "not rewired");
1938 }
1939
1940 #[test]
1941 fn apply_event_rejects_event_from_a_different_run() {
1942 let tmp = TempDir::new().unwrap();
1943 let run_id = "01jxsnap000000000000000000";
1944 let rid = RunId::parse_str(run_id).unwrap();
1945 let dir = crate::run_dir(tmp.path(), &rid);
1946 std::fs::create_dir_all(&dir).unwrap();
1947 let paths = RunPaths::new(dir, run_id).unwrap();
1948
1949 let foreign = event("02jxsnap000000000000000000");
1951 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
1952 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
1953
1954 let mine = event(run_id);
1957 apply_event(&paths, &mine).expect("matching run_id must be accepted");
1958 }
1959
1960 #[test]
1961 fn tmux_identity_from_data_reads_qualified_fields() {
1962 let d = serde_json::json!({
1963 "tmux_socket": "/private/tmp/tmux-501/default",
1964 "tmux_session": "taskfleet",
1965 "tmux_window_id": "@42",
1966 });
1967 let id = tmux_identity_from_data(&d).expect("qualified identity");
1968 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
1969 assert_eq!(id.session, "taskfleet");
1970 assert_eq!(id.window_id, "@42");
1971 assert_eq!(id.pane_id, None);
1973
1974 let d2 = serde_json::json!({
1976 "tmux_socket": null,
1977 "tmux_session": "taskfleet",
1978 "tmux_window_id": "@7",
1979 });
1980 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
1981 assert_eq!(id2.socket, None);
1982 assert_eq!(id2.window_id, "@7");
1983
1984 let d3 = serde_json::json!({
1986 "tmux_session": "taskfleet",
1987 "tmux_window_id": "@42",
1988 "tmux_pane_id": "%7",
1989 });
1990 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
1991 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
1992 assert_eq!(id3.capture_target(), "%7");
1993
1994 let d4 = serde_json::json!({
1997 "tmux_session": "taskfleet",
1998 "tmux_window_id": "@42",
1999 "tmux_pane_id": null,
2000 });
2001 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
2002 assert_eq!(id4.pane_id, None);
2003 assert_eq!(id4.capture_target(), "@42");
2004 }
2005
2006 #[test]
2007 fn tmux_identity_from_data_back_compat_is_none() {
2008 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
2010 assert!(tmux_identity_from_data(&legacy).is_none());
2011 let partial = serde_json::json!({ "tmux_window_id": "@42" });
2013 assert!(tmux_identity_from_data(&partial).is_none());
2014 }
2015
2016 #[test]
2019 fn node_created_populates_tmux_identity() {
2020 let tmp = TempDir::new().unwrap();
2021 let run_id = "01jxsnap000000000000000000";
2022 let rid = RunId::parse_str(run_id).unwrap();
2023 let dir = crate::run_dir(tmp.path(), &rid);
2024 std::fs::create_dir_all(&dir).unwrap();
2025 let paths = RunPaths::new(dir, run_id).unwrap();
2026
2027 let mut ev = event(run_id);
2028 ev.seq = 2;
2029 ev.kind = "node.created".into();
2030 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2031 ev.data = serde_json::json!({
2032 "kind": "spinoff",
2033 "tmux_window": "🚀 wt/x",
2034 "tmux_socket": "/private/tmp/tmux-501/default",
2035 "tmux_session": "taskfleet",
2036 "tmux_window_id": "@42",
2037 });
2038 apply_event(&paths, &ev).expect("node.created applies");
2039 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2040 .unwrap()
2041 .unwrap();
2042 let id = n.tmux_identity.expect("qualified identity recorded");
2043 assert_eq!(id.session, "taskfleet");
2044 assert_eq!(id.window_id, "@42");
2045 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
2046
2047 let run2 = "02jxsnap000000000000000000";
2049 let rid2 = RunId::parse_str(run2).unwrap();
2050 let dir2 = crate::run_dir(tmp.path(), &rid2);
2051 std::fs::create_dir_all(&dir2).unwrap();
2052 let paths2 = RunPaths::new(dir2, run2).unwrap();
2053 let mut ev2 = event(run2);
2054 ev2.seq = 2;
2055 ev2.kind = "node.created".into();
2056 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2057 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
2058 apply_event(&paths2, &ev2).expect("legacy node.created applies");
2059 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
2060 .unwrap()
2061 .unwrap();
2062 assert!(n2.tmux_identity.is_none());
2063 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
2064 }
2065
2066 #[test]
2067 fn node_materialization_populates_missing_manifest_source_branch() {
2068 let tmp = TempDir::new().unwrap();
2069 let run_id = "01jxsnap000000000000000001";
2070 let rid = RunId::parse_str(run_id).unwrap();
2071 let dir = crate::run_dir(tmp.path(), &rid);
2072 std::fs::create_dir_all(&dir).unwrap();
2073 let paths = RunPaths::new(dir, run_id).unwrap();
2074
2075 let mut created = event(run_id);
2076 created.kind = "run.created".into();
2077 created.data = serde_json::json!({
2078 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2079 });
2080 apply_event(&paths, &created).unwrap();
2081 assert!(read_manifest_opt(&paths)
2082 .unwrap()
2083 .unwrap()
2084 .source_branch
2085 .is_none());
2086
2087 let mut node = event(run_id);
2088 node.seq = 2;
2089 node.kind = "node.created".into();
2090 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2091 node.data = serde_json::json!({
2092 "kind": "spinoff",
2093 "source_branch": "main",
2094 "worktree_path": "/tmp/wt/pending"
2095 });
2096 apply_event(&paths, &node).unwrap();
2097
2098 let manifest = read_manifest_opt(&paths).unwrap().unwrap();
2099 assert_eq!(manifest.status, Status::Pending);
2100 assert_eq!(manifest.source_branch.as_deref(), Some("main"));
2101 let node = read_n0001(&paths);
2102 assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
2103 }
2104
2105 #[test]
2106 fn node_materialization_preserves_explicit_manifest_source_branch() {
2107 let tmp = TempDir::new().unwrap();
2108 let run_id = "01jxsnap000000000000000002";
2109 let rid = RunId::parse_str(run_id).unwrap();
2110 let dir = crate::run_dir(tmp.path(), &rid);
2111 std::fs::create_dir_all(&dir).unwrap();
2112 let paths = RunPaths::new(dir, run_id).unwrap();
2113
2114 let mut created = event(run_id);
2115 created.kind = "run.created".into();
2116 created.data = serde_json::json!({
2117 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
2118 "source_branch": "release"
2119 });
2120 apply_event(&paths, &created).unwrap();
2121
2122 let mut node = event(run_id);
2123 node.seq = 2;
2124 node.kind = "node.created".into();
2125 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2126 node.data = serde_json::json!({
2127 "kind": "spinoff", "source_branch": "main"
2128 });
2129 apply_event(&paths, &node).unwrap();
2130
2131 assert_eq!(
2132 read_manifest_opt(&paths)
2133 .unwrap()
2134 .unwrap()
2135 .source_branch
2136 .as_deref(),
2137 Some("release")
2138 );
2139 }
2140
2141 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
2144 let rid = RunId::parse_str(run_id).unwrap();
2145 let dir = crate::run_dir(tmp.path(), &rid);
2146 std::fs::create_dir_all(&dir).unwrap();
2147 let paths = RunPaths::new(dir, run_id).unwrap();
2148
2149 let mut created = event(run_id);
2150 created.kind = "run.created".into();
2151 created.data = serde_json::json!({
2152 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2153 });
2154 apply_event(&paths, &created).expect("run.created applies");
2155
2156 let mut node = event(run_id);
2157 node.seq = 2;
2158 node.kind = "node.created".into();
2159 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2160 node.data = serde_json::json!({ "kind": "spinoff" });
2161 apply_event(&paths, &node).expect("node.created applies");
2162 paths
2163 }
2164
2165 fn read_n0001(paths: &RunPaths) -> Node {
2166 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
2167 .unwrap()
2168 .unwrap()
2169 }
2170
2171 #[test]
2172 fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
2173 let tmp = TempDir::new().unwrap();
2174 let run_id = "01jxwd0000000000000000000w";
2175 let paths = seed_run_with_node(&tmp, run_id);
2176 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2177 let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
2178 let mut open = event(run_id);
2179 open.seq = 3;
2180 open.ts = opened_at;
2181 open.kind = "node.awaiting_input".into();
2182 open.node_id = nid.clone();
2183 open.data = serde_json::json!({ "discussion_items": [{
2184 "topic": "Which scope?",
2185 "options": ["small", "large"],
2186 "recommended_default": "small"
2187 }] });
2188 apply_event(&paths, &open).unwrap();
2189 let first = read_n0001(&paths).awaiting_input.unwrap();
2190 assert_eq!(first.opened_at, opened_at);
2191 assert_eq!(first.event_seq, 3);
2192
2193 let mut duplicate = open.clone();
2195 duplicate.seq = 4;
2196 duplicate.ts = opened_at + chrono::Duration::hours(1);
2197 apply_event(&paths, &duplicate).unwrap();
2198 let still_first = read_n0001(&paths).awaiting_input.unwrap();
2199 assert_eq!(still_first.opened_at, opened_at);
2200 assert_eq!(still_first.event_seq, 3);
2201
2202 let mut stale = event(run_id);
2204 stale.seq = 5;
2205 stale.kind = "node.input_resolved".into();
2206 stale.node_id = nid.clone();
2207 stale.data = serde_json::json!({ "event_seq": 2 });
2208 apply_event(&paths, &stale).unwrap();
2209 assert!(read_n0001(&paths).awaiting_input.is_some());
2210
2211 let mut resolved = stale;
2212 resolved.seq = 6;
2213 resolved.data = serde_json::json!({ "event_seq": 3 });
2214 apply_event(&paths, &resolved).unwrap();
2215 assert!(read_n0001(&paths).awaiting_input.is_none());
2216 }
2217
2218 #[test]
2219 fn awaiting_input_rejects_missing_default_without_mutating_projection() {
2220 let tmp = TempDir::new().unwrap();
2221 let run_id = "01jxwd0000000000000000000x";
2222 let paths = seed_run_with_node(&tmp, run_id);
2223 let mut open = event(run_id);
2224 open.seq = 3;
2225 open.kind = "node.awaiting_input".into();
2226 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2227 open.data = serde_json::json!({ "discussion_items": [{
2228 "topic": "Which scope?", "options": ["small", "large"]
2229 }] });
2230 assert!(reduce_event_to_ops(&paths, &open).is_err());
2231 assert!(read_n0001(&paths).awaiting_input.is_none());
2232 }
2233
2234 #[test]
2235 fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
2236 let tmp = TempDir::new().unwrap();
2237 let run_id = "01jxwd0000000000000000000y";
2238 let paths = seed_run_with_node(&tmp, run_id);
2239 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2240
2241 let mut open = event(run_id);
2242 open.seq = 3;
2243 open.kind = "node.awaiting_input".into();
2244 open.node_id = nid.clone();
2245 open.data = serde_json::json!({ "discussion_items": [{
2246 "topic": "Which scope?", "options": ["small", "large"],
2247 "recommended_default": "small"
2248 }] });
2249 apply_event(&paths, &open).unwrap();
2250
2251 let mut malformed_duplicate = open.clone();
2252 malformed_duplicate.seq = 4;
2253 malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
2254 assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
2255
2256 let mut exited = event(run_id);
2257 exited.seq = 5;
2258 exited.kind = "worker.exited".into();
2259 exited.node_id = nid;
2260 exited.data = serde_json::json!({ "exit_code": 0 });
2261 apply_event(&paths, &exited).unwrap();
2262 let node = read_n0001(&paths);
2263 assert!(node.awaiting_input.is_none());
2264 assert!(node.worker_exit.is_some());
2265
2266 let mut delayed_open = open;
2267 delayed_open.seq = 6;
2268 assert!(reduce_event_to_ops(&paths, &delayed_open)
2269 .unwrap()
2270 .is_empty());
2271 }
2272
2273 #[test]
2274 fn input_resolved_requires_generation_even_when_nothing_is_open() {
2275 let tmp = TempDir::new().unwrap();
2276 let run_id = "01jxwd0000000000000000000z";
2277 let paths = seed_run_with_node(&tmp, run_id);
2278 let mut resolved = event(run_id);
2279 resolved.seq = 3;
2280 resolved.kind = "node.input_resolved".into();
2281 resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2282 resolved.data = serde_json::json!({});
2283 assert!(reduce_event_to_ops(&paths, &resolved).is_err());
2284 }
2285
2286 #[test]
2287 fn awaiting_input_default_must_be_one_of_options() {
2288 let tmp = TempDir::new().unwrap();
2289 let run_id = "01jxwd00000000000000000010";
2290 let paths = seed_run_with_node(&tmp, run_id);
2291 let mut open = event(run_id);
2292 open.seq = 3;
2293 open.kind = "node.awaiting_input".into();
2294 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2295 open.data = serde_json::json!({ "discussion_items": [{
2296 "topic": "Which scope?", "options": ["small", "large"],
2297 "recommended_default": "other"
2298 }] });
2299 assert!(reduce_event_to_ops(&paths, &open).is_err());
2300 }
2301
2302 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
2303 let mut ev = event(run_id);
2304 ev.seq = seq;
2305 ev.kind = KIND_MERGE_STARTED.into();
2306 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2307 ev.data = serde_json::json!({
2308 "op_id": op_id,
2309 "source_branch": "main",
2310 "worker_branch": "wt/worker",
2311 "expected_source_oid": expected,
2312 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2313 "base_sha": null,
2314 "driver_pid": 4242,
2315 "driver_pid_start_secs": null,
2316 "started_at": "2026-08-15T00:00:00Z",
2317 });
2318 ev
2319 }
2320
2321 #[test]
2324 fn merge_started_records_pending_transaction() {
2325 let tmp = TempDir::new().unwrap();
2326 let run_id = "01jxsnap000000000000000000";
2327 let paths = seed_run_with_node(&tmp, run_id);
2328
2329 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2330 let n = read_n0001(&paths);
2331 assert_eq!(
2332 n.status,
2333 Status::Pending,
2334 "recording a merge is not terminal"
2335 );
2336 let txn = n.pending_merge.expect("transaction recorded");
2337 assert_eq!(txn.op_id, "op-1");
2338 assert_eq!(txn.expected_source_oid, "aaa");
2339 }
2340
2341 #[test]
2344 fn merge_aborted_clears_matching_transaction_only() {
2345 let tmp = TempDir::new().unwrap();
2346 let run_id = "01jxsnap000000000000000000";
2347 let paths = seed_run_with_node(&tmp, run_id);
2348 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2349
2350 let mut stale = event(run_id);
2352 stale.seq = 4;
2353 stale.kind = KIND_MERGE_ABORTED.into();
2354 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2355 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2356 apply_event(&paths, &stale).unwrap();
2357 assert!(
2358 read_n0001(&paths).pending_merge.is_some(),
2359 "stale abort is a no-op"
2360 );
2361
2362 let mut abort = event(run_id);
2364 abort.seq = 5;
2365 abort.kind = KIND_MERGE_ABORTED.into();
2366 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2367 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2368 apply_event(&paths, &abort).unwrap();
2369 let n = read_n0001(&paths);
2370 assert!(
2371 n.pending_merge.is_none(),
2372 "matching abort clears the transaction"
2373 );
2374 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2375 }
2376
2377 #[test]
2380 fn terminal_report_clears_pending_merge() {
2381 let tmp = TempDir::new().unwrap();
2382 let run_id = "01jxsnap000000000000000000";
2383 let paths = seed_run_with_node(&tmp, run_id);
2384 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2385
2386 let mut report = event(run_id);
2387 report.seq = 4;
2388 report.kind = "node.report".into();
2389 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2390 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2391 apply_event(&paths, &report).unwrap();
2392 let n = read_n0001(&paths);
2393 assert_eq!(n.status, Status::Done);
2394 assert!(
2395 n.pending_merge.is_none(),
2396 "completed merge clears the transaction"
2397 );
2398 }
2399
2400 #[test]
2409 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2410 let tmp = TempDir::new().unwrap();
2411
2412 let drive = |run_id: &str, report_data: Value| -> Status {
2415 let paths = seed_run_with_node(&tmp, run_id);
2416 let mut fail = event(run_id);
2418 fail.seq = 3;
2419 fail.kind = "node.status".into();
2420 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2421 fail.data = serde_json::json!({ "status": "failed" });
2422 apply_event(&paths, &fail).unwrap();
2423 assert_eq!(read_n0001(&paths).status, Status::Failed);
2424 let mut report = event(run_id);
2426 report.seq = 4;
2427 report.kind = "node.report".into();
2428 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2429 report.data = report_data;
2430 apply_event(&paths, &report).unwrap();
2431 read_n0001(&paths).status
2432 };
2433
2434 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2436 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2437 assert_eq!(
2438 drive("01jxsnap000000000000000001", agent_forged),
2439 Status::Failed,
2440 "an Agent-origin report with a forged via must not be adopted"
2441 );
2442
2443 let malformed = serde_json::json!({
2445 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2446 });
2447 assert_eq!(
2448 drive("01jxsnap000000000000000002", malformed),
2449 Status::Failed,
2450 "a malformed origin must not re-unlock the legacy via adoption path"
2451 );
2452
2453 let mut run_merge = serde_json::json!({ "success": true });
2455 crate::ReportOrigin::RunMerge {
2456 op_id: Some("op-1".into()),
2457 worker_oid: Some("cafebabe".into()),
2458 }
2459 .stamp(&mut run_merge);
2460 assert_eq!(
2461 drive("01jxsnap000000000000000003", run_merge),
2462 Status::Done,
2463 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2464 );
2465
2466 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2469 assert_eq!(
2470 drive("01jxsnap000000000000000004", legacy),
2471 Status::Done,
2472 "a legacy via-only report (no origin field) is still adopted"
2473 );
2474 }
2475
2476 #[test]
2480 fn terminal_node_status_clears_pending_merge() {
2481 let tmp = TempDir::new().unwrap();
2482 let run_id = "01jxsnap000000000000000000";
2483 let paths = seed_run_with_node(&tmp, run_id);
2484 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2485 assert!(read_n0001(&paths).pending_merge.is_some());
2486
2487 let mut status = event(run_id);
2488 status.seq = 4;
2489 status.kind = "node.status".into();
2490 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2491 status.data = serde_json::json!({ "status": "failed" });
2492 apply_event(&paths, &status).unwrap();
2493 let n = read_n0001(&paths);
2494 assert_eq!(n.status, Status::Failed);
2495 assert!(
2496 n.pending_merge.is_none(),
2497 "terminal status clears the transaction"
2498 );
2499 }
2500
2501 #[test]
2505 fn supervisor_attached_sets_supervisor_pid() {
2506 let tmp = TempDir::new().unwrap();
2507 let run_id = "01jxsnap000000000000000000";
2508 let paths = seed_run_with_node(&tmp, run_id);
2509 assert_eq!(read_n0001(&paths).supervisor_pid, None);
2510
2511 let mut ev = event(run_id);
2512 ev.seq = 3;
2513 ev.kind = "supervisor.attached".into();
2514 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2515 ev.data = serde_json::json!({ "pid": 47820 });
2516 apply_event(&paths, &ev).expect("supervisor.attached applies");
2517 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2518 }
2519
2520 #[test]
2523 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2524 let tmp = TempDir::new().unwrap();
2525 let run_id = "01jxsnap000000000000000000";
2526 let paths = seed_run_with_node(&tmp, run_id);
2527
2528 let mut ev = event(run_id);
2529 ev.kind = "supervisor.attached".into();
2530 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2531
2532 ev.seq = 3;
2533 ev.data = serde_json::json!({ "pid": 100 });
2534 apply_event(&paths, &ev).expect("first attach applies");
2535 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2536
2537 ev.seq = 4;
2539 ev.data = serde_json::json!({ "pid": 200 });
2540 apply_event(&paths, &ev).expect("second attach applies");
2541 let after_second = read_n0001(&paths);
2542 assert_eq!(after_second.supervisor_pid, Some(200));
2543
2544 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2547 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2548 apply_event(&paths, &ev).expect("replay applies as no-op");
2549 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2550 }
2551
2552 #[test]
2555 fn supervisor_cursor_advanced_sets_report_cursor() {
2556 let tmp = TempDir::new().unwrap();
2557 let run_id = "01jxsnap000000000000000000";
2558 let paths = seed_run_with_node(&tmp, run_id);
2559 let child = "02jxsnap000000000000000000";
2560
2561 let mut ev = event(run_id);
2562 ev.seq = 3;
2563 ev.kind = "supervisor.cursor_advanced".into();
2564 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2565 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2566 apply_event(&paths, &ev).expect("cursor_advanced applies");
2567
2568 let n = read_n0001(&paths);
2569 assert_eq!(
2570 n.last_processed_report_seq_by_child.get(child),
2571 Some(&Value::from(7u64))
2572 );
2573 }
2574
2575 #[test]
2580 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2581 let tmp = TempDir::new().unwrap();
2582 let run_id = "01jxsnap000000000000000000";
2583 let paths = seed_run_with_node(&tmp, run_id);
2584 let child_a = "02jxsnap000000000000000000";
2585 let child_b = "03jxsnap000000000000000000";
2586
2587 let mut ev = event(run_id);
2588 ev.kind = "supervisor.cursor_advanced".into();
2589 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2590
2591 ev.seq = 3;
2592 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2593 apply_event(&paths, &ev).expect("seq 5 applies");
2594
2595 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2597 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2598
2599 ev.seq = 4;
2601 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2602 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2603 assert!(ops.is_empty(), "older seq must plan no ops");
2604 apply_event(&paths, &ev).expect("older seq applies as no-op");
2605 assert_eq!(
2606 read_n0001(&paths)
2607 .last_processed_report_seq_by_child
2608 .get(child_a),
2609 Some(&Value::from(5u64))
2610 );
2611
2612 ev.seq = 5;
2614 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2615 apply_event(&paths, &ev).expect("higher seq applies");
2616 ev.seq = 6;
2617 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2618 apply_event(&paths, &ev).expect("second child applies");
2619
2620 let n = read_n0001(&paths);
2621 assert_eq!(
2622 n.last_processed_report_seq_by_child.get(child_a),
2623 Some(&Value::from(9u64))
2624 );
2625 assert_eq!(
2626 n.last_processed_report_seq_by_child.get(child_b),
2627 Some(&Value::from(1u64))
2628 );
2629 }
2630
2631 #[test]
2634 fn supervisor_state_events_reject_malformed_payloads() {
2635 let tmp = TempDir::new().unwrap();
2636 let run_id = "01jxsnap000000000000000000";
2637 let paths = seed_run_with_node(&tmp, run_id);
2638 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2639
2640 let mut ev = event(run_id);
2642 ev.seq = 3;
2643 ev.kind = "supervisor.attached".into();
2644 ev.node_id = nid.clone();
2645 ev.data = serde_json::json!({});
2646 assert!(matches!(
2647 reduce_event_to_ops(&paths, &ev),
2648 Err(Error::CorruptEventLog { .. })
2649 ));
2650
2651 ev.node_id = None;
2653 ev.data = serde_json::json!({ "pid": 1 });
2654 assert!(matches!(
2655 reduce_event_to_ops(&paths, &ev),
2656 Err(Error::CorruptEventLog { .. })
2657 ));
2658
2659 let mut ev2 = event(run_id);
2661 ev2.seq = 4;
2662 ev2.kind = "supervisor.cursor_advanced".into();
2663 ev2.node_id = nid.clone();
2664 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
2665 assert!(matches!(
2666 reduce_event_to_ops(&paths, &ev2),
2667 Err(Error::CorruptEventLog { .. })
2668 ));
2669
2670 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2672 assert!(matches!(
2673 reduce_event_to_ops(&paths, &ev2),
2674 Err(Error::CorruptEventLog { .. })
2675 ));
2676 }
2677
2678 #[test]
2685 fn removed_or_garbage_kind_in_created_events_is_rejected() {
2686 let tmp = TempDir::new().unwrap();
2687 let run_id = "01jxsnap000000000000000000";
2688 let rid = RunId::parse_str(run_id).unwrap();
2689 let dir = crate::run_dir(tmp.path(), &rid);
2690 std::fs::create_dir_all(&dir).unwrap();
2691 let paths = RunPaths::new(dir, run_id).unwrap();
2692
2693 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2694 let mut ev = event(run_id);
2695 ev.kind = "run.created".into();
2696 ev.node_id = None;
2697 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2698 assert!(
2699 matches!(
2700 reduce_event_to_ops(&paths, &ev),
2701 Err(Error::CorruptEventLog { .. })
2702 ),
2703 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2704 );
2705 }
2706
2707 let mut ok = event(run_id);
2710 ok.kind = "run.created".into();
2711 ok.node_id = None;
2712 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2713 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2714 }
2715
2716 #[cfg(unix)]
2725 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2726 use std::os::unix::fs::MetadataExt;
2727 let mut consider = vec![paths.manifest()];
2728 for dir in [paths.nodes_dir()] {
2729 if let Ok(rd) = std::fs::read_dir(&dir) {
2730 for ent in rd.flatten() {
2731 let p = ent.path();
2732 if p.extension().and_then(|s| s.to_str()) == Some("json") {
2733 consider.push(p);
2734 }
2735 }
2736 }
2737 }
2738 let mut map = std::collections::BTreeMap::new();
2739 for p in consider {
2740 if let Ok(md) = std::fs::symlink_metadata(&p) {
2741 if md.file_type().is_file() {
2742 map.insert(p, md.ino());
2743 }
2744 }
2745 }
2746 map
2747 }
2748
2749 #[cfg(unix)]
2758 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2759 use std::collections::BTreeSet;
2760 let before = projection_inodes(paths);
2761 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2762 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2763 .into_iter()
2764 .collect();
2765 apply_event(paths, ev)
2766 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2767 let after = projection_inodes(paths);
2768 let touched: BTreeSet<PathBuf> = after
2769 .iter()
2770 .filter(|(p, ino)| before.get(*p) != Some(*ino))
2771 .map(|(p, _)| p.clone())
2772 .collect();
2773 assert_eq!(
2774 planned, touched,
2775 "kind={}: plan_projections must name exactly the files apply_event writes",
2776 ev.kind
2777 );
2778 if expect_writes {
2779 assert!(
2780 !touched.is_empty(),
2781 "kind={}: expected this event to write at least one projection",
2782 ev.kind
2783 );
2784 }
2785 }
2786
2787 #[cfg(unix)]
2794 #[test]
2795 fn plan_projections_matches_apply_for_every_kind() {
2796 let tmp = TempDir::new().unwrap();
2797 let run_id = "01jxsnap000000000000000000";
2798 let rid = RunId::parse_str(run_id).unwrap();
2799 let dir = crate::run_dir(tmp.path(), &rid);
2800 std::fs::create_dir_all(&dir).unwrap();
2801 let paths = RunPaths::new(dir, run_id).unwrap();
2802 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2803 let child = "02jxsnap000000000000000000";
2804
2805 let mut next_seq = 0u64;
2807 let mut at = |kind: &str, node_id, data| {
2808 next_seq += 1;
2809 Event {
2810 ts: Utc::now(),
2811 seq: next_seq,
2812 kind: kind.into(),
2813 run_id: rid.clone(),
2814 node_id,
2815 idempotency_key: None,
2816 data,
2817 }
2818 };
2819
2820 assert_plan_matches_apply(
2822 &paths,
2823 &at(
2824 "run.created",
2825 None,
2826 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2827 ),
2828 true,
2829 );
2830 assert_plan_matches_apply(
2832 &paths,
2833 &at(
2834 "run.status",
2835 None,
2836 serde_json::json!({ "status": "running" }),
2837 ),
2838 true,
2839 );
2840 assert_plan_matches_apply(
2842 &paths,
2843 &at(
2844 "node.created",
2845 nid(),
2846 serde_json::json!({ "kind": "spinoff" }),
2847 ),
2848 true,
2849 );
2850 assert_plan_matches_apply(
2852 &paths,
2853 &at(
2854 "node.status",
2855 nid(),
2856 serde_json::json!({ "status": "running" }),
2857 ),
2858 true,
2859 );
2860 assert_plan_matches_apply(
2862 &paths,
2863 &at(
2864 "supervisor.attached",
2865 nid(),
2866 serde_json::json!({ "pid": 4242 }),
2867 ),
2868 true,
2869 );
2870 assert_plan_matches_apply(
2872 &paths,
2873 &at(
2874 "supervisor.cursor_advanced",
2875 nid(),
2876 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2877 ),
2878 true,
2879 );
2880 assert_plan_matches_apply(
2882 &paths,
2883 &at(
2884 "child.spawned",
2885 nid(),
2886 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2887 ),
2888 true,
2889 );
2890 assert_plan_matches_apply(
2892 &paths,
2893 &at("node.report", nid(), serde_json::json!({ "success": true })),
2894 true,
2895 );
2896 assert_plan_matches_apply(
2899 &paths,
2900 &at(
2901 "node.status",
2902 nid(),
2903 serde_json::json!({ "status": "failed" }),
2904 ),
2905 false,
2906 );
2907 for kind in [
2909 "supervisor.exited",
2910 "orchestrator.decision",
2911 "discuss.critical",
2912 "cleanup.window_missing",
2913 ] {
2914 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
2915 }
2916 }
2917
2918 #[test]
2922 fn worker_exited_records_clean_exit_without_transitioning_status() {
2923 let tmp = TempDir::new().unwrap();
2924 let run_id = "01jxsnap000000000000000000";
2925 let paths = bootstrap_retry_node(&tmp, run_id);
2926
2927 let mut ev = event(run_id);
2928 ev.seq = 3;
2929 ev.kind = "worker.exited".into();
2930 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2931 ev.data = serde_json::json!({ "exit_code": 0 });
2932 apply_event(&paths, &ev).expect("worker.exited applies");
2933
2934 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2935 .unwrap()
2936 .unwrap();
2937 let exit = n.worker_exit.expect("worker_exit recorded");
2938 assert_eq!(exit.code, Some(0));
2939 assert_eq!(exit.signal, None);
2940 assert!(exit.is_clean());
2941 assert_eq!(
2942 n.status,
2943 Status::Pending,
2944 "the exit fact never transitions status"
2945 );
2946 }
2947
2948 #[test]
2952 fn worker_exited_records_signal_and_is_first_write_wins() {
2953 let tmp = TempDir::new().unwrap();
2954 let run_id = "01jxsnap000000000000000000";
2955 let paths = bootstrap_retry_node(&tmp, run_id);
2956 let nid = NodeId::parse_str("n-0001").unwrap();
2957
2958 let mut ev = event(run_id);
2959 ev.seq = 3;
2960 ev.kind = "worker.exited".into();
2961 ev.node_id = Some(nid.clone());
2962 ev.data = serde_json::json!({ "signal": 9 });
2963 apply_event(&paths, &ev).expect("worker.exited applies");
2964
2965 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2966 let exit = n.worker_exit.expect("worker_exit recorded");
2967 assert_eq!(exit.signal, Some(9));
2968 assert!(exit.is_failure());
2969
2970 let mut dup = event(run_id);
2973 dup.seq = 4;
2974 dup.kind = "worker.exited".into();
2975 dup.node_id = Some(nid.clone());
2976 dup.data = serde_json::json!({ "exit_code": 0 });
2977 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
2978 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
2979 assert_eq!(
2980 n2.worker_exit.unwrap().signal,
2981 Some(9),
2982 "first-write-wins: the replayed exit must not overwrite the recorded fact"
2983 );
2984 }
2985
2986 #[test]
2990 fn worker_exited_without_code_or_signal_is_corrupt() {
2991 let tmp = TempDir::new().unwrap();
2992 let run_id = "01jxsnap000000000000000000";
2993 let paths = bootstrap_retry_node(&tmp, run_id);
2994
2995 let mut ev = event(run_id);
2996 ev.seq = 3;
2997 ev.kind = "worker.exited".into();
2998 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2999 ev.data = serde_json::json!({});
3000 match reduce_event_to_ops(&paths, &ev) {
3001 Err(Error::CorruptEventLog { .. }) => {}
3002 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
3003 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3004 }
3005
3006 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
3009 match reduce_event_to_ops(&paths, &ev) {
3010 Err(Error::CorruptEventLog { .. }) => {}
3011 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
3012 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3013 }
3014 }
3015
3016 #[test]
3022 fn node_death_observed_records_first_death_first_write_wins() {
3023 let tmp = TempDir::new().unwrap();
3024 let run_id = "01jxsnap000000000000000000";
3025 let paths = bootstrap_retry_node(&tmp, run_id);
3026 let nid = NodeId::parse_str("n-0001").unwrap();
3027
3028 let mut ev = event(run_id);
3029 ev.seq = 3;
3030 ev.kind = "node.death_observed".into();
3031 ev.node_id = Some(nid.clone());
3032 ev.data = serde_json::json!({});
3033 apply_event(&paths, &ev).expect("node.death_observed applies");
3034 let first = read_node_opt(&paths, &nid)
3035 .unwrap()
3036 .unwrap()
3037 .first_death_at
3038 .expect("first_death_at recorded");
3039 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
3040
3041 let mut later = event(run_id);
3043 later.seq = 4;
3044 later.kind = "node.death_observed".into();
3045 later.node_id = Some(nid.clone());
3046 later.ts = ev.ts + chrono::Duration::seconds(30);
3047 later.data = serde_json::json!({});
3048 apply_event(&paths, &later).expect("re-observation applies as no-op");
3049 assert_eq!(
3050 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3051 Some(first),
3052 "first-write-wins: a re-observation must not reset the anchor"
3053 );
3054 }
3055
3056 #[test]
3061 fn node_death_observed_noop_when_worker_exit_present() {
3062 let tmp = TempDir::new().unwrap();
3063 let run_id = "01jxsnap000000000000000000";
3064 let paths = bootstrap_retry_node(&tmp, run_id);
3065 let nid = NodeId::parse_str("n-0001").unwrap();
3066
3067 let mut exit = event(run_id);
3069 exit.seq = 3;
3070 exit.kind = "worker.exited".into();
3071 exit.node_id = Some(nid.clone());
3072 exit.data = serde_json::json!({ "exit_code": 0 });
3073 apply_event(&paths, &exit).unwrap();
3074
3075 let mut death = event(run_id);
3077 death.seq = 4;
3078 death.kind = "node.death_observed".into();
3079 death.node_id = Some(nid.clone());
3080 death.data = serde_json::json!({});
3081 apply_event(&paths, &death).expect("applies as no-op");
3082 assert_eq!(
3083 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3084 None,
3085 "a told worker.exited makes the crash backstop moot; no anchor recorded"
3086 );
3087 }
3088
3089 #[test]
3094 fn node_retry_clears_worker_exit() {
3095 let tmp = TempDir::new().unwrap();
3096 let run_id = "01jxsnap000000000000000000";
3097 let paths = bootstrap_retry_node(&tmp, run_id);
3098 let nid = NodeId::parse_str("n-0001").unwrap();
3099
3100 let mut exit = event(run_id);
3102 exit.seq = 3;
3103 exit.kind = "worker.exited".into();
3104 exit.node_id = Some(nid.clone());
3105 exit.data = serde_json::json!({ "exit_code": 7 });
3106 apply_event(&paths, &exit).unwrap();
3107 assert!(read_node_opt(&paths, &nid)
3108 .unwrap()
3109 .unwrap()
3110 .worker_exit
3111 .is_some());
3112
3113 let mut retry = event(run_id);
3115 retry.seq = 4;
3116 retry.kind = "node.retry".into();
3117 retry.node_id = Some(nid.clone());
3118 retry.data = serde_json::json!({
3119 "attempt": 1,
3120 "reason": "agent-died",
3121 "branch": "wt/foo",
3122 "worktree_path": "/tmp/new-wt",
3123 "agent_pid": 222,
3124 });
3125 apply_event(&paths, &retry).unwrap();
3126
3127 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3128 assert!(
3129 n.worker_exit.is_none(),
3130 "node.retry must clear the previous attempt's worker_exit"
3131 );
3132 assert_eq!(
3133 n.status,
3134 Status::Pending,
3135 "retry returns the node to Pending"
3136 );
3137 }
3138}