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, IdValidationError, Kind, Lifecycle, Manifest, MergeTxn, Node, NodeId, RunId,
59 Status, TmuxIdentity, 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 "node.death_observed" => reduce_node_death_observed(paths, ev),
300 KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
301 KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
302 "child.spawned" => reduce_child_spawned(paths, ev),
303 "supervisor.attached" => reduce_supervisor_attached(paths, ev),
304 "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
305 "supervisor.exited" => Ok(vec![]),
306 "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
314 "run.notified" => Ok(vec![]),
321 "cleanup.window_missing"
346 | "cleanup.worktree_missing"
347 | "cleanup.branch_remove_failed"
348 | "cleanup.branch_preserved"
349 | "cleanup.session_killed"
350 | "cleanup.session_retained" => Ok(vec![]),
351 "supervisor.child_id_quarantined" => Ok(vec![]),
361 _ => Ok(vec![]),
362 }
363}
364
365fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
371 match op {
372 ProjectionOp::Manifest(_) => paths.manifest(),
373 ProjectionOp::Node(n) => paths.node(&n.node_id),
374 }
375}
376
377pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
399 let ops = reduce_event_to_ops(paths, event)?;
400 Ok(ops.iter().map(|op| op_path(paths, op)).collect())
401}
402
403pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
419 let ops = reduce_event_to_ops(paths, ev)?;
420 commit_ops(paths, ops)
421}
422
423#[cfg(test)]
432pub(crate) fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
433 reduce_event_to_ops(paths, ev).map(|_| ())
434}
435
436fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
440 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
441 path: events_path.to_path_buf(),
442 reason: format!(
443 "event seq={} kind={} missing top-level `node_id`",
444 ev.seq, ev.kind
445 ),
446 })
447}
448
449fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
450 if let Some(existing) = read_manifest_opt(paths)? {
454 if existing.run_id != ev.run_id {
455 return Err(Error::CorruptEventLog {
456 path: paths.manifest(),
457 reason: format!(
458 "run.created run_id={} conflicts with existing manifest run_id={}",
459 ev.run_id, existing.run_id
460 ),
461 });
462 }
463 return Ok(vec![]);
464 }
465 let events_path = paths.events();
466 let d = &ev.data;
467 let kind =
468 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
469 path: events_path.clone(),
470 reason: "run.created missing/invalid `kind`".into(),
471 })?;
472 let lifecycle: Lifecycle = serde_json::from_value(
473 d.get("lifecycle").cloned().unwrap_or(Value::Null),
474 )
475 .map_err(|_| Error::CorruptEventLog {
476 path: events_path.clone(),
477 reason: "run.created missing/invalid `lifecycle`".into(),
478 })?;
479 let title = want_str(&events_path, ev, d, "title")?.to_string();
480 let m = Manifest {
481 schema_version: STATE_SCHEMA_VERSION,
482 applied_seq: 0,
485 run_id: paths.run_id.clone(),
487 kind,
488 lifecycle,
489 title,
490 status: Status::Pending,
491 created_at: ev.ts,
492 updated_at: ev.ts,
493 source_repo: d
494 .get("source_repo")
495 .and_then(Value::as_str)
496 .map(str::to_string),
497 source_branch: d
498 .get("source_branch")
499 .and_then(Value::as_str)
500 .map(str::to_string),
501 worktree_root: d
502 .get("worktree_root")
503 .and_then(Value::as_str)
504 .map(str::to_string),
505 managed_tmux_session: d
506 .get("managed_tmux_session")
507 .and_then(Value::as_str)
508 .map(str::to_string),
509 notify_cmd: d
510 .get("notify_cmd")
511 .and_then(Value::as_str)
512 .map(str::to_string),
513 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
514 node_count: 0,
515 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
516 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
517 };
518 Ok(vec![ProjectionOp::Manifest(m)])
519}
520
521fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
522 let mut m = match read_manifest_opt(paths)? {
523 Some(m) => m,
524 None => return Ok(vec![]),
525 };
526 let new_status = require_status(ev, paths.events())?;
527 if m.status.is_terminal() {
530 trace_terminal_noop(ev, m.status, new_status);
531 return Ok(vec![]);
532 }
533 if m.status == new_status {
534 return Ok(vec![]);
535 }
536 m.status = new_status;
537 m.updated_at = ev.ts;
538 Ok(vec![ProjectionOp::Manifest(m)])
539}
540
541fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
552 let nonempty = |key| {
553 d.get(key)
554 .and_then(Value::as_str)
555 .map(str::trim)
556 .filter(|s| !s.is_empty())
557 .map(str::to_string)
558 };
559 let session = nonempty("tmux_session")?;
560 let window_id = nonempty("tmux_window_id")?;
561 Some(TmuxIdentity {
562 socket: nonempty("tmux_socket"),
563 session,
564 window_id,
565 pane_id: nonempty("tmux_pane_id"),
568 })
569}
570
571fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
572 let events_path = paths.events();
573 let node_id = require_envelope_node_id(&events_path, ev)?;
576 if read_node_opt(paths, &node_id)?.is_some() {
578 return Ok(vec![]);
579 }
580 let d = &ev.data;
581 let kind =
582 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
583 path: events_path.clone(),
584 reason: format!(
585 "event seq={} kind=node.created missing/invalid `kind`",
586 ev.seq
587 ),
588 })?;
589 let n = Node {
590 schema_version: STATE_SCHEMA_VERSION,
591 node_id,
592 run_id: paths.run_id.clone(),
594 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
595 kind,
596 status: Status::Pending,
597 task: d.get("task").and_then(Value::as_str).map(str::to_string),
598 worktree_path: d
599 .get("worktree_path")
600 .and_then(Value::as_str)
601 .map(str::to_string),
602 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
603 base_sha: d
604 .get("base_sha")
605 .and_then(Value::as_str)
606 .filter(|s| !s.is_empty())
607 .map(str::to_string),
608 tmux_window: d
609 .get("tmux_window")
610 .and_then(Value::as_str)
611 .map(str::to_string),
612 tmux_identity: tmux_identity_from_data(d),
613 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
614 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
615 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
616 children: Vec::new(),
617 started_at: Some(ev.ts),
618 updated_at: ev.ts,
619 last_report: None,
620 last_processed_report_seq_by_child: serde_json::Map::default(),
621 retry_attempts: 0,
622 worker_exit: None,
623 pending_merge: None,
624 first_death_at: None,
625 };
626 let mut ops = vec![ProjectionOp::Node(n)];
627 if let Some(mut m) = read_manifest_opt(paths)? {
628 m.updated_at = ev.ts;
633 ops.push(ProjectionOp::Manifest(m));
634 }
635 Ok(ops)
636}
637
638fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
657 let events_path = paths.events();
658 let node_id = require_envelope_node_id(&events_path, ev)?;
659 let mut n = match read_node_opt(paths, &node_id)? {
660 Some(n) => n,
661 None => return Ok(vec![]),
662 };
663 if n.status.is_terminal() {
666 tracing::debug!(
667 target: "octl_core::reducer",
668 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
669 "no-op: node.retry against terminal node"
670 );
671 return Ok(vec![]);
672 }
673 let d = &ev.data;
674 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
677 n.base_sha = d
678 .get("base_sha")
679 .and_then(Value::as_str)
680 .filter(|s| !s.is_empty())
681 .map(str::to_string);
682 n.worktree_path = d
683 .get("worktree_path")
684 .and_then(Value::as_str)
685 .map(str::to_string);
686 n.tmux_window = d
687 .get("tmux_window")
688 .and_then(Value::as_str)
689 .map(str::to_string);
690 n.tmux_identity = tmux_identity_from_data(d);
691 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
692 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
693 n.status = Status::Pending;
694 n.started_at = Some(ev.ts);
695 n.updated_at = ev.ts;
696 n.last_report = None;
697 n.pending_merge = None;
703 n.worker_exit = None;
708 n.first_death_at = None;
713 n.retry_attempts = d
721 .get("attempt")
722 .and_then(Value::as_u64)
723 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
724 Ok(vec![ProjectionOp::Node(n)])
725}
726
727fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
728 let events_path = paths.events();
729 let node_id = require_envelope_node_id(&events_path, ev)?;
730 let mut n = match read_node_opt(paths, &node_id)? {
731 Some(n) => n,
732 None => return Ok(vec![]),
733 };
734 let new_status = require_status(ev, events_path)?;
735 if n.status.is_terminal() {
738 trace_terminal_noop(ev, n.status, new_status);
739 return Ok(vec![]);
740 }
741 if n.status == new_status {
742 return Ok(vec![]);
743 }
744 n.status = new_status;
745 if new_status.is_terminal() {
751 n.pending_merge = None;
752 }
753 n.updated_at = ev.ts;
754 Ok(vec![ProjectionOp::Node(n)])
755}
756
757fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
758 let events_path = paths.events();
759 let node_id = require_envelope_node_id(&events_path, ev)?;
760 let mut n = match read_node_opt(paths, &node_id)? {
761 Some(n) => n,
762 None => return Ok(vec![]),
763 };
764 if n.status.is_terminal() {
776 if matches!(n.status, Status::Failed | Status::Done)
815 && report_is_confirmed_explicit_merge(&ev.data)
816 {
817 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
818 return Ok(vec![]);
819 }
820 tracing::info!(
821 target: "octl_core::reducer",
822 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
823 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
824 );
825 n.last_report = Some(ev.data.clone());
826 n.status = Status::Done;
830 n.pending_merge = None;
834 n.updated_at = ev.ts;
835 return Ok(vec![ProjectionOp::Node(n)]);
836 }
837 tracing::debug!(
838 target: "octl_core::reducer",
839 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
840 "no-op: node.report against terminal node"
841 );
842 return Ok(vec![]);
843 }
844 let new_status = report_terminal_status(&events_path, ev)?;
852 n.last_report = Some(ev.data.clone());
853 n.status = new_status;
854 n.pending_merge = None;
859 n.updated_at = ev.ts;
860 Ok(vec![ProjectionOp::Node(n)])
861}
862
863fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
878 let events_path = paths.events();
879 let node_id = require_envelope_node_id(&events_path, ev)?;
880 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
881 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
882 match (code, signal) {
887 (Some(_), None) | (None, Some(_)) => {}
888 _ => {
889 return Err(Error::CorruptEventLog {
890 path: events_path,
891 reason: format!(
892 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
893 ev.seq
894 ),
895 });
896 }
897 }
898 let mut n = match read_node_opt(paths, &node_id)? {
899 Some(n) => n,
900 None => return Ok(vec![]),
907 };
908 if n.worker_exit.is_some() {
912 return Ok(vec![]);
913 }
914 n.worker_exit = Some(WorkerExit {
915 code,
916 signal,
917 at: ev.ts,
918 });
919 n.updated_at = ev.ts;
920 Ok(vec![ProjectionOp::Node(n)])
921}
922
923fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
936 let events_path = paths.events();
937 let node_id = require_envelope_node_id(&events_path, ev)?;
938 let mut n = match read_node_opt(paths, &node_id)? {
939 Some(n) => n,
940 None => return Ok(vec![]),
941 };
942 if n.first_death_at.is_some()
948 || n.status.is_terminal()
949 || n.worker_exit.is_some()
950 || n.last_report.is_some()
951 || n.pending_merge.is_some()
952 {
953 return Ok(vec![]);
954 }
955 n.first_death_at = Some(ev.ts);
956 n.updated_at = ev.ts;
957 Ok(vec![ProjectionOp::Node(n)])
958}
959
960pub const KIND_MERGE_STARTED: &str = "merge.started";
964
965pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
970
971fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
985 let events_path = paths.events();
986 let node_id = require_envelope_node_id(&events_path, ev)?;
987 let txn: MergeTxn =
988 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
989 path: events_path.clone(),
990 reason: format!(
991 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
992 ev.seq
993 ),
994 })?;
995 let mut n = match read_node_opt(paths, &node_id)? {
996 Some(n) => n,
997 None => return Ok(vec![]),
998 };
999 if n.status.is_terminal() {
1003 return Ok(vec![]);
1004 }
1005 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1008 return Ok(vec![]);
1009 }
1010 n.pending_merge = Some(Box::new(txn));
1011 n.updated_at = ev.ts;
1012 Ok(vec![ProjectionOp::Node(n)])
1013}
1014
1015fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1027 let events_path = paths.events();
1028 let node_id = require_envelope_node_id(&events_path, ev)?;
1029 let op_id = ev
1030 .data
1031 .get("op_id")
1032 .and_then(Value::as_str)
1033 .ok_or_else(|| Error::CorruptEventLog {
1034 path: events_path.clone(),
1035 reason: format!(
1036 "event seq={} kind=merge.aborted is missing string `op_id`",
1037 ev.seq
1038 ),
1039 })?;
1040 let mut n = match read_node_opt(paths, &node_id)? {
1041 Some(n) => n,
1042 None => return Ok(vec![]),
1043 };
1044 match n.pending_merge.as_ref() {
1047 Some(t) if t.op_id == op_id => {}
1048 _ => return Ok(vec![]),
1049 }
1050 n.pending_merge = None;
1051 n.updated_at = ev.ts;
1052 Ok(vec![ProjectionOp::Node(n)])
1053}
1054
1055fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1063 if current == incoming {
1064 tracing::debug!(
1065 target: "octl_core::reducer",
1066 seq = ev.seq, kind = %ev.kind, status = ?current,
1067 "no-op: status re-applied to terminal target"
1068 );
1069 } else {
1070 tracing::warn!(
1071 target: "octl_core::reducer",
1072 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1073 "no-op: ignored conflicting transition from terminal target"
1074 );
1075 }
1076}
1077
1078fn report_is_confirmed_explicit_merge(data: &Value) -> bool {
1095 ReportOrigin::report_is_confirmed_merge(data)
1096}
1097
1098fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1107 let corrupt = |reason: String| Error::CorruptEventLog {
1108 path: events_path.to_path_buf(),
1109 reason,
1110 };
1111 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1112 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1113 if cancelled {
1114 if success == Some(true) {
1115 return Err(corrupt(format!(
1116 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1117 ev.seq
1118 )));
1119 }
1120 Ok(Status::Cancelled)
1121 } else {
1122 match success {
1123 Some(true) => Ok(Status::Done),
1124 Some(false) => Ok(Status::Failed),
1125 None => Err(corrupt(format!(
1126 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1127 ev.seq
1128 ))),
1129 }
1130 }
1131}
1132
1133fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1134 let events_path = paths.events();
1137 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1138 path: events_path.clone(),
1139 reason: format!(
1140 "event seq={} kind=child.spawned missing parent `node_id`",
1141 ev.seq
1142 ),
1143 })?;
1144 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1145 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1146 let child_node_id = NodeId::parse_str(
1147 ev.data
1148 .get("child_node_id")
1149 .and_then(Value::as_str)
1150 .unwrap_or("n-0001"),
1151 )
1152 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1153 let mut n = match read_node_opt(paths, &parent_node_id)? {
1154 Some(n) => n,
1155 None => return Ok(vec![]),
1156 };
1157 let new_ref = ChildRef {
1158 run_id: child_run_id,
1159 node_id: child_node_id,
1160 };
1161 if n.children.iter().any(|c| c == &new_ref) {
1162 return Ok(vec![]);
1165 }
1166 n.children.push(new_ref);
1167 n.updated_at = ev.ts;
1168 Ok(vec![ProjectionOp::Node(n)])
1169}
1170
1171fn reduce_supervisor_attached(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 raw = ev
1185 .data
1186 .get("pid")
1187 .and_then(Value::as_i64)
1188 .ok_or_else(|| Error::CorruptEventLog {
1189 path: events_path.clone(),
1190 reason: format!(
1191 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1192 ev.seq
1193 ),
1194 })?;
1195 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1196 path: events_path.clone(),
1197 reason: format!(
1198 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1199 ev.seq
1200 ),
1201 })?;
1202 let mut n = match read_node_opt(paths, &node_id)? {
1203 Some(n) => n,
1204 None => return Ok(vec![]),
1205 };
1206 if n.supervisor_pid == Some(pid) {
1207 return Ok(vec![]);
1208 }
1209 n.supervisor_pid = Some(pid);
1210 n.updated_at = ev.ts;
1211 Ok(vec![ProjectionOp::Node(n)])
1212}
1213
1214fn reduce_supervisor_cursor_advanced(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 child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1228 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1232 let report_seq = ev
1233 .data
1234 .get("report_seq")
1235 .and_then(Value::as_u64)
1236 .ok_or_else(|| Error::CorruptEventLog {
1237 path: events_path.clone(),
1238 reason: format!(
1239 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1240 ev.seq
1241 ),
1242 })?;
1243 let mut n = match read_node_opt(paths, &node_id)? {
1244 Some(n) => n,
1245 None => return Ok(vec![]),
1246 };
1247 if let Some(prev) = n
1248 .last_processed_report_seq_by_child
1249 .get(child_run_id)
1250 .and_then(Value::as_u64)
1251 {
1252 if report_seq <= prev {
1253 return Ok(vec![]);
1254 }
1255 }
1256 n.last_processed_report_seq_by_child
1257 .insert(child_run_id.to_string(), Value::from(report_seq));
1258 n.updated_at = ev.ts;
1259 Ok(vec![ProjectionOp::Node(n)])
1260}
1261
1262#[cfg(test)]
1263mod tests {
1264 use super::*;
1265 use crate::schema::Event;
1266 use chrono::Utc;
1267 use tempfile::TempDir;
1268
1269 fn event(run_id: &str) -> Event {
1270 Event {
1271 ts: Utc::now(),
1272 seq: 1,
1273 kind: "run.status".into(),
1274 run_id: RunId::parse_str(run_id).unwrap(),
1275 node_id: None,
1276 idempotency_key: None,
1277 data: serde_json::json!({ "status": "running" }),
1278 }
1279 }
1280
1281 #[test]
1282 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1283 let tmp = TempDir::new().unwrap();
1287 let run_id = "01jxsnap000000000000000000";
1288 let rid = RunId::parse_str(run_id).unwrap();
1289 let dir = crate::run_dir(tmp.path(), &rid);
1290 std::fs::create_dir_all(&dir).unwrap();
1291 let paths = RunPaths::new(dir, run_id).unwrap();
1292
1293 let mut created = event(run_id);
1296 created.kind = "run.created".into();
1297 created.data = serde_json::json!({
1298 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1299 });
1300 apply_event(&paths, &created).expect("run.created applies");
1301 let manifest_before = std::fs::read(paths.manifest()).unwrap();
1302
1303 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1304 let mut ev = event(run_id);
1305 ev.seq = seq;
1306 ev.kind = kind.into();
1307 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1309 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1310 assert!(ops.is_empty(), "{kind} must plan no projection ops");
1311 apply_event(&paths, &ev).expect("audit kind applies as no-op");
1313 }
1314
1315 assert_eq!(
1317 std::fs::read(paths.manifest()).unwrap(),
1318 manifest_before,
1319 "audit events must not mutate the manifest"
1320 );
1321 assert!(!paths.nodes_dir().exists(), "no node projection created");
1322 }
1323
1324 #[test]
1325 fn run_created_folds_harness_when_present_and_defaults_none() {
1326 let tmp = TempDir::new().unwrap();
1327
1328 let run_id = "01jxhrnsaa0000000000000001";
1330 let rid = RunId::parse_str(run_id).unwrap();
1331 let dir = crate::run_dir(tmp.path(), &rid);
1332 std::fs::create_dir_all(&dir).unwrap();
1333 let paths = RunPaths::new(dir, run_id).unwrap();
1334 let mut created = event(run_id);
1335 created.kind = "run.created".into();
1336 created.data = serde_json::json!({
1337 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1338 "harness": "pi", "harness_source": "flag",
1339 });
1340 apply_event(&paths, &created).expect("run.created applies");
1341 let m = read_manifest_opt(&paths).unwrap().unwrap();
1342 assert_eq!(m.harness.as_deref(), Some("pi"));
1343
1344 let run_id2 = "01jxhrnsaa0000000000000002";
1346 let rid2 = RunId::parse_str(run_id2).unwrap();
1347 let dir2 = crate::run_dir(tmp.path(), &rid2);
1348 std::fs::create_dir_all(&dir2).unwrap();
1349 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1350 let mut created2 = event(run_id2);
1351 created2.kind = "run.created".into();
1352 created2.data =
1353 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1354 apply_event(&paths2, &created2).expect("run.created applies");
1355 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1356 assert_eq!(m2.harness, None);
1357 }
1358
1359 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1362 let rid = RunId::parse_str(run_id).unwrap();
1363 let dir = crate::run_dir(tmp.path(), &rid);
1364 std::fs::create_dir_all(&dir).unwrap();
1365 let paths = RunPaths::new(dir, run_id).unwrap();
1366 let mut created = event(run_id);
1367 created.kind = "run.created".into();
1368 created.data =
1369 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1370 apply_event(&paths, &created).expect("run.created applies");
1371 let mut node = event(run_id);
1372 node.seq = 2;
1373 node.kind = "node.created".into();
1374 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1375 node.data = serde_json::json!({
1376 "kind": "spinoff",
1377 "branch": "wt/foo",
1378 "worktree_path": "/tmp/old-wt",
1379 "agent_pid": 111,
1380 });
1381 apply_event(&paths, &node).expect("node.created applies");
1382 paths
1383 }
1384
1385 #[test]
1389 fn node_retry_rewires_node_and_increments_attempts() {
1390 let tmp = TempDir::new().unwrap();
1391 let run_id = "01jxsnap000000000000000000";
1392 let paths = bootstrap_retry_node(&tmp, run_id);
1393
1394 let mut retry = event(run_id);
1395 retry.seq = 3;
1396 retry.kind = "node.retry".into();
1397 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1398 retry.data = serde_json::json!({
1399 "attempt": 1,
1400 "reason": "agent-died",
1401 "branch": "wt/foo-r1",
1402 "base_sha": "a".repeat(40),
1403 "worktree_path": "/tmp/new-wt",
1404 "agent_pid": 222,
1405 "tmux_session": "s",
1406 "tmux_window_id": "@9",
1407 });
1408 apply_event(&paths, &retry).expect("node.retry applies");
1409
1410 let n = read_n0001(&paths);
1411 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1412 assert_eq!(
1413 n.branch.as_deref(),
1414 Some("wt/foo-r1"),
1415 "rewired to new branch"
1416 );
1417 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1418 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1419 assert_eq!(n.status, Status::Pending, "node returns to pending");
1420 assert!(n.last_report.is_none());
1421 assert_eq!(
1422 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1423 Some("@9"),
1424 "rewired tmux identity"
1425 );
1426
1427 let mut retry2 = event(run_id);
1429 retry2.seq = 4;
1430 retry2.kind = "node.retry".into();
1431 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1432 retry2.data = serde_json::json!({
1433 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1434 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1435 });
1436 apply_event(&paths, &retry2).expect("node.retry applies");
1437 assert_eq!(read_n0001(&paths).retry_attempts, 2);
1438 }
1439
1440 #[test]
1444 fn node_retry_against_terminal_node_is_noop() {
1445 let tmp = TempDir::new().unwrap();
1446 let run_id = "01jxsnap000000000000000000";
1447 let paths = bootstrap_retry_node(&tmp, run_id);
1448
1449 let mut report = event(run_id);
1451 report.seq = 3;
1452 report.kind = "node.report".into();
1453 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1454 report.data = serde_json::json!({ "success": true });
1455 apply_event(&paths, &report).expect("node.report applies");
1456 assert_eq!(read_n0001(&paths).status, Status::Done);
1457
1458 let mut retry = event(run_id);
1459 retry.seq = 4;
1460 retry.kind = "node.retry".into();
1461 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1462 retry.data = serde_json::json!({
1463 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1464 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1465 });
1466 apply_event(&paths, &retry).expect("node.retry applies as no-op");
1467
1468 let n = read_n0001(&paths);
1469 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
1470 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
1471 assert_eq!(n.agent_pid, Some(111), "not rewired");
1472 }
1473
1474 #[test]
1475 fn apply_event_rejects_event_from_a_different_run() {
1476 let tmp = TempDir::new().unwrap();
1477 let run_id = "01jxsnap000000000000000000";
1478 let rid = RunId::parse_str(run_id).unwrap();
1479 let dir = crate::run_dir(tmp.path(), &rid);
1480 std::fs::create_dir_all(&dir).unwrap();
1481 let paths = RunPaths::new(dir, run_id).unwrap();
1482
1483 let foreign = event("02jxsnap000000000000000000");
1485 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
1486 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
1487
1488 let mine = event(run_id);
1491 apply_event(&paths, &mine).expect("matching run_id must be accepted");
1492 }
1493
1494 #[test]
1495 fn tmux_identity_from_data_reads_qualified_fields() {
1496 let d = serde_json::json!({
1497 "tmux_socket": "/private/tmp/tmux-501/default",
1498 "tmux_session": "octl",
1499 "tmux_window_id": "@42",
1500 });
1501 let id = tmux_identity_from_data(&d).expect("qualified identity");
1502 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
1503 assert_eq!(id.session, "octl");
1504 assert_eq!(id.window_id, "@42");
1505 assert_eq!(id.pane_id, None);
1507
1508 let d2 = serde_json::json!({
1510 "tmux_socket": null,
1511 "tmux_session": "octl",
1512 "tmux_window_id": "@7",
1513 });
1514 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
1515 assert_eq!(id2.socket, None);
1516 assert_eq!(id2.window_id, "@7");
1517
1518 let d3 = serde_json::json!({
1520 "tmux_session": "octl",
1521 "tmux_window_id": "@42",
1522 "tmux_pane_id": "%7",
1523 });
1524 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
1525 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
1526 assert_eq!(id3.capture_target(), "%7");
1527
1528 let d4 = serde_json::json!({
1531 "tmux_session": "octl",
1532 "tmux_window_id": "@42",
1533 "tmux_pane_id": null,
1534 });
1535 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
1536 assert_eq!(id4.pane_id, None);
1537 assert_eq!(id4.capture_target(), "@42");
1538 }
1539
1540 #[test]
1541 fn tmux_identity_from_data_back_compat_is_none() {
1542 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
1544 assert!(tmux_identity_from_data(&legacy).is_none());
1545 let partial = serde_json::json!({ "tmux_window_id": "@42" });
1547 assert!(tmux_identity_from_data(&partial).is_none());
1548 }
1549
1550 #[test]
1553 fn node_created_populates_tmux_identity() {
1554 let tmp = TempDir::new().unwrap();
1555 let run_id = "01jxsnap000000000000000000";
1556 let rid = RunId::parse_str(run_id).unwrap();
1557 let dir = crate::run_dir(tmp.path(), &rid);
1558 std::fs::create_dir_all(&dir).unwrap();
1559 let paths = RunPaths::new(dir, run_id).unwrap();
1560
1561 let mut ev = event(run_id);
1562 ev.seq = 2;
1563 ev.kind = "node.created".into();
1564 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1565 ev.data = serde_json::json!({
1566 "kind": "spinoff",
1567 "tmux_window": "🚀 wt/x",
1568 "tmux_socket": "/private/tmp/tmux-501/default",
1569 "tmux_session": "octl",
1570 "tmux_window_id": "@42",
1571 });
1572 apply_event(&paths, &ev).expect("node.created applies");
1573 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
1574 .unwrap()
1575 .unwrap();
1576 let id = n.tmux_identity.expect("qualified identity recorded");
1577 assert_eq!(id.session, "octl");
1578 assert_eq!(id.window_id, "@42");
1579 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
1580
1581 let run2 = "02jxsnap000000000000000000";
1583 let rid2 = RunId::parse_str(run2).unwrap();
1584 let dir2 = crate::run_dir(tmp.path(), &rid2);
1585 std::fs::create_dir_all(&dir2).unwrap();
1586 let paths2 = RunPaths::new(dir2, run2).unwrap();
1587 let mut ev2 = event(run2);
1588 ev2.seq = 2;
1589 ev2.kind = "node.created".into();
1590 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1591 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
1592 apply_event(&paths2, &ev2).expect("legacy node.created applies");
1593 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
1594 .unwrap()
1595 .unwrap();
1596 assert!(n2.tmux_identity.is_none());
1597 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
1598 }
1599
1600 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1603 let rid = RunId::parse_str(run_id).unwrap();
1604 let dir = crate::run_dir(tmp.path(), &rid);
1605 std::fs::create_dir_all(&dir).unwrap();
1606 let paths = RunPaths::new(dir, run_id).unwrap();
1607
1608 let mut created = event(run_id);
1609 created.kind = "run.created".into();
1610 created.data = serde_json::json!({
1611 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1612 });
1613 apply_event(&paths, &created).expect("run.created applies");
1614
1615 let mut node = event(run_id);
1616 node.seq = 2;
1617 node.kind = "node.created".into();
1618 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1619 node.data = serde_json::json!({ "kind": "spinoff" });
1620 apply_event(&paths, &node).expect("node.created applies");
1621 paths
1622 }
1623
1624 fn read_n0001(paths: &RunPaths) -> Node {
1625 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
1626 .unwrap()
1627 .unwrap()
1628 }
1629
1630 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
1631 let mut ev = event(run_id);
1632 ev.seq = seq;
1633 ev.kind = KIND_MERGE_STARTED.into();
1634 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1635 ev.data = serde_json::json!({
1636 "op_id": op_id,
1637 "source_branch": "main",
1638 "worker_branch": "wt/worker",
1639 "expected_source_oid": expected,
1640 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
1641 "base_sha": null,
1642 "driver_pid": 4242,
1643 "driver_pid_start_secs": null,
1644 "started_at": "2026-08-15T00:00:00Z",
1645 });
1646 ev
1647 }
1648
1649 #[test]
1652 fn merge_started_records_pending_transaction() {
1653 let tmp = TempDir::new().unwrap();
1654 let run_id = "01jxsnap000000000000000000";
1655 let paths = seed_run_with_node(&tmp, run_id);
1656
1657 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
1658 let n = read_n0001(&paths);
1659 assert_eq!(
1660 n.status,
1661 Status::Pending,
1662 "recording a merge is not terminal"
1663 );
1664 let txn = n.pending_merge.expect("transaction recorded");
1665 assert_eq!(txn.op_id, "op-1");
1666 assert_eq!(txn.expected_source_oid, "aaa");
1667 }
1668
1669 #[test]
1672 fn merge_aborted_clears_matching_transaction_only() {
1673 let tmp = TempDir::new().unwrap();
1674 let run_id = "01jxsnap000000000000000000";
1675 let paths = seed_run_with_node(&tmp, run_id);
1676 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
1677
1678 let mut stale = event(run_id);
1680 stale.seq = 4;
1681 stale.kind = KIND_MERGE_ABORTED.into();
1682 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1683 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
1684 apply_event(&paths, &stale).unwrap();
1685 assert!(
1686 read_n0001(&paths).pending_merge.is_some(),
1687 "stale abort is a no-op"
1688 );
1689
1690 let mut abort = event(run_id);
1692 abort.seq = 5;
1693 abort.kind = KIND_MERGE_ABORTED.into();
1694 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1695 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
1696 apply_event(&paths, &abort).unwrap();
1697 let n = read_n0001(&paths);
1698 assert!(
1699 n.pending_merge.is_none(),
1700 "matching abort clears the transaction"
1701 );
1702 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
1703 }
1704
1705 #[test]
1708 fn terminal_report_clears_pending_merge() {
1709 let tmp = TempDir::new().unwrap();
1710 let run_id = "01jxsnap000000000000000000";
1711 let paths = seed_run_with_node(&tmp, run_id);
1712 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
1713
1714 let mut report = event(run_id);
1715 report.seq = 4;
1716 report.kind = "node.report".into();
1717 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1718 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
1719 apply_event(&paths, &report).unwrap();
1720 let n = read_n0001(&paths);
1721 assert_eq!(n.status, Status::Done);
1722 assert!(
1723 n.pending_merge.is_none(),
1724 "completed merge clears the transaction"
1725 );
1726 }
1727
1728 #[test]
1737 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
1738 let tmp = TempDir::new().unwrap();
1739
1740 let drive = |run_id: &str, report_data: Value| -> Status {
1743 let paths = seed_run_with_node(&tmp, run_id);
1744 let mut fail = event(run_id);
1746 fail.seq = 3;
1747 fail.kind = "node.status".into();
1748 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1749 fail.data = serde_json::json!({ "status": "failed" });
1750 apply_event(&paths, &fail).unwrap();
1751 assert_eq!(read_n0001(&paths).status, Status::Failed);
1752 let mut report = event(run_id);
1754 report.seq = 4;
1755 report.kind = "node.report".into();
1756 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1757 report.data = report_data;
1758 apply_event(&paths, &report).unwrap();
1759 read_n0001(&paths).status
1760 };
1761
1762 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
1764 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
1765 assert_eq!(
1766 drive("01jxsnap000000000000000001", agent_forged),
1767 Status::Failed,
1768 "an Agent-origin report with a forged via must not be adopted"
1769 );
1770
1771 let malformed = serde_json::json!({
1773 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
1774 });
1775 assert_eq!(
1776 drive("01jxsnap000000000000000002", malformed),
1777 Status::Failed,
1778 "a malformed origin must not re-unlock the legacy via adoption path"
1779 );
1780
1781 let mut run_merge = serde_json::json!({ "success": true });
1783 crate::ReportOrigin::RunMerge {
1784 op_id: Some("op-1".into()),
1785 worker_oid: Some("cafebabe".into()),
1786 }
1787 .stamp(&mut run_merge);
1788 assert_eq!(
1789 drive("01jxsnap000000000000000003", run_merge),
1790 Status::Done,
1791 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
1792 );
1793
1794 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
1797 assert_eq!(
1798 drive("01jxsnap000000000000000004", legacy),
1799 Status::Done,
1800 "a legacy via-only report (no origin field) is still adopted"
1801 );
1802 }
1803
1804 #[test]
1808 fn terminal_node_status_clears_pending_merge() {
1809 let tmp = TempDir::new().unwrap();
1810 let run_id = "01jxsnap000000000000000000";
1811 let paths = seed_run_with_node(&tmp, run_id);
1812 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
1813 assert!(read_n0001(&paths).pending_merge.is_some());
1814
1815 let mut status = event(run_id);
1816 status.seq = 4;
1817 status.kind = "node.status".into();
1818 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1819 status.data = serde_json::json!({ "status": "failed" });
1820 apply_event(&paths, &status).unwrap();
1821 let n = read_n0001(&paths);
1822 assert_eq!(n.status, Status::Failed);
1823 assert!(
1824 n.pending_merge.is_none(),
1825 "terminal status clears the transaction"
1826 );
1827 }
1828
1829 #[test]
1833 fn supervisor_attached_sets_supervisor_pid() {
1834 let tmp = TempDir::new().unwrap();
1835 let run_id = "01jxsnap000000000000000000";
1836 let paths = seed_run_with_node(&tmp, run_id);
1837 assert_eq!(read_n0001(&paths).supervisor_pid, None);
1838
1839 let mut ev = event(run_id);
1840 ev.seq = 3;
1841 ev.kind = "supervisor.attached".into();
1842 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1843 ev.data = serde_json::json!({ "pid": 47820 });
1844 apply_event(&paths, &ev).expect("supervisor.attached applies");
1845 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
1846 }
1847
1848 #[test]
1851 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
1852 let tmp = TempDir::new().unwrap();
1853 let run_id = "01jxsnap000000000000000000";
1854 let paths = seed_run_with_node(&tmp, run_id);
1855
1856 let mut ev = event(run_id);
1857 ev.kind = "supervisor.attached".into();
1858 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1859
1860 ev.seq = 3;
1861 ev.data = serde_json::json!({ "pid": 100 });
1862 apply_event(&paths, &ev).expect("first attach applies");
1863 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
1864
1865 ev.seq = 4;
1867 ev.data = serde_json::json!({ "pid": 200 });
1868 apply_event(&paths, &ev).expect("second attach applies");
1869 let after_second = read_n0001(&paths);
1870 assert_eq!(after_second.supervisor_pid, Some(200));
1871
1872 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
1875 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
1876 apply_event(&paths, &ev).expect("replay applies as no-op");
1877 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
1878 }
1879
1880 #[test]
1883 fn supervisor_cursor_advanced_sets_report_cursor() {
1884 let tmp = TempDir::new().unwrap();
1885 let run_id = "01jxsnap000000000000000000";
1886 let paths = seed_run_with_node(&tmp, run_id);
1887 let child = "02jxsnap000000000000000000";
1888
1889 let mut ev = event(run_id);
1890 ev.seq = 3;
1891 ev.kind = "supervisor.cursor_advanced".into();
1892 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1893 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
1894 apply_event(&paths, &ev).expect("cursor_advanced applies");
1895
1896 let n = read_n0001(&paths);
1897 assert_eq!(
1898 n.last_processed_report_seq_by_child.get(child),
1899 Some(&Value::from(7u64))
1900 );
1901 }
1902
1903 #[test]
1908 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
1909 let tmp = TempDir::new().unwrap();
1910 let run_id = "01jxsnap000000000000000000";
1911 let paths = seed_run_with_node(&tmp, run_id);
1912 let child_a = "02jxsnap000000000000000000";
1913 let child_b = "03jxsnap000000000000000000";
1914
1915 let mut ev = event(run_id);
1916 ev.kind = "supervisor.cursor_advanced".into();
1917 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1918
1919 ev.seq = 3;
1920 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
1921 apply_event(&paths, &ev).expect("seq 5 applies");
1922
1923 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
1925 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
1926
1927 ev.seq = 4;
1929 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
1930 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
1931 assert!(ops.is_empty(), "older seq must plan no ops");
1932 apply_event(&paths, &ev).expect("older seq applies as no-op");
1933 assert_eq!(
1934 read_n0001(&paths)
1935 .last_processed_report_seq_by_child
1936 .get(child_a),
1937 Some(&Value::from(5u64))
1938 );
1939
1940 ev.seq = 5;
1942 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
1943 apply_event(&paths, &ev).expect("higher seq applies");
1944 ev.seq = 6;
1945 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
1946 apply_event(&paths, &ev).expect("second child applies");
1947
1948 let n = read_n0001(&paths);
1949 assert_eq!(
1950 n.last_processed_report_seq_by_child.get(child_a),
1951 Some(&Value::from(9u64))
1952 );
1953 assert_eq!(
1954 n.last_processed_report_seq_by_child.get(child_b),
1955 Some(&Value::from(1u64))
1956 );
1957 }
1958
1959 #[test]
1962 fn supervisor_state_events_reject_malformed_payloads() {
1963 let tmp = TempDir::new().unwrap();
1964 let run_id = "01jxsnap000000000000000000";
1965 let paths = seed_run_with_node(&tmp, run_id);
1966 let nid = Some(NodeId::parse_str("n-0001").unwrap());
1967
1968 let mut ev = event(run_id);
1970 ev.seq = 3;
1971 ev.kind = "supervisor.attached".into();
1972 ev.node_id = nid.clone();
1973 ev.data = serde_json::json!({});
1974 assert!(matches!(
1975 reduce_event_to_ops(&paths, &ev),
1976 Err(Error::CorruptEventLog { .. })
1977 ));
1978
1979 ev.node_id = None;
1981 ev.data = serde_json::json!({ "pid": 1 });
1982 assert!(matches!(
1983 reduce_event_to_ops(&paths, &ev),
1984 Err(Error::CorruptEventLog { .. })
1985 ));
1986
1987 let mut ev2 = event(run_id);
1989 ev2.seq = 4;
1990 ev2.kind = "supervisor.cursor_advanced".into();
1991 ev2.node_id = nid.clone();
1992 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
1993 assert!(matches!(
1994 reduce_event_to_ops(&paths, &ev2),
1995 Err(Error::CorruptEventLog { .. })
1996 ));
1997
1998 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2000 assert!(matches!(
2001 reduce_event_to_ops(&paths, &ev2),
2002 Err(Error::CorruptEventLog { .. })
2003 ));
2004 }
2005
2006 #[test]
2013 fn removed_or_garbage_kind_in_created_events_is_rejected() {
2014 let tmp = TempDir::new().unwrap();
2015 let run_id = "01jxsnap000000000000000000";
2016 let rid = RunId::parse_str(run_id).unwrap();
2017 let dir = crate::run_dir(tmp.path(), &rid);
2018 std::fs::create_dir_all(&dir).unwrap();
2019 let paths = RunPaths::new(dir, run_id).unwrap();
2020
2021 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2022 let mut ev = event(run_id);
2023 ev.kind = "run.created".into();
2024 ev.node_id = None;
2025 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2026 assert!(
2027 matches!(
2028 reduce_event_to_ops(&paths, &ev),
2029 Err(Error::CorruptEventLog { .. })
2030 ),
2031 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2032 );
2033 }
2034
2035 let mut ok = event(run_id);
2038 ok.kind = "run.created".into();
2039 ok.node_id = None;
2040 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2041 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2042 }
2043
2044 #[cfg(unix)]
2053 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2054 use std::os::unix::fs::MetadataExt;
2055 let mut consider = vec![paths.manifest()];
2056 for dir in [paths.nodes_dir()] {
2057 if let Ok(rd) = std::fs::read_dir(&dir) {
2058 for ent in rd.flatten() {
2059 let p = ent.path();
2060 if p.extension().and_then(|s| s.to_str()) == Some("json") {
2061 consider.push(p);
2062 }
2063 }
2064 }
2065 }
2066 let mut map = std::collections::BTreeMap::new();
2067 for p in consider {
2068 if let Ok(md) = std::fs::symlink_metadata(&p) {
2069 if md.file_type().is_file() {
2070 map.insert(p, md.ino());
2071 }
2072 }
2073 }
2074 map
2075 }
2076
2077 #[cfg(unix)]
2086 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2087 use std::collections::BTreeSet;
2088 let before = projection_inodes(paths);
2089 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2090 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2091 .into_iter()
2092 .collect();
2093 apply_event(paths, ev)
2094 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2095 let after = projection_inodes(paths);
2096 let touched: BTreeSet<PathBuf> = after
2097 .iter()
2098 .filter(|(p, ino)| before.get(*p) != Some(*ino))
2099 .map(|(p, _)| p.clone())
2100 .collect();
2101 assert_eq!(
2102 planned, touched,
2103 "kind={}: plan_projections must name exactly the files apply_event writes",
2104 ev.kind
2105 );
2106 if expect_writes {
2107 assert!(
2108 !touched.is_empty(),
2109 "kind={}: expected this event to write at least one projection",
2110 ev.kind
2111 );
2112 }
2113 }
2114
2115 #[cfg(unix)]
2122 #[test]
2123 fn plan_projections_matches_apply_for_every_kind() {
2124 let tmp = TempDir::new().unwrap();
2125 let run_id = "01jxsnap000000000000000000";
2126 let rid = RunId::parse_str(run_id).unwrap();
2127 let dir = crate::run_dir(tmp.path(), &rid);
2128 std::fs::create_dir_all(&dir).unwrap();
2129 let paths = RunPaths::new(dir, run_id).unwrap();
2130 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2131 let child = "02jxsnap000000000000000000";
2132
2133 let mut next_seq = 0u64;
2135 let mut at = |kind: &str, node_id, data| {
2136 next_seq += 1;
2137 Event {
2138 ts: Utc::now(),
2139 seq: next_seq,
2140 kind: kind.into(),
2141 run_id: rid.clone(),
2142 node_id,
2143 idempotency_key: None,
2144 data,
2145 }
2146 };
2147
2148 assert_plan_matches_apply(
2150 &paths,
2151 &at(
2152 "run.created",
2153 None,
2154 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2155 ),
2156 true,
2157 );
2158 assert_plan_matches_apply(
2160 &paths,
2161 &at(
2162 "run.status",
2163 None,
2164 serde_json::json!({ "status": "running" }),
2165 ),
2166 true,
2167 );
2168 assert_plan_matches_apply(
2170 &paths,
2171 &at(
2172 "node.created",
2173 nid(),
2174 serde_json::json!({ "kind": "spinoff" }),
2175 ),
2176 true,
2177 );
2178 assert_plan_matches_apply(
2180 &paths,
2181 &at(
2182 "node.status",
2183 nid(),
2184 serde_json::json!({ "status": "running" }),
2185 ),
2186 true,
2187 );
2188 assert_plan_matches_apply(
2190 &paths,
2191 &at(
2192 "supervisor.attached",
2193 nid(),
2194 serde_json::json!({ "pid": 4242 }),
2195 ),
2196 true,
2197 );
2198 assert_plan_matches_apply(
2200 &paths,
2201 &at(
2202 "supervisor.cursor_advanced",
2203 nid(),
2204 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2205 ),
2206 true,
2207 );
2208 assert_plan_matches_apply(
2210 &paths,
2211 &at(
2212 "child.spawned",
2213 nid(),
2214 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2215 ),
2216 true,
2217 );
2218 assert_plan_matches_apply(
2220 &paths,
2221 &at("node.report", nid(), serde_json::json!({ "success": true })),
2222 true,
2223 );
2224 assert_plan_matches_apply(
2227 &paths,
2228 &at(
2229 "node.status",
2230 nid(),
2231 serde_json::json!({ "status": "failed" }),
2232 ),
2233 false,
2234 );
2235 for kind in [
2237 "supervisor.exited",
2238 "orchestrator.decision",
2239 "discuss.critical",
2240 "cleanup.window_missing",
2241 ] {
2242 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
2243 }
2244 }
2245
2246 #[test]
2250 fn worker_exited_records_clean_exit_without_transitioning_status() {
2251 let tmp = TempDir::new().unwrap();
2252 let run_id = "01jxsnap000000000000000000";
2253 let paths = bootstrap_retry_node(&tmp, run_id);
2254
2255 let mut ev = event(run_id);
2256 ev.seq = 3;
2257 ev.kind = "worker.exited".into();
2258 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2259 ev.data = serde_json::json!({ "exit_code": 0 });
2260 apply_event(&paths, &ev).expect("worker.exited applies");
2261
2262 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2263 .unwrap()
2264 .unwrap();
2265 let exit = n.worker_exit.expect("worker_exit recorded");
2266 assert_eq!(exit.code, Some(0));
2267 assert_eq!(exit.signal, None);
2268 assert!(exit.is_clean());
2269 assert_eq!(
2270 n.status,
2271 Status::Pending,
2272 "the exit fact never transitions status"
2273 );
2274 }
2275
2276 #[test]
2280 fn worker_exited_records_signal_and_is_first_write_wins() {
2281 let tmp = TempDir::new().unwrap();
2282 let run_id = "01jxsnap000000000000000000";
2283 let paths = bootstrap_retry_node(&tmp, run_id);
2284 let nid = NodeId::parse_str("n-0001").unwrap();
2285
2286 let mut ev = event(run_id);
2287 ev.seq = 3;
2288 ev.kind = "worker.exited".into();
2289 ev.node_id = Some(nid.clone());
2290 ev.data = serde_json::json!({ "signal": 9 });
2291 apply_event(&paths, &ev).expect("worker.exited applies");
2292
2293 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2294 let exit = n.worker_exit.expect("worker_exit recorded");
2295 assert_eq!(exit.signal, Some(9));
2296 assert!(exit.is_failure());
2297
2298 let mut dup = event(run_id);
2301 dup.seq = 4;
2302 dup.kind = "worker.exited".into();
2303 dup.node_id = Some(nid.clone());
2304 dup.data = serde_json::json!({ "exit_code": 0 });
2305 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
2306 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
2307 assert_eq!(
2308 n2.worker_exit.unwrap().signal,
2309 Some(9),
2310 "first-write-wins: the replayed exit must not overwrite the recorded fact"
2311 );
2312 }
2313
2314 #[test]
2318 fn worker_exited_without_code_or_signal_is_corrupt() {
2319 let tmp = TempDir::new().unwrap();
2320 let run_id = "01jxsnap000000000000000000";
2321 let paths = bootstrap_retry_node(&tmp, run_id);
2322
2323 let mut ev = event(run_id);
2324 ev.seq = 3;
2325 ev.kind = "worker.exited".into();
2326 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2327 ev.data = serde_json::json!({});
2328 match reduce_event_to_ops(&paths, &ev) {
2329 Err(Error::CorruptEventLog { .. }) => {}
2330 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
2331 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2332 }
2333
2334 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
2337 match reduce_event_to_ops(&paths, &ev) {
2338 Err(Error::CorruptEventLog { .. }) => {}
2339 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
2340 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2341 }
2342 }
2343
2344 #[test]
2350 fn node_death_observed_records_first_death_first_write_wins() {
2351 let tmp = TempDir::new().unwrap();
2352 let run_id = "01jxsnap000000000000000000";
2353 let paths = bootstrap_retry_node(&tmp, run_id);
2354 let nid = NodeId::parse_str("n-0001").unwrap();
2355
2356 let mut ev = event(run_id);
2357 ev.seq = 3;
2358 ev.kind = "node.death_observed".into();
2359 ev.node_id = Some(nid.clone());
2360 ev.data = serde_json::json!({});
2361 apply_event(&paths, &ev).expect("node.death_observed applies");
2362 let first = read_node_opt(&paths, &nid)
2363 .unwrap()
2364 .unwrap()
2365 .first_death_at
2366 .expect("first_death_at recorded");
2367 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
2368
2369 let mut later = event(run_id);
2371 later.seq = 4;
2372 later.kind = "node.death_observed".into();
2373 later.node_id = Some(nid.clone());
2374 later.ts = ev.ts + chrono::Duration::seconds(30);
2375 later.data = serde_json::json!({});
2376 apply_event(&paths, &later).expect("re-observation applies as no-op");
2377 assert_eq!(
2378 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2379 Some(first),
2380 "first-write-wins: a re-observation must not reset the anchor"
2381 );
2382 }
2383
2384 #[test]
2389 fn node_death_observed_noop_when_worker_exit_present() {
2390 let tmp = TempDir::new().unwrap();
2391 let run_id = "01jxsnap000000000000000000";
2392 let paths = bootstrap_retry_node(&tmp, run_id);
2393 let nid = NodeId::parse_str("n-0001").unwrap();
2394
2395 let mut exit = event(run_id);
2397 exit.seq = 3;
2398 exit.kind = "worker.exited".into();
2399 exit.node_id = Some(nid.clone());
2400 exit.data = serde_json::json!({ "exit_code": 0 });
2401 apply_event(&paths, &exit).unwrap();
2402
2403 let mut death = event(run_id);
2405 death.seq = 4;
2406 death.kind = "node.death_observed".into();
2407 death.node_id = Some(nid.clone());
2408 death.data = serde_json::json!({});
2409 apply_event(&paths, &death).expect("applies as no-op");
2410 assert_eq!(
2411 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2412 None,
2413 "a told worker.exited makes the crash backstop moot; no anchor recorded"
2414 );
2415 }
2416
2417 #[test]
2422 fn node_retry_clears_worker_exit() {
2423 let tmp = TempDir::new().unwrap();
2424 let run_id = "01jxsnap000000000000000000";
2425 let paths = bootstrap_retry_node(&tmp, run_id);
2426 let nid = NodeId::parse_str("n-0001").unwrap();
2427
2428 let mut exit = event(run_id);
2430 exit.seq = 3;
2431 exit.kind = "worker.exited".into();
2432 exit.node_id = Some(nid.clone());
2433 exit.data = serde_json::json!({ "exit_code": 7 });
2434 apply_event(&paths, &exit).unwrap();
2435 assert!(read_node_opt(&paths, &nid)
2436 .unwrap()
2437 .unwrap()
2438 .worker_exit
2439 .is_some());
2440
2441 let mut retry = event(run_id);
2443 retry.seq = 4;
2444 retry.kind = "node.retry".into();
2445 retry.node_id = Some(nid.clone());
2446 retry.data = serde_json::json!({
2447 "attempt": 1,
2448 "reason": "agent-died",
2449 "branch": "wt/foo",
2450 "worktree_path": "/tmp/new-wt",
2451 "agent_pid": 222,
2452 });
2453 apply_event(&paths, &retry).unwrap();
2454
2455 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2456 assert!(
2457 n.worker_exit.is_none(),
2458 "node.retry must clear the previous attempt's worker_exit"
2459 );
2460 assert_eq!(
2461 n.status,
2462 Status::Pending,
2463 "retry returns the node to Pending"
2464 );
2465 }
2466}