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 "worker.display.retained" => reduce_worker_display_retained(paths, ev),
302 "worker.display.expired" => reduce_worker_display_expired(paths, ev),
303 "worker.display.unavailable" => reduce_worker_display_unavailable(paths, ev),
304 "node.death_observed" => reduce_node_death_observed(paths, ev),
305 "node.awaiting_input" => reduce_node_awaiting_input(paths, ev),
306 "node.input_resolved" => reduce_node_input_resolved(paths, ev),
307 KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
308 KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
309 "child.spawned" => reduce_child_spawned(paths, ev),
310 "supervisor.attached" => reduce_supervisor_attached(paths, ev),
311 "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
312 "supervisor.exited" => Ok(vec![]),
313 "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
321 "run.notified" | "run.awaiting_input_notified" => Ok(vec![]),
328 "cleanup.window_missing"
353 | "cleanup.worktree_missing"
354 | "cleanup.branch_remove_failed"
355 | "cleanup.branch_preserved"
356 | "cleanup.discard_authorized"
357 | "cleanup.session_killed"
358 | "cleanup.session_retained" => Ok(vec![]),
359 "supervisor.child_id_quarantined" => Ok(vec![]),
369 _ => Ok(vec![]),
370 }
371}
372
373fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
379 match op {
380 ProjectionOp::Manifest(_) => paths.manifest(),
381 ProjectionOp::Node(n) => paths.node(&n.node_id),
382 }
383}
384
385pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
407 let ops = reduce_event_to_ops(paths, event)?;
408 Ok(ops.iter().map(|op| op_path(paths, op)).collect())
409}
410
411pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
427 let ops = reduce_event_to_ops(paths, ev)?;
428 commit_ops(paths, ops)
429}
430
431pub fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
437 reduce_event_to_ops(paths, ev).map(|_| ())
438}
439
440fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
444 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
445 path: events_path.to_path_buf(),
446 reason: format!(
447 "event seq={} kind={} missing top-level `node_id`",
448 ev.seq, ev.kind
449 ),
450 })
451}
452
453fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
454 if let Some(existing) = read_manifest_opt(paths)? {
458 if existing.run_id != ev.run_id {
459 return Err(Error::CorruptEventLog {
460 path: paths.manifest(),
461 reason: format!(
462 "run.created run_id={} conflicts with existing manifest run_id={}",
463 ev.run_id, existing.run_id
464 ),
465 });
466 }
467 return Ok(vec![]);
468 }
469 let events_path = paths.events();
470 let d = &ev.data;
471 let kind =
472 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
473 path: events_path.clone(),
474 reason: "run.created missing/invalid `kind`".into(),
475 })?;
476 let lifecycle: Lifecycle = serde_json::from_value(
477 d.get("lifecycle").cloned().unwrap_or(Value::Null),
478 )
479 .map_err(|_| Error::CorruptEventLog {
480 path: events_path.clone(),
481 reason: "run.created missing/invalid `lifecycle`".into(),
482 })?;
483 let title = want_str(&events_path, ev, d, "title")?.to_string();
484 let agent_selection: Option<crate::schema::AgentSelection> = d
485 .get("agent_selection")
486 .cloned()
487 .map(serde_json::from_value)
488 .transpose()
489 .map_err(|e| Error::CorruptEventLog {
490 path: events_path.clone(),
491 reason: format!("run.created invalid `agent_selection`: {e}"),
492 })?;
493 if let Some(selection) = &agent_selection {
494 selection
495 .validate()
496 .map_err(|reason| Error::CorruptEventLog {
497 path: events_path.clone(),
498 reason: format!("run.created invalid `agent_selection`: {reason}"),
499 })?;
500 }
501 let m = Manifest {
502 schema_version: STATE_SCHEMA_VERSION,
503 applied_seq: 0,
506 run_id: paths.run_id.clone(),
508 kind,
509 lifecycle,
510 title,
511 status: Status::Pending,
512 created_at: ev.ts,
513 updated_at: ev.ts,
514 source_repo: d
515 .get("source_repo")
516 .and_then(Value::as_str)
517 .map(str::to_string),
518 source_branch: d
519 .get("source_branch")
520 .and_then(Value::as_str)
521 .map(str::to_string),
522 worktree_root: d
523 .get("worktree_root")
524 .and_then(Value::as_str)
525 .map(str::to_string),
526 managed_tmux_session: d
527 .get("managed_tmux_session")
528 .and_then(Value::as_str)
529 .map(str::to_string),
530 tmux_retention: retention_policy_from_data(&events_path, d)?,
531 notify_cmd: d
532 .get("notify_cmd")
533 .and_then(Value::as_str)
534 .map(str::to_string),
535 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
536 agent_selection,
537 node_count: 0,
538 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
539 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
540 };
541 Ok(vec![ProjectionOp::Manifest(m)])
542}
543
544fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
545 let mut m = match read_manifest_opt(paths)? {
546 Some(m) => m,
547 None => return Ok(vec![]),
548 };
549 let new_status = require_status(ev, paths.events())?;
550 if m.status.is_terminal() {
556 let recovery_seq = ev
557 .data
558 .get("recovery_merge_report_seq")
559 .and_then(Value::as_u64);
560 let children_successful =
561 ev.data.get("children_successful").and_then(Value::as_bool) == Some(true);
562 let recovery_allowed = m.status == Status::Failed
563 && new_status == Status::Done
564 && children_successful
565 && recovery_seq.is_some_and(|wanted| {
566 crate::cancel::read_node_status_facts(paths, Some(ev.seq))
567 .ok()
568 .filter(|facts| {
569 crate::aggregate_terminal_status(facts.iter().map(|fact| fact.status))
570 == Some(Status::Done)
571 })
572 .is_some_and(|facts| {
573 facts
574 .iter()
575 .any(|fact| fact.confirmed_merge_seq == Some(wanted))
576 })
577 });
578 if !recovery_allowed {
579 trace_terminal_noop(ev, m.status, new_status);
580 return Ok(vec![]);
581 }
582 }
583 if m.status == new_status {
584 return Ok(vec![]);
585 }
586 m.status = new_status;
587 m.updated_at = ev.ts;
588 Ok(vec![ProjectionOp::Manifest(m)])
589}
590
591fn retention_policy_from_data(
602 events_path: &Path,
603 data: &Value,
604) -> Result<Option<Box<crate::schema::TmuxRetentionPolicy>>> {
605 let Some(value) = data.get("tmux_retention") else {
606 return Ok(None);
607 };
608 let policy: crate::schema::TmuxRetentionPolicy = serde_json::from_value(value.clone())
609 .map_err(|e| Error::CorruptEventLog {
610 path: events_path.to_path_buf(),
611 reason: format!("run.created invalid `tmux_retention`: {e}"),
612 })?;
613 if !policy.persistent
614 || policy.completed_window_ttl_secs == 0
615 || policy.completed_window_max == 0
616 {
617 return Err(Error::CorruptEventLog {
618 path: events_path.to_path_buf(),
619 reason: "run.created tmux_retention must be persistent with positive ttl/max".into(),
620 });
621 }
622 Ok(Some(Box::new(policy)))
623}
624
625fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
626 let nonempty = |key| {
627 d.get(key)
628 .and_then(Value::as_str)
629 .map(str::trim)
630 .filter(|s| !s.is_empty())
631 .map(str::to_string)
632 };
633 let session = nonempty("tmux_session")?;
634 let window_id = nonempty("tmux_window_id")?;
635 Some(TmuxIdentity {
636 socket: nonempty("tmux_socket"),
637 session,
638 window_id,
639 pane_id: nonempty("tmux_pane_id"),
642 server_pid: d
643 .get("tmux_server_pid")
644 .and_then(Value::as_u64)
645 .and_then(|value| u32::try_from(value).ok()),
646 server_pid_start_secs: d.get("tmux_server_pid_start_secs").and_then(Value::as_u64),
647 server_marker: nonempty("tmux_server_marker"),
648 })
649}
650
651fn worker_evidence_from_spawn_data(
652 events_path: &Path,
653 ev: &Event,
654 data: &Value,
655) -> Result<Option<WorkerEvidence>> {
656 let Some(session_id) = data.get("pi_session_id") else {
657 if data.get("pi_session_path").is_some() || data.get("pi_session_cwd").is_some() {
658 return Err(Error::CorruptEventLog {
659 path: events_path.to_path_buf(),
660 reason: format!(
661 "event seq={} Pi evidence fields must be all-or-none",
662 ev.seq
663 ),
664 });
665 }
666 return Ok(None);
667 };
668 if session_id.is_null() {
669 if data
670 .get("pi_session_path")
671 .is_some_and(|value| !value.is_null())
672 || data
673 .get("pi_session_cwd")
674 .is_some_and(|value| !value.is_null())
675 {
676 return Err(Error::CorruptEventLog {
677 path: events_path.to_path_buf(),
678 reason: format!(
679 "event seq={} Pi evidence fields must be all-or-none",
680 ev.seq
681 ),
682 });
683 }
684 return Ok(None);
685 }
686 let session_id = session_id
687 .as_str()
688 .filter(|v| !v.is_empty())
689 .ok_or_else(|| Error::CorruptEventLog {
690 path: events_path.to_path_buf(),
691 reason: format!(
692 "event seq={} pi_session_id must be a non-empty string",
693 ev.seq
694 ),
695 })?;
696 if session_id.len() != 36
697 || !session_id.bytes().enumerate().all(|(index, byte)| {
698 if matches!(index, 8 | 13 | 18 | 23) {
699 byte == b'-'
700 } else {
701 byte.is_ascii_hexdigit()
702 }
703 })
704 {
705 return Err(Error::CorruptEventLog {
706 path: events_path.to_path_buf(),
707 reason: format!("event seq={} pi_session_id must be a UUID", ev.seq),
708 });
709 }
710 let original_cwd = want_str(events_path, ev, data, "pi_session_cwd")?;
711 let live_session_path = want_str(events_path, ev, data, "pi_session_path")?;
712 let expected_live_path = format!(
713 ".creating/pi-sessions/{}/pi-session-{session_id}.jsonl",
714 ev.run_id.as_str()
715 );
716 if live_session_path != expected_live_path {
717 return Err(Error::CorruptEventLog {
718 path: events_path.to_path_buf(),
719 reason: format!(
720 "event seq={} pi_session_path is not the canonical state-relative path",
721 ev.seq
722 ),
723 });
724 }
725 let attempt = match data.get("attempt") {
726 None | Some(Value::Null) => 0,
727 Some(value) => value
728 .as_u64()
729 .and_then(|raw| u32::try_from(raw).ok())
730 .ok_or_else(|| Error::CorruptEventLog {
731 path: events_path.to_path_buf(),
732 reason: format!("event seq={} attempt must be a u32", ev.seq),
733 })?,
734 };
735 Ok(Some(WorkerEvidence {
736 attempt,
737 session_id: session_id.to_string(),
738 original_cwd: original_cwd.to_string(),
739 live_session_path: live_session_path.to_string(),
740 status: EvidenceStatus::Pending,
741 transcript_path: None,
742 resume_path: None,
743 pane_path: None,
744 report_path: None,
745 transcript_sha256: None,
746 error: None,
747 }))
748}
749
750fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
751 let events_path = paths.events();
752 let node_id = require_envelope_node_id(&events_path, ev)?;
755 let is_default_node = node_id.as_str() == "n-0001";
756 if read_node_opt(paths, &node_id)?.is_some() {
758 return Ok(vec![]);
759 }
760 let d = &ev.data;
761 let kind =
762 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
763 path: events_path.clone(),
764 reason: format!(
765 "event seq={} kind=node.created missing/invalid `kind`",
766 ev.seq
767 ),
768 })?;
769 let n = Node {
770 schema_version: STATE_SCHEMA_VERSION,
771 node_id,
772 run_id: paths.run_id.clone(),
774 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
775 kind,
776 status: Status::Pending,
777 task: d.get("task").and_then(Value::as_str).map(str::to_string),
778 worktree_path: d
779 .get("worktree_path")
780 .and_then(Value::as_str)
781 .map(str::to_string),
782 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
783 base_sha: d
784 .get("base_sha")
785 .and_then(Value::as_str)
786 .filter(|s| !s.is_empty())
787 .map(str::to_string),
788 tmux_window: d
789 .get("tmux_window")
790 .and_then(Value::as_str)
791 .map(str::to_string),
792 tmux_identity: tmux_identity_from_data(d).map(Box::new),
793 evidence: worker_evidence_from_spawn_data(&events_path, ev, d)?,
794 retained_display: None,
795 retention_unavailable: None,
796 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
797 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
798 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
799 children: Vec::new(),
800 started_at: Some(ev.ts),
801 updated_at: ev.ts,
802 last_report: None,
803 last_processed_report_seq_by_child: serde_json::Map::default(),
804 retry_attempts: 0,
805 worker_exit: None,
806 pending_merge: None,
807 first_death_at: None,
808 awaiting_input: None,
809 };
810 let mut ops = vec![ProjectionOp::Node(n)];
811 if let Some(mut m) = read_manifest_opt(paths)? {
812 if is_default_node && m.source_branch.is_none() {
817 m.source_branch = d
818 .get("source_branch")
819 .and_then(Value::as_str)
820 .filter(|branch| !branch.is_empty())
821 .map(str::to_string);
822 }
823 m.updated_at = ev.ts;
828 ops.push(ProjectionOp::Manifest(m));
829 }
830 Ok(ops)
831}
832
833fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
852 let events_path = paths.events();
853 let node_id = require_envelope_node_id(&events_path, ev)?;
854 let mut n = match read_node_opt(paths, &node_id)? {
855 Some(n) => n,
856 None => return Ok(vec![]),
857 };
858 if n.status.is_terminal() {
861 tracing::debug!(
862 target: "taskfleet_core::reducer",
863 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
864 "no-op: node.retry against terminal node"
865 );
866 return Ok(vec![]);
867 }
868 let d = &ev.data;
869 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
872 n.base_sha = d
873 .get("base_sha")
874 .and_then(Value::as_str)
875 .filter(|s| !s.is_empty())
876 .map(str::to_string);
877 n.worktree_path = d
878 .get("worktree_path")
879 .and_then(Value::as_str)
880 .map(str::to_string);
881 n.tmux_window = d
882 .get("tmux_window")
883 .and_then(Value::as_str)
884 .map(str::to_string);
885 n.tmux_identity = tmux_identity_from_data(d).map(Box::new);
886 n.evidence = worker_evidence_from_spawn_data(&events_path, ev, d)?;
887 n.retained_display = None;
888 n.retention_unavailable = None;
889 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
890 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
891 n.status = Status::Pending;
892 n.started_at = Some(ev.ts);
893 n.updated_at = ev.ts;
894 n.last_report = None;
895 n.pending_merge = None;
901 n.worker_exit = None;
906 n.first_death_at = None;
911 n.awaiting_input = None;
914 n.retry_attempts = d
922 .get("attempt")
923 .and_then(Value::as_u64)
924 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
925 Ok(vec![ProjectionOp::Node(n)])
926}
927
928fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
929 let events_path = paths.events();
930 let node_id = require_envelope_node_id(&events_path, ev)?;
931 let mut n = match read_node_opt(paths, &node_id)? {
932 Some(n) => n,
933 None => return Ok(vec![]),
934 };
935 let new_status = require_status(ev, events_path)?;
936 if n.status.is_terminal() {
939 trace_terminal_noop(ev, n.status, new_status);
940 return Ok(vec![]);
941 }
942 if n.status == new_status {
943 return Ok(vec![]);
944 }
945 n.status = new_status;
946 if new_status.is_terminal() {
952 n.pending_merge = None;
953 n.awaiting_input = None;
954 }
955 n.updated_at = ev.ts;
956 Ok(vec![ProjectionOp::Node(n)])
957}
958
959fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
960 let events_path = paths.events();
961 let node_id = require_envelope_node_id(&events_path, ev)?;
962 let mut n = match read_node_opt(paths, &node_id)? {
963 Some(n) => n,
964 None => return Ok(vec![]),
965 };
966 if n.status.is_terminal() {
978 if ReportOrigin::permits_terminal_merge_recovery(n.status, &ev.data) {
1011 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
1012 return Ok(vec![]);
1013 }
1014 tracing::info!(
1015 target: "taskfleet_core::reducer",
1016 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
1017 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
1018 );
1019 n.last_report = Some(ev.data.clone());
1020 n.status = Status::Done;
1024 n.pending_merge = None;
1028 n.awaiting_input = None;
1029 n.updated_at = ev.ts;
1030 return Ok(vec![ProjectionOp::Node(n)]);
1031 }
1032 tracing::debug!(
1033 target: "taskfleet_core::reducer",
1034 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
1035 "no-op: node.report against terminal node"
1036 );
1037 return Ok(vec![]);
1038 }
1039 let new_status = report_terminal_status(&events_path, ev)?;
1047 n.last_report = Some(ev.data.clone());
1048 n.status = new_status;
1049 n.awaiting_input = None;
1053 n.pending_merge = None;
1058 n.updated_at = ev.ts;
1059 Ok(vec![ProjectionOp::Node(n)])
1060}
1061
1062fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1077 let events_path = paths.events();
1078 let node_id = require_envelope_node_id(&events_path, ev)?;
1079 let attempt = match ev.data.get("attempt") {
1080 None | Some(Value::Null) => 0,
1081 Some(_) => evidence_attempt(&events_path, ev)?,
1082 };
1083 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
1084 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
1085 match (code, signal) {
1090 (Some(_), None) | (None, Some(_)) => {}
1091 _ => {
1092 return Err(Error::CorruptEventLog {
1093 path: events_path,
1094 reason: format!(
1095 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
1096 ev.seq
1097 ),
1098 });
1099 }
1100 }
1101 let mut n = match read_node_opt(paths, &node_id)? {
1102 Some(n) => n,
1103 None => return Ok(vec![]),
1110 };
1111 if attempt != n.retry_attempts || n.worker_exit.is_some() {
1115 return Ok(vec![]);
1116 }
1117 n.worker_exit = Some(WorkerExit {
1118 code,
1119 signal,
1120 at: ev.ts,
1121 });
1122 n.awaiting_input = None;
1126 n.updated_at = ev.ts;
1127 Ok(vec![ProjectionOp::Node(n)])
1128}
1129
1130fn evidence_attempt(events_path: &Path, ev: &Event) -> Result<u32> {
1131 ev.data
1132 .get("attempt")
1133 .and_then(Value::as_u64)
1134 .and_then(|raw| u32::try_from(raw).ok())
1135 .ok_or_else(|| Error::CorruptEventLog {
1136 path: events_path.to_path_buf(),
1137 reason: format!("event seq={} attempt must be a u32", ev.seq),
1138 })
1139}
1140
1141fn reduce_worker_evidence_archived(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1142 let events_path = paths.events();
1143 let node_id = require_envelope_node_id(&events_path, ev)?;
1144 let attempt = evidence_attempt(&events_path, ev)?;
1145 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1146 let mut values = Vec::new();
1147 for field in [
1148 "transcript_path",
1149 "resume_path",
1150 "pane_path",
1151 "report_path",
1152 "transcript_sha256",
1153 ] {
1154 values.push(want_str(&events_path, ev, &ev.data, field)?.to_string());
1155 }
1156 let expected_prefix = format!("evidence/{}/", node_id.as_str());
1157 for (index, suffix) in [
1158 "pi-session.original.jsonl",
1159 "pi-session.resume.jsonl",
1160 "final-pane.log",
1161 "terminal-report.json",
1162 ]
1163 .iter()
1164 .enumerate()
1165 {
1166 if values[index] != format!("{expected_prefix}{suffix}") {
1167 return Err(Error::CorruptEventLog {
1168 path: events_path,
1169 reason: format!(
1170 "event seq={} evidence artifact path is not canonical",
1171 ev.seq
1172 ),
1173 });
1174 }
1175 }
1176 if values[4].len() != 64
1177 || !values[4]
1178 .bytes()
1179 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1180 {
1181 return Err(Error::CorruptEventLog {
1182 path: events_path,
1183 reason: format!(
1184 "event seq={} transcript_sha256 is not lowercase SHA-256",
1185 ev.seq
1186 ),
1187 });
1188 }
1189 let mut node = match read_node_opt(paths, &node_id)? {
1190 Some(node) => node,
1191 None => return Ok(vec![]),
1192 };
1193 let Some(evidence) = node.evidence.as_mut() else {
1194 return Err(Error::CorruptEventLog {
1195 path: events_path,
1196 reason: format!(
1197 "event seq={} evidence archive has no recorded Pi session",
1198 ev.seq
1199 ),
1200 });
1201 };
1202 if evidence.attempt != attempt || evidence.session_id != session_id {
1203 return Ok(vec![]);
1204 }
1205 if evidence.status == EvidenceStatus::Complete {
1206 return Ok(vec![]);
1207 }
1208 evidence.transcript_path = Some(values.remove(0));
1209 evidence.resume_path = Some(values.remove(0));
1210 evidence.pane_path = Some(values.remove(0));
1211 evidence.report_path = Some(values.remove(0));
1212 evidence.transcript_sha256 = Some(values.remove(0));
1213 evidence.status = EvidenceStatus::Complete;
1214 evidence.error = None;
1215 node.updated_at = ev.ts;
1216 Ok(vec![ProjectionOp::Node(node)])
1217}
1218
1219fn reduce_worker_display_retained(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1220 let events_path = paths.events();
1221 let node_id = require_envelope_node_id(&events_path, ev)?;
1222 let display: crate::schema::RetainedDisplay =
1223 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1224 path: events_path.clone(),
1225 reason: format!("event seq={} invalid retained display: {e}", ev.seq),
1226 })?;
1227 let mut node = match read_node_opt(paths, &node_id)? {
1228 Some(node) => node,
1229 None => return Ok(vec![]),
1230 };
1231 if display.attempt != node.retry_attempts || display.ownership_marker.is_empty() {
1232 return Ok(vec![]);
1233 }
1234 if node.retained_display.is_some() {
1235 return Ok(vec![]);
1236 }
1237 node.retained_display = Some(Box::new(display));
1238 node.updated_at = ev.ts;
1239 Ok(vec![ProjectionOp::Node(node)])
1240}
1241
1242fn reduce_worker_display_expired(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1243 let events_path = paths.events();
1244 let node_id = require_envelope_node_id(&events_path, ev)?;
1245 let marker = want_str(&events_path, ev, &ev.data, "ownership_marker")?;
1246 let attempt = evidence_attempt(&events_path, ev)?;
1247 let mut node = match read_node_opt(paths, &node_id)? {
1248 Some(node) => node,
1249 None => return Ok(vec![]),
1250 };
1251 let Some(display) = node.retained_display.as_mut() else {
1252 return Ok(vec![]);
1253 };
1254 if display.attempt != attempt
1255 || display.ownership_marker != marker
1256 || display.expired_at.is_some()
1257 {
1258 return Ok(vec![]);
1259 }
1260 display.expired_at = Some(ev.ts);
1261 node.updated_at = ev.ts;
1262 Ok(vec![ProjectionOp::Node(node)])
1263}
1264
1265fn reduce_worker_display_unavailable(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1266 let events_path = paths.events();
1267 let node_id = require_envelope_node_id(&events_path, ev)?;
1268 let attempt = evidence_attempt(&events_path, ev)?;
1269 let reason = want_str(&events_path, ev, &ev.data, "reason")?;
1270 let mut node = match read_node_opt(paths, &node_id)? {
1271 Some(node) => node,
1272 None => return Ok(vec![]),
1273 };
1274 if attempt != node.retry_attempts || node.retained_display.is_some() {
1275 return Ok(vec![]);
1276 }
1277 if node.retention_unavailable.is_none() {
1278 node.retention_unavailable = Some(reason.to_string());
1279 node.updated_at = ev.ts;
1280 }
1281 Ok(vec![ProjectionOp::Node(node)])
1282}
1283
1284fn reduce_worker_evidence_failed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1285 let events_path = paths.events();
1286 let node_id = require_envelope_node_id(&events_path, ev)?;
1287 let attempt = evidence_attempt(&events_path, ev)?;
1288 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1289 let detail = want_str(&events_path, ev, &ev.data, "error")?.to_string();
1290 let mut node = match read_node_opt(paths, &node_id)? {
1291 Some(node) => node,
1292 None => return Ok(vec![]),
1293 };
1294 let Some(evidence) = node.evidence.as_mut() else {
1295 return Err(Error::CorruptEventLog {
1296 path: events_path,
1297 reason: format!(
1298 "event seq={} evidence failure has no recorded Pi session",
1299 ev.seq
1300 ),
1301 });
1302 };
1303 if evidence.attempt != attempt || evidence.session_id != session_id {
1304 return Ok(vec![]);
1305 }
1306 if evidence.status == EvidenceStatus::Complete {
1307 return Ok(vec![]);
1308 }
1309 evidence.status = EvidenceStatus::Failed;
1310 evidence.error = Some(detail);
1311 node.updated_at = ev.ts;
1312 Ok(vec![ProjectionOp::Node(node)])
1313}
1314
1315fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1328 let events_path = paths.events();
1329 let node_id = require_envelope_node_id(&events_path, ev)?;
1330 let mut n = match read_node_opt(paths, &node_id)? {
1331 Some(n) => n,
1332 None => return Ok(vec![]),
1333 };
1334 if n.first_death_at.is_some()
1340 || n.status.is_terminal()
1341 || n.worker_exit.is_some()
1342 || n.last_report.is_some()
1343 || n.pending_merge.is_some()
1344 {
1345 return Ok(vec![]);
1346 }
1347 n.first_death_at = Some(ev.ts);
1348 n.updated_at = ev.ts;
1349 Ok(vec![ProjectionOp::Node(n)])
1350}
1351
1352fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1361 const MAX_ITEMS: usize = 8;
1362 const MAX_TOPIC_CHARS: usize = 512;
1363 const MAX_OPTIONS: usize = 16;
1364 const MAX_OPTION_CHARS: usize = 256;
1365
1366 let events_path = paths.events();
1367 let node_id = require_envelope_node_id(&events_path, ev)?;
1368 let items = ev
1371 .data
1372 .get("discussion_items")
1373 .and_then(Value::as_array)
1374 .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1375 .ok_or_else(|| Error::CorruptEventLog {
1376 path: events_path.clone(),
1377 reason: format!(
1378 "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1379 ev.seq
1380 ),
1381 })?;
1382 for (index, item) in items.iter().enumerate() {
1383 let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1384 path: events_path.clone(),
1385 reason: format!(
1386 "event seq={} discussion_items[{index}] must be an object",
1387 ev.seq
1388 ),
1389 })?;
1390 let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1391 let default = obj
1392 .get("recommended_default")
1393 .and_then(Value::as_str)
1394 .unwrap_or("");
1395 let options = obj.get("options").and_then(Value::as_array);
1396 let options_valid = options.is_some_and(|values| {
1397 !values.is_empty()
1398 && values.len() <= MAX_OPTIONS
1399 && values.iter().all(|v| {
1400 v.as_str().is_some_and(|s| {
1401 !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1402 })
1403 })
1404 && values.iter().any(|v| v.as_str() == Some(default))
1405 });
1406 if topic.trim().is_empty()
1407 || topic.chars().count() > MAX_TOPIC_CHARS
1408 || default.trim().is_empty()
1409 || !options_valid
1410 {
1411 return Err(Error::CorruptEventLog {
1412 path: events_path.clone(),
1413 reason: format!(
1414 "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1415 ev.seq
1416 ),
1417 });
1418 }
1419 }
1420
1421 let mut n = match read_node_opt(paths, &node_id)? {
1422 Some(n) => n,
1423 None => return Ok(vec![]),
1424 };
1425 if n.status.is_terminal() || n.worker_exit.is_some() {
1426 return Ok(vec![]);
1427 }
1428 if let Some(open) = n.awaiting_input.as_mut() {
1429 if open.discussion_items.len() + items.len() > MAX_ITEMS {
1432 return Err(Error::CorruptEventLog {
1433 path: events_path,
1434 reason: format!(
1435 "event seq={} would exceed {MAX_ITEMS} open discussion items",
1436 ev.seq
1437 ),
1438 });
1439 }
1440 open.discussion_items.extend(items.iter().cloned());
1441 } else {
1442 n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1443 opened_at: ev.ts,
1444 event_seq: ev.seq,
1445 discussion_items: items.clone(),
1446 }));
1447 }
1448 n.updated_at = ev.ts;
1449 Ok(vec![ProjectionOp::Node(n)])
1450}
1451
1452fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1456 let events_path = paths.events();
1457 let node_id = require_envelope_node_id(&events_path, ev)?;
1458 let seq = ev
1461 .data
1462 .get("event_seq")
1463 .and_then(Value::as_u64)
1464 .ok_or_else(|| Error::CorruptEventLog {
1465 path: events_path.clone(),
1466 reason: format!(
1467 "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1468 ev.seq
1469 ),
1470 })?;
1471 let mut n = match read_node_opt(paths, &node_id)? {
1472 Some(n) => n,
1473 None => return Ok(vec![]),
1474 };
1475 let Some(open) = n.awaiting_input.as_ref() else {
1476 return Ok(vec![]);
1477 };
1478 if seq != open.event_seq {
1479 return Ok(vec![]);
1480 }
1481 n.awaiting_input = None;
1482 n.updated_at = ev.ts;
1483 Ok(vec![ProjectionOp::Node(n)])
1484}
1485
1486pub const KIND_MERGE_STARTED: &str = "merge.started";
1490
1491pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1496
1497fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1511 let events_path = paths.events();
1512 let node_id = require_envelope_node_id(&events_path, ev)?;
1513 let txn: MergeTxn =
1514 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1515 path: events_path.clone(),
1516 reason: format!(
1517 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1518 ev.seq
1519 ),
1520 })?;
1521 let mut n = match read_node_opt(paths, &node_id)? {
1522 Some(n) => n,
1523 None => return Ok(vec![]),
1524 };
1525 if n.status.is_terminal() {
1529 return Ok(vec![]);
1530 }
1531 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1534 return Ok(vec![]);
1535 }
1536 n.pending_merge = Some(Box::new(txn));
1537 n.updated_at = ev.ts;
1538 Ok(vec![ProjectionOp::Node(n)])
1539}
1540
1541fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1553 let events_path = paths.events();
1554 let node_id = require_envelope_node_id(&events_path, ev)?;
1555 let op_id = ev
1556 .data
1557 .get("op_id")
1558 .and_then(Value::as_str)
1559 .ok_or_else(|| Error::CorruptEventLog {
1560 path: events_path.clone(),
1561 reason: format!(
1562 "event seq={} kind=merge.aborted is missing string `op_id`",
1563 ev.seq
1564 ),
1565 })?;
1566 let mut n = match read_node_opt(paths, &node_id)? {
1567 Some(n) => n,
1568 None => return Ok(vec![]),
1569 };
1570 match n.pending_merge.as_ref() {
1573 Some(t) if t.op_id == op_id => {}
1574 _ => return Ok(vec![]),
1575 }
1576 n.pending_merge = None;
1577 n.updated_at = ev.ts;
1578 Ok(vec![ProjectionOp::Node(n)])
1579}
1580
1581fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1589 if current == incoming {
1590 tracing::debug!(
1591 target: "taskfleet_core::reducer",
1592 seq = ev.seq, kind = %ev.kind, status = ?current,
1593 "no-op: status re-applied to terminal target"
1594 );
1595 } else {
1596 tracing::warn!(
1597 target: "taskfleet_core::reducer",
1598 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1599 "no-op: ignored conflicting transition from terminal target"
1600 );
1601 }
1602}
1603
1604fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1613 let corrupt = |reason: String| Error::CorruptEventLog {
1614 path: events_path.to_path_buf(),
1615 reason,
1616 };
1617 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1618 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1619 if cancelled {
1620 if success == Some(true) {
1621 return Err(corrupt(format!(
1622 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1623 ev.seq
1624 )));
1625 }
1626 Ok(Status::Cancelled)
1627 } else {
1628 match success {
1629 Some(true) => Ok(Status::Done),
1630 Some(false) => Ok(Status::Failed),
1631 None => Err(corrupt(format!(
1632 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1633 ev.seq
1634 ))),
1635 }
1636 }
1637}
1638
1639fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1640 let events_path = paths.events();
1643 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1644 path: events_path.clone(),
1645 reason: format!(
1646 "event seq={} kind=child.spawned missing parent `node_id`",
1647 ev.seq
1648 ),
1649 })?;
1650 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1651 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1652 let child_node_id = NodeId::parse_str(
1653 ev.data
1654 .get("child_node_id")
1655 .and_then(Value::as_str)
1656 .unwrap_or("n-0001"),
1657 )
1658 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1659 let mut n = match read_node_opt(paths, &parent_node_id)? {
1660 Some(n) => n,
1661 None => return Ok(vec![]),
1662 };
1663 let new_ref = ChildRef {
1664 run_id: child_run_id,
1665 node_id: child_node_id,
1666 };
1667 if n.children.iter().any(|c| c == &new_ref) {
1668 return Ok(vec![]);
1671 }
1672 n.children.push(new_ref);
1673 n.updated_at = ev.ts;
1674 Ok(vec![ProjectionOp::Node(n)])
1675}
1676
1677fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1688 let events_path = paths.events();
1689 let node_id = require_envelope_node_id(&events_path, ev)?;
1690 let raw = ev
1691 .data
1692 .get("pid")
1693 .and_then(Value::as_i64)
1694 .ok_or_else(|| Error::CorruptEventLog {
1695 path: events_path.clone(),
1696 reason: format!(
1697 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1698 ev.seq
1699 ),
1700 })?;
1701 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1702 path: events_path.clone(),
1703 reason: format!(
1704 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1705 ev.seq
1706 ),
1707 })?;
1708 let mut n = match read_node_opt(paths, &node_id)? {
1709 Some(n) => n,
1710 None => return Ok(vec![]),
1711 };
1712 if n.supervisor_pid == Some(pid) {
1713 return Ok(vec![]);
1714 }
1715 n.supervisor_pid = Some(pid);
1716 n.updated_at = ev.ts;
1717 Ok(vec![ProjectionOp::Node(n)])
1718}
1719
1720fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1731 let events_path = paths.events();
1732 let node_id = require_envelope_node_id(&events_path, ev)?;
1733 let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1734 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1738 let report_seq = ev
1739 .data
1740 .get("report_seq")
1741 .and_then(Value::as_u64)
1742 .ok_or_else(|| Error::CorruptEventLog {
1743 path: events_path.clone(),
1744 reason: format!(
1745 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1746 ev.seq
1747 ),
1748 })?;
1749 let mut n = match read_node_opt(paths, &node_id)? {
1750 Some(n) => n,
1751 None => return Ok(vec![]),
1752 };
1753 if let Some(prev) = n
1754 .last_processed_report_seq_by_child
1755 .get(child_run_id)
1756 .and_then(Value::as_u64)
1757 {
1758 if report_seq <= prev {
1759 return Ok(vec![]);
1760 }
1761 }
1762 n.last_processed_report_seq_by_child
1763 .insert(child_run_id.to_string(), Value::from(report_seq));
1764 n.updated_at = ev.ts;
1765 Ok(vec![ProjectionOp::Node(n)])
1766}
1767
1768#[cfg(test)]
1769mod tests {
1770 use super::*;
1771 use crate::schema::Event;
1772 use chrono::Utc;
1773 use tempfile::TempDir;
1774
1775 fn event(run_id: &str) -> Event {
1776 Event {
1777 ts: Utc::now(),
1778 seq: 1,
1779 kind: "run.status".into(),
1780 run_id: RunId::parse_str(run_id).unwrap(),
1781 node_id: None,
1782 idempotency_key: None,
1783 data: serde_json::json!({ "status": "running" }),
1784 }
1785 }
1786
1787 #[test]
1788 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1789 let tmp = TempDir::new().unwrap();
1793 let run_id = "01jxsnap000000000000000000";
1794 let rid = RunId::parse_str(run_id).unwrap();
1795 let dir = crate::run_dir(tmp.path(), &rid);
1796 std::fs::create_dir_all(&dir).unwrap();
1797 let paths = RunPaths::new(dir, run_id).unwrap();
1798
1799 let mut created = event(run_id);
1802 created.kind = "run.created".into();
1803 created.data = serde_json::json!({
1804 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1805 });
1806 apply_event(&paths, &created).expect("run.created applies");
1807 let manifest_before = std::fs::read(paths.manifest()).unwrap();
1808
1809 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1810 let mut ev = event(run_id);
1811 ev.seq = seq;
1812 ev.kind = kind.into();
1813 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1815 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1816 assert!(ops.is_empty(), "{kind} must plan no projection ops");
1817 apply_event(&paths, &ev).expect("audit kind applies as no-op");
1819 }
1820
1821 assert_eq!(
1823 std::fs::read(paths.manifest()).unwrap(),
1824 manifest_before,
1825 "audit events must not mutate the manifest"
1826 );
1827 assert!(!paths.nodes_dir().exists(), "no node projection created");
1828 }
1829
1830 #[test]
1831 fn run_created_folds_harness_when_present_and_defaults_none() {
1832 let tmp = TempDir::new().unwrap();
1833
1834 let run_id = "01jxhrnsaa0000000000000001";
1836 let rid = RunId::parse_str(run_id).unwrap();
1837 let dir = crate::run_dir(tmp.path(), &rid);
1838 std::fs::create_dir_all(&dir).unwrap();
1839 let paths = RunPaths::new(dir, run_id).unwrap();
1840 let mut created = event(run_id);
1841 created.kind = "run.created".into();
1842 created.data = serde_json::json!({
1843 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1844 "harness": "pi", "harness_source": "flag",
1845 });
1846 apply_event(&paths, &created).expect("run.created applies");
1847 let m = read_manifest_opt(&paths).unwrap().unwrap();
1848 assert_eq!(m.harness.as_deref(), Some("pi"));
1849
1850 let run_id2 = "01jxhrnsaa0000000000000002";
1852 let rid2 = RunId::parse_str(run_id2).unwrap();
1853 let dir2 = crate::run_dir(tmp.path(), &rid2);
1854 std::fs::create_dir_all(&dir2).unwrap();
1855 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1856 let mut created2 = event(run_id2);
1857 created2.kind = "run.created".into();
1858 created2.data =
1859 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1860 apply_event(&paths2, &created2).expect("run.created applies");
1861 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1862 assert_eq!(m2.harness, None);
1863 }
1864
1865 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1868 let rid = RunId::parse_str(run_id).unwrap();
1869 let dir = crate::run_dir(tmp.path(), &rid);
1870 std::fs::create_dir_all(&dir).unwrap();
1871 let paths = RunPaths::new(dir, run_id).unwrap();
1872 let mut created = event(run_id);
1873 created.kind = "run.created".into();
1874 created.data =
1875 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1876 apply_event(&paths, &created).expect("run.created applies");
1877 let mut node = event(run_id);
1878 node.seq = 2;
1879 node.kind = "node.created".into();
1880 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1881 node.data = serde_json::json!({
1882 "kind": "spinoff",
1883 "branch": "wt/foo",
1884 "worktree_path": "/tmp/old-wt",
1885 "agent_pid": 111,
1886 });
1887 apply_event(&paths, &node).expect("node.created applies");
1888 paths
1889 }
1890
1891 #[test]
1895 fn node_retry_rewires_node_and_increments_attempts() {
1896 let tmp = TempDir::new().unwrap();
1897 let run_id = "01jxsnap000000000000000000";
1898 let paths = bootstrap_retry_node(&tmp, run_id);
1899
1900 let mut retry = event(run_id);
1901 retry.seq = 3;
1902 retry.kind = "node.retry".into();
1903 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1904 retry.data = serde_json::json!({
1905 "attempt": 1,
1906 "reason": "agent-died",
1907 "branch": "wt/foo-r1",
1908 "base_sha": "a".repeat(40),
1909 "worktree_path": "/tmp/new-wt",
1910 "agent_pid": 222,
1911 "tmux_session": "s",
1912 "tmux_window_id": "@9",
1913 });
1914 apply_event(&paths, &retry).expect("node.retry applies");
1915
1916 let n = read_n0001(&paths);
1917 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1918 assert_eq!(
1919 n.branch.as_deref(),
1920 Some("wt/foo-r1"),
1921 "rewired to new branch"
1922 );
1923 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1924 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1925 assert_eq!(n.status, Status::Pending, "node returns to pending");
1926 assert!(n.last_report.is_none());
1927 assert_eq!(
1928 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1929 Some("@9"),
1930 "rewired tmux identity"
1931 );
1932
1933 let mut retry2 = event(run_id);
1935 retry2.seq = 4;
1936 retry2.kind = "node.retry".into();
1937 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1938 retry2.data = serde_json::json!({
1939 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1940 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1941 });
1942 apply_event(&paths, &retry2).expect("node.retry applies");
1943 assert_eq!(read_n0001(&paths).retry_attempts, 2);
1944 }
1945
1946 #[test]
1947 fn stale_worker_evidence_and_exit_do_not_cross_retry_generation() {
1948 let tmp = TempDir::new().unwrap();
1949 let run_id = "01jxsnap000000000000000000";
1950 let paths = bootstrap_retry_node(&tmp, run_id);
1951 let nid = NodeId::parse_str("n-0001").unwrap();
1952 let current_session = "018f5f64-b137-7d44-b2b4-4f02c3f646e8";
1953
1954 let mut retry = event(run_id);
1955 retry.seq = 3;
1956 retry.kind = "node.retry".into();
1957 retry.node_id = Some(nid.clone());
1958 retry.data = serde_json::json!({
1959 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1960 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1961 "pi_session_id": current_session,
1962 "pi_session_path": format!(".creating/pi-sessions/{run_id}/pi-session-{current_session}.jsonl"),
1963 "pi_session_cwd": "/tmp/new-wt"
1964 });
1965 apply_event(&paths, &retry).unwrap();
1966
1967 let mut stale_failure = event(run_id);
1968 stale_failure.seq = 4;
1969 stale_failure.kind = "worker.evidence.failed".into();
1970 stale_failure.node_id = Some(nid.clone());
1971 stale_failure.data = serde_json::json!({
1972 "attempt": 0,
1973 "session_id": "118f5f64-b137-7d44-b2b4-4f02c3f646e8",
1974 "error": "old attempt"
1975 });
1976 apply_event(&paths, &stale_failure).unwrap();
1977 assert_eq!(
1978 read_n0001(&paths).evidence.unwrap().status,
1979 EvidenceStatus::Pending
1980 );
1981
1982 let mut stale_exit = event(run_id);
1983 stale_exit.seq = 5;
1984 stale_exit.kind = "worker.exited".into();
1985 stale_exit.node_id = Some(nid.clone());
1986 stale_exit.data = serde_json::json!({"attempt":0,"exit_code":9});
1987 apply_event(&paths, &stale_exit).unwrap();
1988 assert!(read_n0001(&paths).worker_exit.is_none());
1989
1990 let mut archived = event(run_id);
1991 archived.seq = 6;
1992 archived.kind = "worker.evidence.archived".into();
1993 archived.node_id = Some(nid);
1994 archived.data = serde_json::json!({
1995 "attempt":1, "session_id":current_session,
1996 "transcript_path":"evidence/n-0001/pi-session.original.jsonl",
1997 "resume_path":"evidence/n-0001/pi-session.resume.jsonl",
1998 "pane_path":"evidence/n-0001/final-pane.log",
1999 "report_path":"evidence/n-0001/terminal-report.json",
2000 "transcript_sha256":"0".repeat(64)
2001 });
2002 apply_event(&paths, &archived).unwrap();
2003 assert_eq!(
2004 read_n0001(&paths).evidence.unwrap().status,
2005 EvidenceStatus::Complete
2006 );
2007 }
2008
2009 #[test]
2013 fn node_retry_against_terminal_node_is_noop() {
2014 let tmp = TempDir::new().unwrap();
2015 let run_id = "01jxsnap000000000000000000";
2016 let paths = bootstrap_retry_node(&tmp, run_id);
2017
2018 let mut report = event(run_id);
2020 report.seq = 3;
2021 report.kind = "node.report".into();
2022 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2023 report.data = serde_json::json!({ "success": true });
2024 apply_event(&paths, &report).expect("node.report applies");
2025 assert_eq!(read_n0001(&paths).status, Status::Done);
2026
2027 let mut retry = event(run_id);
2028 retry.seq = 4;
2029 retry.kind = "node.retry".into();
2030 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2031 retry.data = serde_json::json!({
2032 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
2033 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
2034 });
2035 apply_event(&paths, &retry).expect("node.retry applies as no-op");
2036
2037 let n = read_n0001(&paths);
2038 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
2039 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
2040 assert_eq!(n.agent_pid, Some(111), "not rewired");
2041 }
2042
2043 #[test]
2044 fn apply_event_rejects_event_from_a_different_run() {
2045 let tmp = TempDir::new().unwrap();
2046 let run_id = "01jxsnap000000000000000000";
2047 let rid = RunId::parse_str(run_id).unwrap();
2048 let dir = crate::run_dir(tmp.path(), &rid);
2049 std::fs::create_dir_all(&dir).unwrap();
2050 let paths = RunPaths::new(dir, run_id).unwrap();
2051
2052 let foreign = event("02jxsnap000000000000000000");
2054 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
2055 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
2056
2057 let mine = event(run_id);
2060 apply_event(&paths, &mine).expect("matching run_id must be accepted");
2061 }
2062
2063 #[test]
2064 fn tmux_identity_from_data_reads_qualified_fields() {
2065 let d = serde_json::json!({
2066 "tmux_socket": "/private/tmp/tmux-501/default",
2067 "tmux_session": "taskfleet",
2068 "tmux_window_id": "@42",
2069 });
2070 let id = tmux_identity_from_data(&d).expect("qualified identity");
2071 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
2072 assert_eq!(id.session, "taskfleet");
2073 assert_eq!(id.window_id, "@42");
2074 assert_eq!(id.pane_id, None);
2076
2077 let d2 = serde_json::json!({
2079 "tmux_socket": null,
2080 "tmux_session": "taskfleet",
2081 "tmux_window_id": "@7",
2082 });
2083 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
2084 assert_eq!(id2.socket, None);
2085 assert_eq!(id2.window_id, "@7");
2086
2087 let d3 = serde_json::json!({
2089 "tmux_session": "taskfleet",
2090 "tmux_window_id": "@42",
2091 "tmux_pane_id": "%7",
2092 });
2093 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
2094 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
2095 assert_eq!(id3.capture_target(), "%7");
2096
2097 let d4 = serde_json::json!({
2100 "tmux_session": "taskfleet",
2101 "tmux_window_id": "@42",
2102 "tmux_pane_id": null,
2103 });
2104 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
2105 assert_eq!(id4.pane_id, None);
2106 assert_eq!(id4.capture_target(), "@42");
2107 }
2108
2109 #[test]
2110 fn tmux_identity_from_data_back_compat_is_none() {
2111 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
2113 assert!(tmux_identity_from_data(&legacy).is_none());
2114 let partial = serde_json::json!({ "tmux_window_id": "@42" });
2116 assert!(tmux_identity_from_data(&partial).is_none());
2117 }
2118
2119 #[test]
2122 fn node_created_populates_tmux_identity() {
2123 let tmp = TempDir::new().unwrap();
2124 let run_id = "01jxsnap000000000000000000";
2125 let rid = RunId::parse_str(run_id).unwrap();
2126 let dir = crate::run_dir(tmp.path(), &rid);
2127 std::fs::create_dir_all(&dir).unwrap();
2128 let paths = RunPaths::new(dir, run_id).unwrap();
2129
2130 let mut ev = event(run_id);
2131 ev.seq = 2;
2132 ev.kind = "node.created".into();
2133 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2134 ev.data = serde_json::json!({
2135 "kind": "spinoff",
2136 "tmux_window": "🚀 wt/x",
2137 "tmux_socket": "/private/tmp/tmux-501/default",
2138 "tmux_session": "taskfleet",
2139 "tmux_window_id": "@42",
2140 });
2141 apply_event(&paths, &ev).expect("node.created applies");
2142 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2143 .unwrap()
2144 .unwrap();
2145 let id = n.tmux_identity.expect("qualified identity recorded");
2146 assert_eq!(id.session, "taskfleet");
2147 assert_eq!(id.window_id, "@42");
2148 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
2149
2150 let run2 = "02jxsnap000000000000000000";
2152 let rid2 = RunId::parse_str(run2).unwrap();
2153 let dir2 = crate::run_dir(tmp.path(), &rid2);
2154 std::fs::create_dir_all(&dir2).unwrap();
2155 let paths2 = RunPaths::new(dir2, run2).unwrap();
2156 let mut ev2 = event(run2);
2157 ev2.seq = 2;
2158 ev2.kind = "node.created".into();
2159 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2160 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
2161 apply_event(&paths2, &ev2).expect("legacy node.created applies");
2162 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
2163 .unwrap()
2164 .unwrap();
2165 assert!(n2.tmux_identity.is_none());
2166 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
2167 }
2168
2169 #[test]
2170 fn node_materialization_populates_missing_manifest_source_branch() {
2171 let tmp = TempDir::new().unwrap();
2172 let run_id = "01jxsnap000000000000000001";
2173 let rid = RunId::parse_str(run_id).unwrap();
2174 let dir = crate::run_dir(tmp.path(), &rid);
2175 std::fs::create_dir_all(&dir).unwrap();
2176 let paths = RunPaths::new(dir, run_id).unwrap();
2177
2178 let mut created = event(run_id);
2179 created.kind = "run.created".into();
2180 created.data = serde_json::json!({
2181 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2182 });
2183 apply_event(&paths, &created).unwrap();
2184 assert!(read_manifest_opt(&paths)
2185 .unwrap()
2186 .unwrap()
2187 .source_branch
2188 .is_none());
2189
2190 let mut node = event(run_id);
2191 node.seq = 2;
2192 node.kind = "node.created".into();
2193 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2194 node.data = serde_json::json!({
2195 "kind": "spinoff",
2196 "source_branch": "main",
2197 "worktree_path": "/tmp/wt/pending"
2198 });
2199 apply_event(&paths, &node).unwrap();
2200
2201 let manifest = read_manifest_opt(&paths).unwrap().unwrap();
2202 assert_eq!(manifest.status, Status::Pending);
2203 assert_eq!(manifest.source_branch.as_deref(), Some("main"));
2204 let node = read_n0001(&paths);
2205 assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
2206 }
2207
2208 #[test]
2209 fn node_materialization_preserves_explicit_manifest_source_branch() {
2210 let tmp = TempDir::new().unwrap();
2211 let run_id = "01jxsnap000000000000000002";
2212 let rid = RunId::parse_str(run_id).unwrap();
2213 let dir = crate::run_dir(tmp.path(), &rid);
2214 std::fs::create_dir_all(&dir).unwrap();
2215 let paths = RunPaths::new(dir, run_id).unwrap();
2216
2217 let mut created = event(run_id);
2218 created.kind = "run.created".into();
2219 created.data = serde_json::json!({
2220 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
2221 "source_branch": "release"
2222 });
2223 apply_event(&paths, &created).unwrap();
2224
2225 let mut node = event(run_id);
2226 node.seq = 2;
2227 node.kind = "node.created".into();
2228 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2229 node.data = serde_json::json!({
2230 "kind": "spinoff", "source_branch": "main"
2231 });
2232 apply_event(&paths, &node).unwrap();
2233
2234 assert_eq!(
2235 read_manifest_opt(&paths)
2236 .unwrap()
2237 .unwrap()
2238 .source_branch
2239 .as_deref(),
2240 Some("release")
2241 );
2242 }
2243
2244 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
2247 let rid = RunId::parse_str(run_id).unwrap();
2248 let dir = crate::run_dir(tmp.path(), &rid);
2249 std::fs::create_dir_all(&dir).unwrap();
2250 let paths = RunPaths::new(dir, run_id).unwrap();
2251
2252 let mut created = event(run_id);
2253 created.kind = "run.created".into();
2254 created.data = serde_json::json!({
2255 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2256 });
2257 apply_event(&paths, &created).expect("run.created applies");
2258
2259 let mut node = event(run_id);
2260 node.seq = 2;
2261 node.kind = "node.created".into();
2262 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2263 node.data = serde_json::json!({ "kind": "spinoff" });
2264 apply_event(&paths, &node).expect("node.created applies");
2265 paths
2266 }
2267
2268 fn read_n0001(paths: &RunPaths) -> Node {
2269 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
2270 .unwrap()
2271 .unwrap()
2272 }
2273
2274 #[test]
2275 fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
2276 let tmp = TempDir::new().unwrap();
2277 let run_id = "01jxwd0000000000000000000w";
2278 let paths = seed_run_with_node(&tmp, run_id);
2279 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2280 let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
2281 let mut open = event(run_id);
2282 open.seq = 3;
2283 open.ts = opened_at;
2284 open.kind = "node.awaiting_input".into();
2285 open.node_id = nid.clone();
2286 open.data = serde_json::json!({ "discussion_items": [{
2287 "topic": "Which scope?",
2288 "options": ["small", "large"],
2289 "recommended_default": "small"
2290 }] });
2291 apply_event(&paths, &open).unwrap();
2292 let first = read_n0001(&paths).awaiting_input.unwrap();
2293 assert_eq!(first.opened_at, opened_at);
2294 assert_eq!(first.event_seq, 3);
2295
2296 let mut duplicate = open.clone();
2298 duplicate.seq = 4;
2299 duplicate.ts = opened_at + chrono::Duration::hours(1);
2300 apply_event(&paths, &duplicate).unwrap();
2301 let still_first = read_n0001(&paths).awaiting_input.unwrap();
2302 assert_eq!(still_first.opened_at, opened_at);
2303 assert_eq!(still_first.event_seq, 3);
2304
2305 let mut stale = event(run_id);
2307 stale.seq = 5;
2308 stale.kind = "node.input_resolved".into();
2309 stale.node_id = nid.clone();
2310 stale.data = serde_json::json!({ "event_seq": 2 });
2311 apply_event(&paths, &stale).unwrap();
2312 assert!(read_n0001(&paths).awaiting_input.is_some());
2313
2314 let mut resolved = stale;
2315 resolved.seq = 6;
2316 resolved.data = serde_json::json!({ "event_seq": 3 });
2317 apply_event(&paths, &resolved).unwrap();
2318 assert!(read_n0001(&paths).awaiting_input.is_none());
2319 }
2320
2321 #[test]
2322 fn awaiting_input_rejects_missing_default_without_mutating_projection() {
2323 let tmp = TempDir::new().unwrap();
2324 let run_id = "01jxwd0000000000000000000x";
2325 let paths = seed_run_with_node(&tmp, run_id);
2326 let mut open = event(run_id);
2327 open.seq = 3;
2328 open.kind = "node.awaiting_input".into();
2329 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2330 open.data = serde_json::json!({ "discussion_items": [{
2331 "topic": "Which scope?", "options": ["small", "large"]
2332 }] });
2333 assert!(reduce_event_to_ops(&paths, &open).is_err());
2334 assert!(read_n0001(&paths).awaiting_input.is_none());
2335 }
2336
2337 #[test]
2338 fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
2339 let tmp = TempDir::new().unwrap();
2340 let run_id = "01jxwd0000000000000000000y";
2341 let paths = seed_run_with_node(&tmp, run_id);
2342 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2343
2344 let mut open = event(run_id);
2345 open.seq = 3;
2346 open.kind = "node.awaiting_input".into();
2347 open.node_id = nid.clone();
2348 open.data = serde_json::json!({ "discussion_items": [{
2349 "topic": "Which scope?", "options": ["small", "large"],
2350 "recommended_default": "small"
2351 }] });
2352 apply_event(&paths, &open).unwrap();
2353
2354 let mut malformed_duplicate = open.clone();
2355 malformed_duplicate.seq = 4;
2356 malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
2357 assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
2358
2359 let mut exited = event(run_id);
2360 exited.seq = 5;
2361 exited.kind = "worker.exited".into();
2362 exited.node_id = nid;
2363 exited.data = serde_json::json!({ "exit_code": 0 });
2364 apply_event(&paths, &exited).unwrap();
2365 let node = read_n0001(&paths);
2366 assert!(node.awaiting_input.is_none());
2367 assert!(node.worker_exit.is_some());
2368
2369 let mut delayed_open = open;
2370 delayed_open.seq = 6;
2371 assert!(reduce_event_to_ops(&paths, &delayed_open)
2372 .unwrap()
2373 .is_empty());
2374 }
2375
2376 #[test]
2377 fn input_resolved_requires_generation_even_when_nothing_is_open() {
2378 let tmp = TempDir::new().unwrap();
2379 let run_id = "01jxwd0000000000000000000z";
2380 let paths = seed_run_with_node(&tmp, run_id);
2381 let mut resolved = event(run_id);
2382 resolved.seq = 3;
2383 resolved.kind = "node.input_resolved".into();
2384 resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2385 resolved.data = serde_json::json!({});
2386 assert!(reduce_event_to_ops(&paths, &resolved).is_err());
2387 }
2388
2389 #[test]
2390 fn awaiting_input_default_must_be_one_of_options() {
2391 let tmp = TempDir::new().unwrap();
2392 let run_id = "01jxwd00000000000000000010";
2393 let paths = seed_run_with_node(&tmp, run_id);
2394 let mut open = event(run_id);
2395 open.seq = 3;
2396 open.kind = "node.awaiting_input".into();
2397 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2398 open.data = serde_json::json!({ "discussion_items": [{
2399 "topic": "Which scope?", "options": ["small", "large"],
2400 "recommended_default": "other"
2401 }] });
2402 assert!(reduce_event_to_ops(&paths, &open).is_err());
2403 }
2404
2405 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
2406 let mut ev = event(run_id);
2407 ev.seq = seq;
2408 ev.kind = KIND_MERGE_STARTED.into();
2409 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2410 ev.data = serde_json::json!({
2411 "op_id": op_id,
2412 "source_branch": "main",
2413 "worker_branch": "wt/worker",
2414 "expected_source_oid": expected,
2415 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2416 "base_sha": null,
2417 "driver_pid": 4242,
2418 "driver_pid_start_secs": null,
2419 "started_at": "2026-08-15T00:00:00Z",
2420 });
2421 ev
2422 }
2423
2424 #[test]
2427 fn merge_started_records_pending_transaction() {
2428 let tmp = TempDir::new().unwrap();
2429 let run_id = "01jxsnap000000000000000000";
2430 let paths = seed_run_with_node(&tmp, run_id);
2431
2432 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2433 let n = read_n0001(&paths);
2434 assert_eq!(
2435 n.status,
2436 Status::Pending,
2437 "recording a merge is not terminal"
2438 );
2439 let txn = n.pending_merge.expect("transaction recorded");
2440 assert_eq!(txn.op_id, "op-1");
2441 assert_eq!(txn.expected_source_oid, "aaa");
2442 }
2443
2444 #[test]
2447 fn merge_aborted_clears_matching_transaction_only() {
2448 let tmp = TempDir::new().unwrap();
2449 let run_id = "01jxsnap000000000000000000";
2450 let paths = seed_run_with_node(&tmp, run_id);
2451 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2452
2453 let mut stale = event(run_id);
2455 stale.seq = 4;
2456 stale.kind = KIND_MERGE_ABORTED.into();
2457 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2458 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2459 apply_event(&paths, &stale).unwrap();
2460 assert!(
2461 read_n0001(&paths).pending_merge.is_some(),
2462 "stale abort is a no-op"
2463 );
2464
2465 let mut abort = event(run_id);
2467 abort.seq = 5;
2468 abort.kind = KIND_MERGE_ABORTED.into();
2469 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2470 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2471 apply_event(&paths, &abort).unwrap();
2472 let n = read_n0001(&paths);
2473 assert!(
2474 n.pending_merge.is_none(),
2475 "matching abort clears the transaction"
2476 );
2477 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2478 }
2479
2480 #[test]
2483 fn terminal_report_clears_pending_merge() {
2484 let tmp = TempDir::new().unwrap();
2485 let run_id = "01jxsnap000000000000000000";
2486 let paths = seed_run_with_node(&tmp, run_id);
2487 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2488
2489 let mut report = event(run_id);
2490 report.seq = 4;
2491 report.kind = "node.report".into();
2492 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2493 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2494 apply_event(&paths, &report).unwrap();
2495 let n = read_n0001(&paths);
2496 assert_eq!(n.status, Status::Done);
2497 assert!(
2498 n.pending_merge.is_none(),
2499 "completed merge clears the transaction"
2500 );
2501 }
2502
2503 #[test]
2512 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2513 let tmp = TempDir::new().unwrap();
2514
2515 let drive = |run_id: &str, report_data: Value| -> Status {
2518 let paths = seed_run_with_node(&tmp, run_id);
2519 let mut fail = event(run_id);
2521 fail.seq = 3;
2522 fail.kind = "node.status".into();
2523 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2524 fail.data = serde_json::json!({ "status": "failed" });
2525 apply_event(&paths, &fail).unwrap();
2526 assert_eq!(read_n0001(&paths).status, Status::Failed);
2527 let mut report = event(run_id);
2529 report.seq = 4;
2530 report.kind = "node.report".into();
2531 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2532 report.data = report_data;
2533 apply_event(&paths, &report).unwrap();
2534 read_n0001(&paths).status
2535 };
2536
2537 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2539 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2540 assert_eq!(
2541 drive("01jxsnap000000000000000001", agent_forged),
2542 Status::Failed,
2543 "an Agent-origin report with a forged via must not be adopted"
2544 );
2545
2546 let malformed = serde_json::json!({
2548 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2549 });
2550 assert_eq!(
2551 drive("01jxsnap000000000000000002", malformed),
2552 Status::Failed,
2553 "a malformed origin must not re-unlock the legacy via adoption path"
2554 );
2555
2556 let mut run_merge = serde_json::json!({ "success": true });
2558 crate::ReportOrigin::RunMerge {
2559 op_id: Some("op-1".into()),
2560 worker_oid: Some("cafebabe".into()),
2561 }
2562 .stamp(&mut run_merge);
2563 assert_eq!(
2564 drive("01jxsnap000000000000000003", run_merge),
2565 Status::Done,
2566 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2567 );
2568
2569 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2572 assert_eq!(
2573 drive("01jxsnap000000000000000004", legacy),
2574 Status::Done,
2575 "a legacy via-only report (no origin field) is still adopted"
2576 );
2577 }
2578
2579 #[test]
2583 fn terminal_node_status_clears_pending_merge() {
2584 let tmp = TempDir::new().unwrap();
2585 let run_id = "01jxsnap000000000000000000";
2586 let paths = seed_run_with_node(&tmp, run_id);
2587 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2588 assert!(read_n0001(&paths).pending_merge.is_some());
2589
2590 let mut status = event(run_id);
2591 status.seq = 4;
2592 status.kind = "node.status".into();
2593 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2594 status.data = serde_json::json!({ "status": "failed" });
2595 apply_event(&paths, &status).unwrap();
2596 let n = read_n0001(&paths);
2597 assert_eq!(n.status, Status::Failed);
2598 assert!(
2599 n.pending_merge.is_none(),
2600 "terminal status clears the transaction"
2601 );
2602 }
2603
2604 #[test]
2608 fn supervisor_attached_sets_supervisor_pid() {
2609 let tmp = TempDir::new().unwrap();
2610 let run_id = "01jxsnap000000000000000000";
2611 let paths = seed_run_with_node(&tmp, run_id);
2612 assert_eq!(read_n0001(&paths).supervisor_pid, None);
2613
2614 let mut ev = event(run_id);
2615 ev.seq = 3;
2616 ev.kind = "supervisor.attached".into();
2617 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2618 ev.data = serde_json::json!({ "pid": 47820 });
2619 apply_event(&paths, &ev).expect("supervisor.attached applies");
2620 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2621 }
2622
2623 #[test]
2626 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2627 let tmp = TempDir::new().unwrap();
2628 let run_id = "01jxsnap000000000000000000";
2629 let paths = seed_run_with_node(&tmp, run_id);
2630
2631 let mut ev = event(run_id);
2632 ev.kind = "supervisor.attached".into();
2633 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2634
2635 ev.seq = 3;
2636 ev.data = serde_json::json!({ "pid": 100 });
2637 apply_event(&paths, &ev).expect("first attach applies");
2638 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2639
2640 ev.seq = 4;
2642 ev.data = serde_json::json!({ "pid": 200 });
2643 apply_event(&paths, &ev).expect("second attach applies");
2644 let after_second = read_n0001(&paths);
2645 assert_eq!(after_second.supervisor_pid, Some(200));
2646
2647 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2650 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2651 apply_event(&paths, &ev).expect("replay applies as no-op");
2652 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2653 }
2654
2655 #[test]
2658 fn supervisor_cursor_advanced_sets_report_cursor() {
2659 let tmp = TempDir::new().unwrap();
2660 let run_id = "01jxsnap000000000000000000";
2661 let paths = seed_run_with_node(&tmp, run_id);
2662 let child = "02jxsnap000000000000000000";
2663
2664 let mut ev = event(run_id);
2665 ev.seq = 3;
2666 ev.kind = "supervisor.cursor_advanced".into();
2667 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2668 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2669 apply_event(&paths, &ev).expect("cursor_advanced applies");
2670
2671 let n = read_n0001(&paths);
2672 assert_eq!(
2673 n.last_processed_report_seq_by_child.get(child),
2674 Some(&Value::from(7u64))
2675 );
2676 }
2677
2678 #[test]
2683 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2684 let tmp = TempDir::new().unwrap();
2685 let run_id = "01jxsnap000000000000000000";
2686 let paths = seed_run_with_node(&tmp, run_id);
2687 let child_a = "02jxsnap000000000000000000";
2688 let child_b = "03jxsnap000000000000000000";
2689
2690 let mut ev = event(run_id);
2691 ev.kind = "supervisor.cursor_advanced".into();
2692 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2693
2694 ev.seq = 3;
2695 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2696 apply_event(&paths, &ev).expect("seq 5 applies");
2697
2698 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2700 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2701
2702 ev.seq = 4;
2704 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2705 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2706 assert!(ops.is_empty(), "older seq must plan no ops");
2707 apply_event(&paths, &ev).expect("older seq applies as no-op");
2708 assert_eq!(
2709 read_n0001(&paths)
2710 .last_processed_report_seq_by_child
2711 .get(child_a),
2712 Some(&Value::from(5u64))
2713 );
2714
2715 ev.seq = 5;
2717 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2718 apply_event(&paths, &ev).expect("higher seq applies");
2719 ev.seq = 6;
2720 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2721 apply_event(&paths, &ev).expect("second child applies");
2722
2723 let n = read_n0001(&paths);
2724 assert_eq!(
2725 n.last_processed_report_seq_by_child.get(child_a),
2726 Some(&Value::from(9u64))
2727 );
2728 assert_eq!(
2729 n.last_processed_report_seq_by_child.get(child_b),
2730 Some(&Value::from(1u64))
2731 );
2732 }
2733
2734 #[test]
2737 fn supervisor_state_events_reject_malformed_payloads() {
2738 let tmp = TempDir::new().unwrap();
2739 let run_id = "01jxsnap000000000000000000";
2740 let paths = seed_run_with_node(&tmp, run_id);
2741 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2742
2743 let mut ev = event(run_id);
2745 ev.seq = 3;
2746 ev.kind = "supervisor.attached".into();
2747 ev.node_id = nid.clone();
2748 ev.data = serde_json::json!({});
2749 assert!(matches!(
2750 reduce_event_to_ops(&paths, &ev),
2751 Err(Error::CorruptEventLog { .. })
2752 ));
2753
2754 ev.node_id = None;
2756 ev.data = serde_json::json!({ "pid": 1 });
2757 assert!(matches!(
2758 reduce_event_to_ops(&paths, &ev),
2759 Err(Error::CorruptEventLog { .. })
2760 ));
2761
2762 let mut ev2 = event(run_id);
2764 ev2.seq = 4;
2765 ev2.kind = "supervisor.cursor_advanced".into();
2766 ev2.node_id = nid.clone();
2767 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
2768 assert!(matches!(
2769 reduce_event_to_ops(&paths, &ev2),
2770 Err(Error::CorruptEventLog { .. })
2771 ));
2772
2773 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2775 assert!(matches!(
2776 reduce_event_to_ops(&paths, &ev2),
2777 Err(Error::CorruptEventLog { .. })
2778 ));
2779 }
2780
2781 #[test]
2788 fn removed_or_garbage_kind_in_created_events_is_rejected() {
2789 let tmp = TempDir::new().unwrap();
2790 let run_id = "01jxsnap000000000000000000";
2791 let rid = RunId::parse_str(run_id).unwrap();
2792 let dir = crate::run_dir(tmp.path(), &rid);
2793 std::fs::create_dir_all(&dir).unwrap();
2794 let paths = RunPaths::new(dir, run_id).unwrap();
2795
2796 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2797 let mut ev = event(run_id);
2798 ev.kind = "run.created".into();
2799 ev.node_id = None;
2800 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2801 assert!(
2802 matches!(
2803 reduce_event_to_ops(&paths, &ev),
2804 Err(Error::CorruptEventLog { .. })
2805 ),
2806 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2807 );
2808 }
2809
2810 let mut ok = event(run_id);
2813 ok.kind = "run.created".into();
2814 ok.node_id = None;
2815 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2816 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2817 }
2818
2819 #[cfg(unix)]
2828 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2829 use std::os::unix::fs::MetadataExt;
2830 let mut consider = vec![paths.manifest()];
2831 for dir in [paths.nodes_dir()] {
2832 if let Ok(rd) = std::fs::read_dir(&dir) {
2833 for ent in rd.flatten() {
2834 let p = ent.path();
2835 if p.extension().and_then(|s| s.to_str()) == Some("json") {
2836 consider.push(p);
2837 }
2838 }
2839 }
2840 }
2841 let mut map = std::collections::BTreeMap::new();
2842 for p in consider {
2843 if let Ok(md) = std::fs::symlink_metadata(&p) {
2844 if md.file_type().is_file() {
2845 map.insert(p, md.ino());
2846 }
2847 }
2848 }
2849 map
2850 }
2851
2852 #[cfg(unix)]
2861 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2862 use std::collections::BTreeSet;
2863 let before = projection_inodes(paths);
2864 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2865 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2866 .into_iter()
2867 .collect();
2868 apply_event(paths, ev)
2869 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2870 let after = projection_inodes(paths);
2871 let touched: BTreeSet<PathBuf> = after
2872 .iter()
2873 .filter(|(p, ino)| before.get(*p) != Some(*ino))
2874 .map(|(p, _)| p.clone())
2875 .collect();
2876 assert_eq!(
2877 planned, touched,
2878 "kind={}: plan_projections must name exactly the files apply_event writes",
2879 ev.kind
2880 );
2881 if expect_writes {
2882 assert!(
2883 !touched.is_empty(),
2884 "kind={}: expected this event to write at least one projection",
2885 ev.kind
2886 );
2887 }
2888 }
2889
2890 #[cfg(unix)]
2897 #[test]
2898 fn plan_projections_matches_apply_for_every_kind() {
2899 let tmp = TempDir::new().unwrap();
2900 let run_id = "01jxsnap000000000000000000";
2901 let rid = RunId::parse_str(run_id).unwrap();
2902 let dir = crate::run_dir(tmp.path(), &rid);
2903 std::fs::create_dir_all(&dir).unwrap();
2904 let paths = RunPaths::new(dir, run_id).unwrap();
2905 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2906 let child = "02jxsnap000000000000000000";
2907
2908 let mut next_seq = 0u64;
2910 let mut at = |kind: &str, node_id, data| {
2911 next_seq += 1;
2912 Event {
2913 ts: Utc::now(),
2914 seq: next_seq,
2915 kind: kind.into(),
2916 run_id: rid.clone(),
2917 node_id,
2918 idempotency_key: None,
2919 data,
2920 }
2921 };
2922
2923 assert_plan_matches_apply(
2925 &paths,
2926 &at(
2927 "run.created",
2928 None,
2929 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2930 ),
2931 true,
2932 );
2933 assert_plan_matches_apply(
2935 &paths,
2936 &at(
2937 "run.status",
2938 None,
2939 serde_json::json!({ "status": "running" }),
2940 ),
2941 true,
2942 );
2943 assert_plan_matches_apply(
2945 &paths,
2946 &at(
2947 "node.created",
2948 nid(),
2949 serde_json::json!({ "kind": "spinoff" }),
2950 ),
2951 true,
2952 );
2953 assert_plan_matches_apply(
2955 &paths,
2956 &at(
2957 "node.status",
2958 nid(),
2959 serde_json::json!({ "status": "running" }),
2960 ),
2961 true,
2962 );
2963 assert_plan_matches_apply(
2965 &paths,
2966 &at(
2967 "supervisor.attached",
2968 nid(),
2969 serde_json::json!({ "pid": 4242 }),
2970 ),
2971 true,
2972 );
2973 assert_plan_matches_apply(
2975 &paths,
2976 &at(
2977 "supervisor.cursor_advanced",
2978 nid(),
2979 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2980 ),
2981 true,
2982 );
2983 assert_plan_matches_apply(
2985 &paths,
2986 &at(
2987 "child.spawned",
2988 nid(),
2989 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2990 ),
2991 true,
2992 );
2993 assert_plan_matches_apply(
2995 &paths,
2996 &at("node.report", nid(), serde_json::json!({ "success": true })),
2997 true,
2998 );
2999 assert_plan_matches_apply(
3002 &paths,
3003 &at(
3004 "node.status",
3005 nid(),
3006 serde_json::json!({ "status": "failed" }),
3007 ),
3008 false,
3009 );
3010 for kind in [
3012 "supervisor.exited",
3013 "orchestrator.decision",
3014 "discuss.critical",
3015 "cleanup.window_missing",
3016 ] {
3017 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
3018 }
3019 }
3020
3021 #[test]
3025 fn worker_exited_records_clean_exit_without_transitioning_status() {
3026 let tmp = TempDir::new().unwrap();
3027 let run_id = "01jxsnap000000000000000000";
3028 let paths = bootstrap_retry_node(&tmp, run_id);
3029
3030 let mut ev = event(run_id);
3031 ev.seq = 3;
3032 ev.kind = "worker.exited".into();
3033 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
3034 ev.data = serde_json::json!({ "exit_code": 0 });
3035 apply_event(&paths, &ev).expect("worker.exited applies");
3036
3037 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
3038 .unwrap()
3039 .unwrap();
3040 let exit = n.worker_exit.expect("worker_exit recorded");
3041 assert_eq!(exit.code, Some(0));
3042 assert_eq!(exit.signal, None);
3043 assert!(exit.is_clean());
3044 assert_eq!(
3045 n.status,
3046 Status::Pending,
3047 "the exit fact never transitions status"
3048 );
3049 }
3050
3051 #[test]
3055 fn worker_exited_records_signal_and_is_first_write_wins() {
3056 let tmp = TempDir::new().unwrap();
3057 let run_id = "01jxsnap000000000000000000";
3058 let paths = bootstrap_retry_node(&tmp, run_id);
3059 let nid = NodeId::parse_str("n-0001").unwrap();
3060
3061 let mut ev = event(run_id);
3062 ev.seq = 3;
3063 ev.kind = "worker.exited".into();
3064 ev.node_id = Some(nid.clone());
3065 ev.data = serde_json::json!({ "signal": 9 });
3066 apply_event(&paths, &ev).expect("worker.exited applies");
3067
3068 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3069 let exit = n.worker_exit.expect("worker_exit recorded");
3070 assert_eq!(exit.signal, Some(9));
3071 assert!(exit.is_failure());
3072
3073 let mut dup = event(run_id);
3076 dup.seq = 4;
3077 dup.kind = "worker.exited".into();
3078 dup.node_id = Some(nid.clone());
3079 dup.data = serde_json::json!({ "exit_code": 0 });
3080 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
3081 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
3082 assert_eq!(
3083 n2.worker_exit.unwrap().signal,
3084 Some(9),
3085 "first-write-wins: the replayed exit must not overwrite the recorded fact"
3086 );
3087 }
3088
3089 #[test]
3093 fn worker_exited_without_code_or_signal_is_corrupt() {
3094 let tmp = TempDir::new().unwrap();
3095 let run_id = "01jxsnap000000000000000000";
3096 let paths = bootstrap_retry_node(&tmp, run_id);
3097
3098 let mut ev = event(run_id);
3099 ev.seq = 3;
3100 ev.kind = "worker.exited".into();
3101 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
3102 ev.data = serde_json::json!({});
3103 match reduce_event_to_ops(&paths, &ev) {
3104 Err(Error::CorruptEventLog { .. }) => {}
3105 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
3106 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3107 }
3108
3109 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
3112 match reduce_event_to_ops(&paths, &ev) {
3113 Err(Error::CorruptEventLog { .. }) => {}
3114 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
3115 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3116 }
3117 }
3118
3119 #[test]
3125 fn node_death_observed_records_first_death_first_write_wins() {
3126 let tmp = TempDir::new().unwrap();
3127 let run_id = "01jxsnap000000000000000000";
3128 let paths = bootstrap_retry_node(&tmp, run_id);
3129 let nid = NodeId::parse_str("n-0001").unwrap();
3130
3131 let mut ev = event(run_id);
3132 ev.seq = 3;
3133 ev.kind = "node.death_observed".into();
3134 ev.node_id = Some(nid.clone());
3135 ev.data = serde_json::json!({});
3136 apply_event(&paths, &ev).expect("node.death_observed applies");
3137 let first = read_node_opt(&paths, &nid)
3138 .unwrap()
3139 .unwrap()
3140 .first_death_at
3141 .expect("first_death_at recorded");
3142 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
3143
3144 let mut later = event(run_id);
3146 later.seq = 4;
3147 later.kind = "node.death_observed".into();
3148 later.node_id = Some(nid.clone());
3149 later.ts = ev.ts + chrono::Duration::seconds(30);
3150 later.data = serde_json::json!({});
3151 apply_event(&paths, &later).expect("re-observation applies as no-op");
3152 assert_eq!(
3153 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3154 Some(first),
3155 "first-write-wins: a re-observation must not reset the anchor"
3156 );
3157 }
3158
3159 #[test]
3164 fn node_death_observed_noop_when_worker_exit_present() {
3165 let tmp = TempDir::new().unwrap();
3166 let run_id = "01jxsnap000000000000000000";
3167 let paths = bootstrap_retry_node(&tmp, run_id);
3168 let nid = NodeId::parse_str("n-0001").unwrap();
3169
3170 let mut exit = event(run_id);
3172 exit.seq = 3;
3173 exit.kind = "worker.exited".into();
3174 exit.node_id = Some(nid.clone());
3175 exit.data = serde_json::json!({ "exit_code": 0 });
3176 apply_event(&paths, &exit).unwrap();
3177
3178 let mut death = event(run_id);
3180 death.seq = 4;
3181 death.kind = "node.death_observed".into();
3182 death.node_id = Some(nid.clone());
3183 death.data = serde_json::json!({});
3184 apply_event(&paths, &death).expect("applies as no-op");
3185 assert_eq!(
3186 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3187 None,
3188 "a told worker.exited makes the crash backstop moot; no anchor recorded"
3189 );
3190 }
3191
3192 #[test]
3197 fn node_retry_clears_worker_exit() {
3198 let tmp = TempDir::new().unwrap();
3199 let run_id = "01jxsnap000000000000000000";
3200 let paths = bootstrap_retry_node(&tmp, run_id);
3201 let nid = NodeId::parse_str("n-0001").unwrap();
3202
3203 let mut exit = event(run_id);
3205 exit.seq = 3;
3206 exit.kind = "worker.exited".into();
3207 exit.node_id = Some(nid.clone());
3208 exit.data = serde_json::json!({ "exit_code": 7 });
3209 apply_event(&paths, &exit).unwrap();
3210 assert!(read_node_opt(&paths, &nid)
3211 .unwrap()
3212 .unwrap()
3213 .worker_exit
3214 .is_some());
3215
3216 let mut retry = event(run_id);
3218 retry.seq = 4;
3219 retry.kind = "node.retry".into();
3220 retry.node_id = Some(nid.clone());
3221 retry.data = serde_json::json!({
3222 "attempt": 1,
3223 "reason": "agent-died",
3224 "branch": "wt/foo",
3225 "worktree_path": "/tmp/new-wt",
3226 "agent_pid": 222,
3227 });
3228 apply_event(&paths, &retry).unwrap();
3229
3230 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3231 assert!(
3232 n.worker_exit.is_none(),
3233 "node.retry must clear the previous attempt's worker_exit"
3234 );
3235 assert_eq!(
3236 n.status,
3237 Status::Pending,
3238 "retry returns the node to Pending"
3239 );
3240 }
3241}