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 "node.awaiting_input" => reduce_node_awaiting_input(paths, ev),
301 "node.input_resolved" => reduce_node_input_resolved(paths, ev),
302 KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
303 KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
304 "child.spawned" => reduce_child_spawned(paths, ev),
305 "supervisor.attached" => reduce_supervisor_attached(paths, ev),
306 "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
307 "supervisor.exited" => Ok(vec![]),
308 "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
316 "run.notified" | "run.awaiting_input_notified" => Ok(vec![]),
323 "cleanup.window_missing"
348 | "cleanup.worktree_missing"
349 | "cleanup.branch_remove_failed"
350 | "cleanup.branch_preserved"
351 | "cleanup.session_killed"
352 | "cleanup.session_retained" => Ok(vec![]),
353 "supervisor.child_id_quarantined" => Ok(vec![]),
363 _ => Ok(vec![]),
364 }
365}
366
367fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
373 match op {
374 ProjectionOp::Manifest(_) => paths.manifest(),
375 ProjectionOp::Node(n) => paths.node(&n.node_id),
376 }
377}
378
379pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
401 let ops = reduce_event_to_ops(paths, event)?;
402 Ok(ops.iter().map(|op| op_path(paths, op)).collect())
403}
404
405pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
421 let ops = reduce_event_to_ops(paths, ev)?;
422 commit_ops(paths, ops)
423}
424
425#[cfg(test)]
434pub(crate) fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
435 reduce_event_to_ops(paths, ev).map(|_| ())
436}
437
438fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
442 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
443 path: events_path.to_path_buf(),
444 reason: format!(
445 "event seq={} kind={} missing top-level `node_id`",
446 ev.seq, ev.kind
447 ),
448 })
449}
450
451fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
452 if let Some(existing) = read_manifest_opt(paths)? {
456 if existing.run_id != ev.run_id {
457 return Err(Error::CorruptEventLog {
458 path: paths.manifest(),
459 reason: format!(
460 "run.created run_id={} conflicts with existing manifest run_id={}",
461 ev.run_id, existing.run_id
462 ),
463 });
464 }
465 return Ok(vec![]);
466 }
467 let events_path = paths.events();
468 let d = &ev.data;
469 let kind =
470 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
471 path: events_path.clone(),
472 reason: "run.created missing/invalid `kind`".into(),
473 })?;
474 let lifecycle: Lifecycle = serde_json::from_value(
475 d.get("lifecycle").cloned().unwrap_or(Value::Null),
476 )
477 .map_err(|_| Error::CorruptEventLog {
478 path: events_path.clone(),
479 reason: "run.created missing/invalid `lifecycle`".into(),
480 })?;
481 let title = want_str(&events_path, ev, d, "title")?.to_string();
482 let m = Manifest {
483 schema_version: STATE_SCHEMA_VERSION,
484 applied_seq: 0,
487 run_id: paths.run_id.clone(),
489 kind,
490 lifecycle,
491 title,
492 status: Status::Pending,
493 created_at: ev.ts,
494 updated_at: ev.ts,
495 source_repo: d
496 .get("source_repo")
497 .and_then(Value::as_str)
498 .map(str::to_string),
499 source_branch: d
500 .get("source_branch")
501 .and_then(Value::as_str)
502 .map(str::to_string),
503 worktree_root: d
504 .get("worktree_root")
505 .and_then(Value::as_str)
506 .map(str::to_string),
507 managed_tmux_session: d
508 .get("managed_tmux_session")
509 .and_then(Value::as_str)
510 .map(str::to_string),
511 notify_cmd: d
512 .get("notify_cmd")
513 .and_then(Value::as_str)
514 .map(str::to_string),
515 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
516 node_count: 0,
517 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
518 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
519 };
520 Ok(vec![ProjectionOp::Manifest(m)])
521}
522
523fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
524 let mut m = match read_manifest_opt(paths)? {
525 Some(m) => m,
526 None => return Ok(vec![]),
527 };
528 let new_status = require_status(ev, paths.events())?;
529 if m.status.is_terminal() {
532 trace_terminal_noop(ev, m.status, new_status);
533 return Ok(vec![]);
534 }
535 if m.status == new_status {
536 return Ok(vec![]);
537 }
538 m.status = new_status;
539 m.updated_at = ev.ts;
540 Ok(vec![ProjectionOp::Manifest(m)])
541}
542
543fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
554 let nonempty = |key| {
555 d.get(key)
556 .and_then(Value::as_str)
557 .map(str::trim)
558 .filter(|s| !s.is_empty())
559 .map(str::to_string)
560 };
561 let session = nonempty("tmux_session")?;
562 let window_id = nonempty("tmux_window_id")?;
563 Some(TmuxIdentity {
564 socket: nonempty("tmux_socket"),
565 session,
566 window_id,
567 pane_id: nonempty("tmux_pane_id"),
570 })
571}
572
573fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
574 let events_path = paths.events();
575 let node_id = require_envelope_node_id(&events_path, ev)?;
578 let is_default_node = node_id.as_str() == "n-0001";
579 if read_node_opt(paths, &node_id)?.is_some() {
581 return Ok(vec![]);
582 }
583 let d = &ev.data;
584 let kind =
585 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
586 path: events_path.clone(),
587 reason: format!(
588 "event seq={} kind=node.created missing/invalid `kind`",
589 ev.seq
590 ),
591 })?;
592 let n = Node {
593 schema_version: STATE_SCHEMA_VERSION,
594 node_id,
595 run_id: paths.run_id.clone(),
597 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
598 kind,
599 status: Status::Pending,
600 task: d.get("task").and_then(Value::as_str).map(str::to_string),
601 worktree_path: d
602 .get("worktree_path")
603 .and_then(Value::as_str)
604 .map(str::to_string),
605 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
606 base_sha: d
607 .get("base_sha")
608 .and_then(Value::as_str)
609 .filter(|s| !s.is_empty())
610 .map(str::to_string),
611 tmux_window: d
612 .get("tmux_window")
613 .and_then(Value::as_str)
614 .map(str::to_string),
615 tmux_identity: tmux_identity_from_data(d),
616 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
617 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
618 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
619 children: Vec::new(),
620 started_at: Some(ev.ts),
621 updated_at: ev.ts,
622 last_report: None,
623 last_processed_report_seq_by_child: serde_json::Map::default(),
624 retry_attempts: 0,
625 worker_exit: None,
626 pending_merge: None,
627 first_death_at: None,
628 awaiting_input: None,
629 };
630 let mut ops = vec![ProjectionOp::Node(n)];
631 if let Some(mut m) = read_manifest_opt(paths)? {
632 if is_default_node && m.source_branch.is_none() {
637 m.source_branch = d
638 .get("source_branch")
639 .and_then(Value::as_str)
640 .filter(|branch| !branch.is_empty())
641 .map(str::to_string);
642 }
643 m.updated_at = ev.ts;
648 ops.push(ProjectionOp::Manifest(m));
649 }
650 Ok(ops)
651}
652
653fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
672 let events_path = paths.events();
673 let node_id = require_envelope_node_id(&events_path, ev)?;
674 let mut n = match read_node_opt(paths, &node_id)? {
675 Some(n) => n,
676 None => return Ok(vec![]),
677 };
678 if n.status.is_terminal() {
681 tracing::debug!(
682 target: "octl_core::reducer",
683 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
684 "no-op: node.retry against terminal node"
685 );
686 return Ok(vec![]);
687 }
688 let d = &ev.data;
689 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
692 n.base_sha = d
693 .get("base_sha")
694 .and_then(Value::as_str)
695 .filter(|s| !s.is_empty())
696 .map(str::to_string);
697 n.worktree_path = d
698 .get("worktree_path")
699 .and_then(Value::as_str)
700 .map(str::to_string);
701 n.tmux_window = d
702 .get("tmux_window")
703 .and_then(Value::as_str)
704 .map(str::to_string);
705 n.tmux_identity = tmux_identity_from_data(d);
706 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
707 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
708 n.status = Status::Pending;
709 n.started_at = Some(ev.ts);
710 n.updated_at = ev.ts;
711 n.last_report = None;
712 n.pending_merge = None;
718 n.worker_exit = None;
723 n.first_death_at = None;
728 n.awaiting_input = None;
731 n.retry_attempts = d
739 .get("attempt")
740 .and_then(Value::as_u64)
741 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
742 Ok(vec![ProjectionOp::Node(n)])
743}
744
745fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
746 let events_path = paths.events();
747 let node_id = require_envelope_node_id(&events_path, ev)?;
748 let mut n = match read_node_opt(paths, &node_id)? {
749 Some(n) => n,
750 None => return Ok(vec![]),
751 };
752 let new_status = require_status(ev, events_path)?;
753 if n.status.is_terminal() {
756 trace_terminal_noop(ev, n.status, new_status);
757 return Ok(vec![]);
758 }
759 if n.status == new_status {
760 return Ok(vec![]);
761 }
762 n.status = new_status;
763 if new_status.is_terminal() {
769 n.pending_merge = None;
770 n.awaiting_input = None;
771 }
772 n.updated_at = ev.ts;
773 Ok(vec![ProjectionOp::Node(n)])
774}
775
776fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
777 let events_path = paths.events();
778 let node_id = require_envelope_node_id(&events_path, ev)?;
779 let mut n = match read_node_opt(paths, &node_id)? {
780 Some(n) => n,
781 None => return Ok(vec![]),
782 };
783 if n.status.is_terminal() {
795 if matches!(n.status, Status::Failed | Status::Done)
834 && report_is_confirmed_explicit_merge(&ev.data)
835 {
836 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
837 return Ok(vec![]);
838 }
839 tracing::info!(
840 target: "octl_core::reducer",
841 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
842 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
843 );
844 n.last_report = Some(ev.data.clone());
845 n.status = Status::Done;
849 n.pending_merge = None;
853 n.awaiting_input = None;
854 n.updated_at = ev.ts;
855 return Ok(vec![ProjectionOp::Node(n)]);
856 }
857 tracing::debug!(
858 target: "octl_core::reducer",
859 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
860 "no-op: node.report against terminal node"
861 );
862 return Ok(vec![]);
863 }
864 let new_status = report_terminal_status(&events_path, ev)?;
872 n.last_report = Some(ev.data.clone());
873 n.status = new_status;
874 n.awaiting_input = None;
878 n.pending_merge = None;
883 n.updated_at = ev.ts;
884 Ok(vec![ProjectionOp::Node(n)])
885}
886
887fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
902 let events_path = paths.events();
903 let node_id = require_envelope_node_id(&events_path, ev)?;
904 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
905 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
906 match (code, signal) {
911 (Some(_), None) | (None, Some(_)) => {}
912 _ => {
913 return Err(Error::CorruptEventLog {
914 path: events_path,
915 reason: format!(
916 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
917 ev.seq
918 ),
919 });
920 }
921 }
922 let mut n = match read_node_opt(paths, &node_id)? {
923 Some(n) => n,
924 None => return Ok(vec![]),
931 };
932 if n.worker_exit.is_some() {
936 return Ok(vec![]);
937 }
938 n.worker_exit = Some(WorkerExit {
939 code,
940 signal,
941 at: ev.ts,
942 });
943 n.awaiting_input = None;
947 n.updated_at = ev.ts;
948 Ok(vec![ProjectionOp::Node(n)])
949}
950
951fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
964 let events_path = paths.events();
965 let node_id = require_envelope_node_id(&events_path, ev)?;
966 let mut n = match read_node_opt(paths, &node_id)? {
967 Some(n) => n,
968 None => return Ok(vec![]),
969 };
970 if n.first_death_at.is_some()
976 || n.status.is_terminal()
977 || n.worker_exit.is_some()
978 || n.last_report.is_some()
979 || n.pending_merge.is_some()
980 {
981 return Ok(vec![]);
982 }
983 n.first_death_at = Some(ev.ts);
984 n.updated_at = ev.ts;
985 Ok(vec![ProjectionOp::Node(n)])
986}
987
988fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
997 const MAX_ITEMS: usize = 8;
998 const MAX_TOPIC_CHARS: usize = 512;
999 const MAX_OPTIONS: usize = 16;
1000 const MAX_OPTION_CHARS: usize = 256;
1001
1002 let events_path = paths.events();
1003 let node_id = require_envelope_node_id(&events_path, ev)?;
1004 let items = ev
1007 .data
1008 .get("discussion_items")
1009 .and_then(Value::as_array)
1010 .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1011 .ok_or_else(|| Error::CorruptEventLog {
1012 path: events_path.clone(),
1013 reason: format!(
1014 "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1015 ev.seq
1016 ),
1017 })?;
1018 for (index, item) in items.iter().enumerate() {
1019 let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1020 path: events_path.clone(),
1021 reason: format!(
1022 "event seq={} discussion_items[{index}] must be an object",
1023 ev.seq
1024 ),
1025 })?;
1026 let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1027 let default = obj
1028 .get("recommended_default")
1029 .and_then(Value::as_str)
1030 .unwrap_or("");
1031 let options = obj.get("options").and_then(Value::as_array);
1032 let options_valid = options.is_some_and(|values| {
1033 !values.is_empty()
1034 && values.len() <= MAX_OPTIONS
1035 && values.iter().all(|v| {
1036 v.as_str().is_some_and(|s| {
1037 !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1038 })
1039 })
1040 && values.iter().any(|v| v.as_str() == Some(default))
1041 });
1042 if topic.trim().is_empty()
1043 || topic.chars().count() > MAX_TOPIC_CHARS
1044 || default.trim().is_empty()
1045 || !options_valid
1046 {
1047 return Err(Error::CorruptEventLog {
1048 path: events_path.clone(),
1049 reason: format!(
1050 "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1051 ev.seq
1052 ),
1053 });
1054 }
1055 }
1056
1057 let mut n = match read_node_opt(paths, &node_id)? {
1058 Some(n) => n,
1059 None => return Ok(vec![]),
1060 };
1061 if n.status.is_terminal() || n.worker_exit.is_some() {
1062 return Ok(vec![]);
1063 }
1064 if let Some(open) = n.awaiting_input.as_mut() {
1065 if open.discussion_items.len() + items.len() > MAX_ITEMS {
1068 return Err(Error::CorruptEventLog {
1069 path: events_path,
1070 reason: format!(
1071 "event seq={} would exceed {MAX_ITEMS} open discussion items",
1072 ev.seq
1073 ),
1074 });
1075 }
1076 open.discussion_items.extend(items.iter().cloned());
1077 } else {
1078 n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1079 opened_at: ev.ts,
1080 event_seq: ev.seq,
1081 discussion_items: items.clone(),
1082 }));
1083 }
1084 n.updated_at = ev.ts;
1085 Ok(vec![ProjectionOp::Node(n)])
1086}
1087
1088fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1092 let events_path = paths.events();
1093 let node_id = require_envelope_node_id(&events_path, ev)?;
1094 let seq = ev
1097 .data
1098 .get("event_seq")
1099 .and_then(Value::as_u64)
1100 .ok_or_else(|| Error::CorruptEventLog {
1101 path: events_path.clone(),
1102 reason: format!(
1103 "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1104 ev.seq
1105 ),
1106 })?;
1107 let mut n = match read_node_opt(paths, &node_id)? {
1108 Some(n) => n,
1109 None => return Ok(vec![]),
1110 };
1111 let Some(open) = n.awaiting_input.as_ref() else {
1112 return Ok(vec![]);
1113 };
1114 if seq != open.event_seq {
1115 return Ok(vec![]);
1116 }
1117 n.awaiting_input = None;
1118 n.updated_at = ev.ts;
1119 Ok(vec![ProjectionOp::Node(n)])
1120}
1121
1122pub const KIND_MERGE_STARTED: &str = "merge.started";
1126
1127pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1132
1133fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1147 let events_path = paths.events();
1148 let node_id = require_envelope_node_id(&events_path, ev)?;
1149 let txn: MergeTxn =
1150 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1151 path: events_path.clone(),
1152 reason: format!(
1153 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1154 ev.seq
1155 ),
1156 })?;
1157 let mut n = match read_node_opt(paths, &node_id)? {
1158 Some(n) => n,
1159 None => return Ok(vec![]),
1160 };
1161 if n.status.is_terminal() {
1165 return Ok(vec![]);
1166 }
1167 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1170 return Ok(vec![]);
1171 }
1172 n.pending_merge = Some(Box::new(txn));
1173 n.updated_at = ev.ts;
1174 Ok(vec![ProjectionOp::Node(n)])
1175}
1176
1177fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1189 let events_path = paths.events();
1190 let node_id = require_envelope_node_id(&events_path, ev)?;
1191 let op_id = ev
1192 .data
1193 .get("op_id")
1194 .and_then(Value::as_str)
1195 .ok_or_else(|| Error::CorruptEventLog {
1196 path: events_path.clone(),
1197 reason: format!(
1198 "event seq={} kind=merge.aborted is missing string `op_id`",
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 match n.pending_merge.as_ref() {
1209 Some(t) if t.op_id == op_id => {}
1210 _ => return Ok(vec![]),
1211 }
1212 n.pending_merge = None;
1213 n.updated_at = ev.ts;
1214 Ok(vec![ProjectionOp::Node(n)])
1215}
1216
1217fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1225 if current == incoming {
1226 tracing::debug!(
1227 target: "octl_core::reducer",
1228 seq = ev.seq, kind = %ev.kind, status = ?current,
1229 "no-op: status re-applied to terminal target"
1230 );
1231 } else {
1232 tracing::warn!(
1233 target: "octl_core::reducer",
1234 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1235 "no-op: ignored conflicting transition from terminal target"
1236 );
1237 }
1238}
1239
1240fn report_is_confirmed_explicit_merge(data: &Value) -> bool {
1257 ReportOrigin::report_is_confirmed_merge(data)
1258}
1259
1260fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1269 let corrupt = |reason: String| Error::CorruptEventLog {
1270 path: events_path.to_path_buf(),
1271 reason,
1272 };
1273 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1274 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1275 if cancelled {
1276 if success == Some(true) {
1277 return Err(corrupt(format!(
1278 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1279 ev.seq
1280 )));
1281 }
1282 Ok(Status::Cancelled)
1283 } else {
1284 match success {
1285 Some(true) => Ok(Status::Done),
1286 Some(false) => Ok(Status::Failed),
1287 None => Err(corrupt(format!(
1288 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1289 ev.seq
1290 ))),
1291 }
1292 }
1293}
1294
1295fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1296 let events_path = paths.events();
1299 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1300 path: events_path.clone(),
1301 reason: format!(
1302 "event seq={} kind=child.spawned missing parent `node_id`",
1303 ev.seq
1304 ),
1305 })?;
1306 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1307 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1308 let child_node_id = NodeId::parse_str(
1309 ev.data
1310 .get("child_node_id")
1311 .and_then(Value::as_str)
1312 .unwrap_or("n-0001"),
1313 )
1314 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1315 let mut n = match read_node_opt(paths, &parent_node_id)? {
1316 Some(n) => n,
1317 None => return Ok(vec![]),
1318 };
1319 let new_ref = ChildRef {
1320 run_id: child_run_id,
1321 node_id: child_node_id,
1322 };
1323 if n.children.iter().any(|c| c == &new_ref) {
1324 return Ok(vec![]);
1327 }
1328 n.children.push(new_ref);
1329 n.updated_at = ev.ts;
1330 Ok(vec![ProjectionOp::Node(n)])
1331}
1332
1333fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1344 let events_path = paths.events();
1345 let node_id = require_envelope_node_id(&events_path, ev)?;
1346 let raw = ev
1347 .data
1348 .get("pid")
1349 .and_then(Value::as_i64)
1350 .ok_or_else(|| Error::CorruptEventLog {
1351 path: events_path.clone(),
1352 reason: format!(
1353 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1354 ev.seq
1355 ),
1356 })?;
1357 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1358 path: events_path.clone(),
1359 reason: format!(
1360 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1361 ev.seq
1362 ),
1363 })?;
1364 let mut n = match read_node_opt(paths, &node_id)? {
1365 Some(n) => n,
1366 None => return Ok(vec![]),
1367 };
1368 if n.supervisor_pid == Some(pid) {
1369 return Ok(vec![]);
1370 }
1371 n.supervisor_pid = Some(pid);
1372 n.updated_at = ev.ts;
1373 Ok(vec![ProjectionOp::Node(n)])
1374}
1375
1376fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1387 let events_path = paths.events();
1388 let node_id = require_envelope_node_id(&events_path, ev)?;
1389 let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1390 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1394 let report_seq = ev
1395 .data
1396 .get("report_seq")
1397 .and_then(Value::as_u64)
1398 .ok_or_else(|| Error::CorruptEventLog {
1399 path: events_path.clone(),
1400 reason: format!(
1401 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1402 ev.seq
1403 ),
1404 })?;
1405 let mut n = match read_node_opt(paths, &node_id)? {
1406 Some(n) => n,
1407 None => return Ok(vec![]),
1408 };
1409 if let Some(prev) = n
1410 .last_processed_report_seq_by_child
1411 .get(child_run_id)
1412 .and_then(Value::as_u64)
1413 {
1414 if report_seq <= prev {
1415 return Ok(vec![]);
1416 }
1417 }
1418 n.last_processed_report_seq_by_child
1419 .insert(child_run_id.to_string(), Value::from(report_seq));
1420 n.updated_at = ev.ts;
1421 Ok(vec![ProjectionOp::Node(n)])
1422}
1423
1424#[cfg(test)]
1425mod tests {
1426 use super::*;
1427 use crate::schema::Event;
1428 use chrono::Utc;
1429 use tempfile::TempDir;
1430
1431 fn event(run_id: &str) -> Event {
1432 Event {
1433 ts: Utc::now(),
1434 seq: 1,
1435 kind: "run.status".into(),
1436 run_id: RunId::parse_str(run_id).unwrap(),
1437 node_id: None,
1438 idempotency_key: None,
1439 data: serde_json::json!({ "status": "running" }),
1440 }
1441 }
1442
1443 #[test]
1444 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1445 let tmp = TempDir::new().unwrap();
1449 let run_id = "01jxsnap000000000000000000";
1450 let rid = RunId::parse_str(run_id).unwrap();
1451 let dir = crate::run_dir(tmp.path(), &rid);
1452 std::fs::create_dir_all(&dir).unwrap();
1453 let paths = RunPaths::new(dir, run_id).unwrap();
1454
1455 let mut created = event(run_id);
1458 created.kind = "run.created".into();
1459 created.data = serde_json::json!({
1460 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1461 });
1462 apply_event(&paths, &created).expect("run.created applies");
1463 let manifest_before = std::fs::read(paths.manifest()).unwrap();
1464
1465 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1466 let mut ev = event(run_id);
1467 ev.seq = seq;
1468 ev.kind = kind.into();
1469 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1471 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1472 assert!(ops.is_empty(), "{kind} must plan no projection ops");
1473 apply_event(&paths, &ev).expect("audit kind applies as no-op");
1475 }
1476
1477 assert_eq!(
1479 std::fs::read(paths.manifest()).unwrap(),
1480 manifest_before,
1481 "audit events must not mutate the manifest"
1482 );
1483 assert!(!paths.nodes_dir().exists(), "no node projection created");
1484 }
1485
1486 #[test]
1487 fn run_created_folds_harness_when_present_and_defaults_none() {
1488 let tmp = TempDir::new().unwrap();
1489
1490 let run_id = "01jxhrnsaa0000000000000001";
1492 let rid = RunId::parse_str(run_id).unwrap();
1493 let dir = crate::run_dir(tmp.path(), &rid);
1494 std::fs::create_dir_all(&dir).unwrap();
1495 let paths = RunPaths::new(dir, run_id).unwrap();
1496 let mut created = event(run_id);
1497 created.kind = "run.created".into();
1498 created.data = serde_json::json!({
1499 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1500 "harness": "pi", "harness_source": "flag",
1501 });
1502 apply_event(&paths, &created).expect("run.created applies");
1503 let m = read_manifest_opt(&paths).unwrap().unwrap();
1504 assert_eq!(m.harness.as_deref(), Some("pi"));
1505
1506 let run_id2 = "01jxhrnsaa0000000000000002";
1508 let rid2 = RunId::parse_str(run_id2).unwrap();
1509 let dir2 = crate::run_dir(tmp.path(), &rid2);
1510 std::fs::create_dir_all(&dir2).unwrap();
1511 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1512 let mut created2 = event(run_id2);
1513 created2.kind = "run.created".into();
1514 created2.data =
1515 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1516 apply_event(&paths2, &created2).expect("run.created applies");
1517 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1518 assert_eq!(m2.harness, None);
1519 }
1520
1521 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1524 let rid = RunId::parse_str(run_id).unwrap();
1525 let dir = crate::run_dir(tmp.path(), &rid);
1526 std::fs::create_dir_all(&dir).unwrap();
1527 let paths = RunPaths::new(dir, run_id).unwrap();
1528 let mut created = event(run_id);
1529 created.kind = "run.created".into();
1530 created.data =
1531 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1532 apply_event(&paths, &created).expect("run.created applies");
1533 let mut node = event(run_id);
1534 node.seq = 2;
1535 node.kind = "node.created".into();
1536 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1537 node.data = serde_json::json!({
1538 "kind": "spinoff",
1539 "branch": "wt/foo",
1540 "worktree_path": "/tmp/old-wt",
1541 "agent_pid": 111,
1542 });
1543 apply_event(&paths, &node).expect("node.created applies");
1544 paths
1545 }
1546
1547 #[test]
1551 fn node_retry_rewires_node_and_increments_attempts() {
1552 let tmp = TempDir::new().unwrap();
1553 let run_id = "01jxsnap000000000000000000";
1554 let paths = bootstrap_retry_node(&tmp, run_id);
1555
1556 let mut retry = event(run_id);
1557 retry.seq = 3;
1558 retry.kind = "node.retry".into();
1559 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1560 retry.data = serde_json::json!({
1561 "attempt": 1,
1562 "reason": "agent-died",
1563 "branch": "wt/foo-r1",
1564 "base_sha": "a".repeat(40),
1565 "worktree_path": "/tmp/new-wt",
1566 "agent_pid": 222,
1567 "tmux_session": "s",
1568 "tmux_window_id": "@9",
1569 });
1570 apply_event(&paths, &retry).expect("node.retry applies");
1571
1572 let n = read_n0001(&paths);
1573 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1574 assert_eq!(
1575 n.branch.as_deref(),
1576 Some("wt/foo-r1"),
1577 "rewired to new branch"
1578 );
1579 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1580 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1581 assert_eq!(n.status, Status::Pending, "node returns to pending");
1582 assert!(n.last_report.is_none());
1583 assert_eq!(
1584 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1585 Some("@9"),
1586 "rewired tmux identity"
1587 );
1588
1589 let mut retry2 = event(run_id);
1591 retry2.seq = 4;
1592 retry2.kind = "node.retry".into();
1593 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1594 retry2.data = serde_json::json!({
1595 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1596 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1597 });
1598 apply_event(&paths, &retry2).expect("node.retry applies");
1599 assert_eq!(read_n0001(&paths).retry_attempts, 2);
1600 }
1601
1602 #[test]
1606 fn node_retry_against_terminal_node_is_noop() {
1607 let tmp = TempDir::new().unwrap();
1608 let run_id = "01jxsnap000000000000000000";
1609 let paths = bootstrap_retry_node(&tmp, run_id);
1610
1611 let mut report = event(run_id);
1613 report.seq = 3;
1614 report.kind = "node.report".into();
1615 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1616 report.data = serde_json::json!({ "success": true });
1617 apply_event(&paths, &report).expect("node.report applies");
1618 assert_eq!(read_n0001(&paths).status, Status::Done);
1619
1620 let mut retry = event(run_id);
1621 retry.seq = 4;
1622 retry.kind = "node.retry".into();
1623 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1624 retry.data = serde_json::json!({
1625 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1626 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1627 });
1628 apply_event(&paths, &retry).expect("node.retry applies as no-op");
1629
1630 let n = read_n0001(&paths);
1631 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
1632 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
1633 assert_eq!(n.agent_pid, Some(111), "not rewired");
1634 }
1635
1636 #[test]
1637 fn apply_event_rejects_event_from_a_different_run() {
1638 let tmp = TempDir::new().unwrap();
1639 let run_id = "01jxsnap000000000000000000";
1640 let rid = RunId::parse_str(run_id).unwrap();
1641 let dir = crate::run_dir(tmp.path(), &rid);
1642 std::fs::create_dir_all(&dir).unwrap();
1643 let paths = RunPaths::new(dir, run_id).unwrap();
1644
1645 let foreign = event("02jxsnap000000000000000000");
1647 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
1648 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
1649
1650 let mine = event(run_id);
1653 apply_event(&paths, &mine).expect("matching run_id must be accepted");
1654 }
1655
1656 #[test]
1657 fn tmux_identity_from_data_reads_qualified_fields() {
1658 let d = serde_json::json!({
1659 "tmux_socket": "/private/tmp/tmux-501/default",
1660 "tmux_session": "octl",
1661 "tmux_window_id": "@42",
1662 });
1663 let id = tmux_identity_from_data(&d).expect("qualified identity");
1664 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
1665 assert_eq!(id.session, "octl");
1666 assert_eq!(id.window_id, "@42");
1667 assert_eq!(id.pane_id, None);
1669
1670 let d2 = serde_json::json!({
1672 "tmux_socket": null,
1673 "tmux_session": "octl",
1674 "tmux_window_id": "@7",
1675 });
1676 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
1677 assert_eq!(id2.socket, None);
1678 assert_eq!(id2.window_id, "@7");
1679
1680 let d3 = serde_json::json!({
1682 "tmux_session": "octl",
1683 "tmux_window_id": "@42",
1684 "tmux_pane_id": "%7",
1685 });
1686 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
1687 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
1688 assert_eq!(id3.capture_target(), "%7");
1689
1690 let d4 = serde_json::json!({
1693 "tmux_session": "octl",
1694 "tmux_window_id": "@42",
1695 "tmux_pane_id": null,
1696 });
1697 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
1698 assert_eq!(id4.pane_id, None);
1699 assert_eq!(id4.capture_target(), "@42");
1700 }
1701
1702 #[test]
1703 fn tmux_identity_from_data_back_compat_is_none() {
1704 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
1706 assert!(tmux_identity_from_data(&legacy).is_none());
1707 let partial = serde_json::json!({ "tmux_window_id": "@42" });
1709 assert!(tmux_identity_from_data(&partial).is_none());
1710 }
1711
1712 #[test]
1715 fn node_created_populates_tmux_identity() {
1716 let tmp = TempDir::new().unwrap();
1717 let run_id = "01jxsnap000000000000000000";
1718 let rid = RunId::parse_str(run_id).unwrap();
1719 let dir = crate::run_dir(tmp.path(), &rid);
1720 std::fs::create_dir_all(&dir).unwrap();
1721 let paths = RunPaths::new(dir, run_id).unwrap();
1722
1723 let mut ev = event(run_id);
1724 ev.seq = 2;
1725 ev.kind = "node.created".into();
1726 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1727 ev.data = serde_json::json!({
1728 "kind": "spinoff",
1729 "tmux_window": "🚀 wt/x",
1730 "tmux_socket": "/private/tmp/tmux-501/default",
1731 "tmux_session": "octl",
1732 "tmux_window_id": "@42",
1733 });
1734 apply_event(&paths, &ev).expect("node.created applies");
1735 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
1736 .unwrap()
1737 .unwrap();
1738 let id = n.tmux_identity.expect("qualified identity recorded");
1739 assert_eq!(id.session, "octl");
1740 assert_eq!(id.window_id, "@42");
1741 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
1742
1743 let run2 = "02jxsnap000000000000000000";
1745 let rid2 = RunId::parse_str(run2).unwrap();
1746 let dir2 = crate::run_dir(tmp.path(), &rid2);
1747 std::fs::create_dir_all(&dir2).unwrap();
1748 let paths2 = RunPaths::new(dir2, run2).unwrap();
1749 let mut ev2 = event(run2);
1750 ev2.seq = 2;
1751 ev2.kind = "node.created".into();
1752 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1753 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
1754 apply_event(&paths2, &ev2).expect("legacy node.created applies");
1755 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
1756 .unwrap()
1757 .unwrap();
1758 assert!(n2.tmux_identity.is_none());
1759 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
1760 }
1761
1762 #[test]
1763 fn node_materialization_populates_missing_manifest_source_branch() {
1764 let tmp = TempDir::new().unwrap();
1765 let run_id = "01jxsnap000000000000000001";
1766 let rid = RunId::parse_str(run_id).unwrap();
1767 let dir = crate::run_dir(tmp.path(), &rid);
1768 std::fs::create_dir_all(&dir).unwrap();
1769 let paths = RunPaths::new(dir, run_id).unwrap();
1770
1771 let mut created = event(run_id);
1772 created.kind = "run.created".into();
1773 created.data = serde_json::json!({
1774 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1775 });
1776 apply_event(&paths, &created).unwrap();
1777 assert!(read_manifest_opt(&paths)
1778 .unwrap()
1779 .unwrap()
1780 .source_branch
1781 .is_none());
1782
1783 let mut node = event(run_id);
1784 node.seq = 2;
1785 node.kind = "node.created".into();
1786 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1787 node.data = serde_json::json!({
1788 "kind": "spinoff",
1789 "source_branch": "main",
1790 "worktree_path": "/tmp/wt/pending"
1791 });
1792 apply_event(&paths, &node).unwrap();
1793
1794 let manifest = read_manifest_opt(&paths).unwrap().unwrap();
1795 assert_eq!(manifest.status, Status::Pending);
1796 assert_eq!(manifest.source_branch.as_deref(), Some("main"));
1797 let node = read_n0001(&paths);
1798 assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
1799 }
1800
1801 #[test]
1802 fn node_materialization_preserves_explicit_manifest_source_branch() {
1803 let tmp = TempDir::new().unwrap();
1804 let run_id = "01jxsnap000000000000000002";
1805 let rid = RunId::parse_str(run_id).unwrap();
1806 let dir = crate::run_dir(tmp.path(), &rid);
1807 std::fs::create_dir_all(&dir).unwrap();
1808 let paths = RunPaths::new(dir, run_id).unwrap();
1809
1810 let mut created = event(run_id);
1811 created.kind = "run.created".into();
1812 created.data = serde_json::json!({
1813 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1814 "source_branch": "release"
1815 });
1816 apply_event(&paths, &created).unwrap();
1817
1818 let mut node = event(run_id);
1819 node.seq = 2;
1820 node.kind = "node.created".into();
1821 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1822 node.data = serde_json::json!({
1823 "kind": "spinoff", "source_branch": "main"
1824 });
1825 apply_event(&paths, &node).unwrap();
1826
1827 assert_eq!(
1828 read_manifest_opt(&paths)
1829 .unwrap()
1830 .unwrap()
1831 .source_branch
1832 .as_deref(),
1833 Some("release")
1834 );
1835 }
1836
1837 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1840 let rid = RunId::parse_str(run_id).unwrap();
1841 let dir = crate::run_dir(tmp.path(), &rid);
1842 std::fs::create_dir_all(&dir).unwrap();
1843 let paths = RunPaths::new(dir, run_id).unwrap();
1844
1845 let mut created = event(run_id);
1846 created.kind = "run.created".into();
1847 created.data = serde_json::json!({
1848 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1849 });
1850 apply_event(&paths, &created).expect("run.created applies");
1851
1852 let mut node = event(run_id);
1853 node.seq = 2;
1854 node.kind = "node.created".into();
1855 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1856 node.data = serde_json::json!({ "kind": "spinoff" });
1857 apply_event(&paths, &node).expect("node.created applies");
1858 paths
1859 }
1860
1861 fn read_n0001(paths: &RunPaths) -> Node {
1862 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
1863 .unwrap()
1864 .unwrap()
1865 }
1866
1867 #[test]
1868 fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
1869 let tmp = TempDir::new().unwrap();
1870 let run_id = "01jxwd0000000000000000000w";
1871 let paths = seed_run_with_node(&tmp, run_id);
1872 let nid = Some(NodeId::parse_str("n-0001").unwrap());
1873 let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
1874 let mut open = event(run_id);
1875 open.seq = 3;
1876 open.ts = opened_at;
1877 open.kind = "node.awaiting_input".into();
1878 open.node_id = nid.clone();
1879 open.data = serde_json::json!({ "discussion_items": [{
1880 "topic": "Which scope?",
1881 "options": ["small", "large"],
1882 "recommended_default": "small"
1883 }] });
1884 apply_event(&paths, &open).unwrap();
1885 let first = read_n0001(&paths).awaiting_input.unwrap();
1886 assert_eq!(first.opened_at, opened_at);
1887 assert_eq!(first.event_seq, 3);
1888
1889 let mut duplicate = open.clone();
1891 duplicate.seq = 4;
1892 duplicate.ts = opened_at + chrono::Duration::hours(1);
1893 apply_event(&paths, &duplicate).unwrap();
1894 let still_first = read_n0001(&paths).awaiting_input.unwrap();
1895 assert_eq!(still_first.opened_at, opened_at);
1896 assert_eq!(still_first.event_seq, 3);
1897
1898 let mut stale = event(run_id);
1900 stale.seq = 5;
1901 stale.kind = "node.input_resolved".into();
1902 stale.node_id = nid.clone();
1903 stale.data = serde_json::json!({ "event_seq": 2 });
1904 apply_event(&paths, &stale).unwrap();
1905 assert!(read_n0001(&paths).awaiting_input.is_some());
1906
1907 let mut resolved = stale;
1908 resolved.seq = 6;
1909 resolved.data = serde_json::json!({ "event_seq": 3 });
1910 apply_event(&paths, &resolved).unwrap();
1911 assert!(read_n0001(&paths).awaiting_input.is_none());
1912 }
1913
1914 #[test]
1915 fn awaiting_input_rejects_missing_default_without_mutating_projection() {
1916 let tmp = TempDir::new().unwrap();
1917 let run_id = "01jxwd0000000000000000000x";
1918 let paths = seed_run_with_node(&tmp, run_id);
1919 let mut open = event(run_id);
1920 open.seq = 3;
1921 open.kind = "node.awaiting_input".into();
1922 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1923 open.data = serde_json::json!({ "discussion_items": [{
1924 "topic": "Which scope?", "options": ["small", "large"]
1925 }] });
1926 assert!(reduce_event_to_ops(&paths, &open).is_err());
1927 assert!(read_n0001(&paths).awaiting_input.is_none());
1928 }
1929
1930 #[test]
1931 fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
1932 let tmp = TempDir::new().unwrap();
1933 let run_id = "01jxwd0000000000000000000y";
1934 let paths = seed_run_with_node(&tmp, run_id);
1935 let nid = Some(NodeId::parse_str("n-0001").unwrap());
1936
1937 let mut open = event(run_id);
1938 open.seq = 3;
1939 open.kind = "node.awaiting_input".into();
1940 open.node_id = nid.clone();
1941 open.data = serde_json::json!({ "discussion_items": [{
1942 "topic": "Which scope?", "options": ["small", "large"],
1943 "recommended_default": "small"
1944 }] });
1945 apply_event(&paths, &open).unwrap();
1946
1947 let mut malformed_duplicate = open.clone();
1948 malformed_duplicate.seq = 4;
1949 malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
1950 assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
1951
1952 let mut exited = event(run_id);
1953 exited.seq = 5;
1954 exited.kind = "worker.exited".into();
1955 exited.node_id = nid;
1956 exited.data = serde_json::json!({ "exit_code": 0 });
1957 apply_event(&paths, &exited).unwrap();
1958 let node = read_n0001(&paths);
1959 assert!(node.awaiting_input.is_none());
1960 assert!(node.worker_exit.is_some());
1961
1962 let mut delayed_open = open;
1963 delayed_open.seq = 6;
1964 assert!(reduce_event_to_ops(&paths, &delayed_open)
1965 .unwrap()
1966 .is_empty());
1967 }
1968
1969 #[test]
1970 fn input_resolved_requires_generation_even_when_nothing_is_open() {
1971 let tmp = TempDir::new().unwrap();
1972 let run_id = "01jxwd0000000000000000000z";
1973 let paths = seed_run_with_node(&tmp, run_id);
1974 let mut resolved = event(run_id);
1975 resolved.seq = 3;
1976 resolved.kind = "node.input_resolved".into();
1977 resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1978 resolved.data = serde_json::json!({});
1979 assert!(reduce_event_to_ops(&paths, &resolved).is_err());
1980 }
1981
1982 #[test]
1983 fn awaiting_input_default_must_be_one_of_options() {
1984 let tmp = TempDir::new().unwrap();
1985 let run_id = "01jxwd00000000000000000010";
1986 let paths = seed_run_with_node(&tmp, run_id);
1987 let mut open = event(run_id);
1988 open.seq = 3;
1989 open.kind = "node.awaiting_input".into();
1990 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1991 open.data = serde_json::json!({ "discussion_items": [{
1992 "topic": "Which scope?", "options": ["small", "large"],
1993 "recommended_default": "other"
1994 }] });
1995 assert!(reduce_event_to_ops(&paths, &open).is_err());
1996 }
1997
1998 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
1999 let mut ev = event(run_id);
2000 ev.seq = seq;
2001 ev.kind = KIND_MERGE_STARTED.into();
2002 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2003 ev.data = serde_json::json!({
2004 "op_id": op_id,
2005 "source_branch": "main",
2006 "worker_branch": "wt/worker",
2007 "expected_source_oid": expected,
2008 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2009 "base_sha": null,
2010 "driver_pid": 4242,
2011 "driver_pid_start_secs": null,
2012 "started_at": "2026-08-15T00:00:00Z",
2013 });
2014 ev
2015 }
2016
2017 #[test]
2020 fn merge_started_records_pending_transaction() {
2021 let tmp = TempDir::new().unwrap();
2022 let run_id = "01jxsnap000000000000000000";
2023 let paths = seed_run_with_node(&tmp, run_id);
2024
2025 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2026 let n = read_n0001(&paths);
2027 assert_eq!(
2028 n.status,
2029 Status::Pending,
2030 "recording a merge is not terminal"
2031 );
2032 let txn = n.pending_merge.expect("transaction recorded");
2033 assert_eq!(txn.op_id, "op-1");
2034 assert_eq!(txn.expected_source_oid, "aaa");
2035 }
2036
2037 #[test]
2040 fn merge_aborted_clears_matching_transaction_only() {
2041 let tmp = TempDir::new().unwrap();
2042 let run_id = "01jxsnap000000000000000000";
2043 let paths = seed_run_with_node(&tmp, run_id);
2044 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2045
2046 let mut stale = event(run_id);
2048 stale.seq = 4;
2049 stale.kind = KIND_MERGE_ABORTED.into();
2050 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2051 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2052 apply_event(&paths, &stale).unwrap();
2053 assert!(
2054 read_n0001(&paths).pending_merge.is_some(),
2055 "stale abort is a no-op"
2056 );
2057
2058 let mut abort = event(run_id);
2060 abort.seq = 5;
2061 abort.kind = KIND_MERGE_ABORTED.into();
2062 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2063 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2064 apply_event(&paths, &abort).unwrap();
2065 let n = read_n0001(&paths);
2066 assert!(
2067 n.pending_merge.is_none(),
2068 "matching abort clears the transaction"
2069 );
2070 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2071 }
2072
2073 #[test]
2076 fn terminal_report_clears_pending_merge() {
2077 let tmp = TempDir::new().unwrap();
2078 let run_id = "01jxsnap000000000000000000";
2079 let paths = seed_run_with_node(&tmp, run_id);
2080 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2081
2082 let mut report = event(run_id);
2083 report.seq = 4;
2084 report.kind = "node.report".into();
2085 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2086 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2087 apply_event(&paths, &report).unwrap();
2088 let n = read_n0001(&paths);
2089 assert_eq!(n.status, Status::Done);
2090 assert!(
2091 n.pending_merge.is_none(),
2092 "completed merge clears the transaction"
2093 );
2094 }
2095
2096 #[test]
2105 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2106 let tmp = TempDir::new().unwrap();
2107
2108 let drive = |run_id: &str, report_data: Value| -> Status {
2111 let paths = seed_run_with_node(&tmp, run_id);
2112 let mut fail = event(run_id);
2114 fail.seq = 3;
2115 fail.kind = "node.status".into();
2116 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2117 fail.data = serde_json::json!({ "status": "failed" });
2118 apply_event(&paths, &fail).unwrap();
2119 assert_eq!(read_n0001(&paths).status, Status::Failed);
2120 let mut report = event(run_id);
2122 report.seq = 4;
2123 report.kind = "node.report".into();
2124 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2125 report.data = report_data;
2126 apply_event(&paths, &report).unwrap();
2127 read_n0001(&paths).status
2128 };
2129
2130 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2132 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2133 assert_eq!(
2134 drive("01jxsnap000000000000000001", agent_forged),
2135 Status::Failed,
2136 "an Agent-origin report with a forged via must not be adopted"
2137 );
2138
2139 let malformed = serde_json::json!({
2141 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2142 });
2143 assert_eq!(
2144 drive("01jxsnap000000000000000002", malformed),
2145 Status::Failed,
2146 "a malformed origin must not re-unlock the legacy via adoption path"
2147 );
2148
2149 let mut run_merge = serde_json::json!({ "success": true });
2151 crate::ReportOrigin::RunMerge {
2152 op_id: Some("op-1".into()),
2153 worker_oid: Some("cafebabe".into()),
2154 }
2155 .stamp(&mut run_merge);
2156 assert_eq!(
2157 drive("01jxsnap000000000000000003", run_merge),
2158 Status::Done,
2159 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2160 );
2161
2162 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2165 assert_eq!(
2166 drive("01jxsnap000000000000000004", legacy),
2167 Status::Done,
2168 "a legacy via-only report (no origin field) is still adopted"
2169 );
2170 }
2171
2172 #[test]
2176 fn terminal_node_status_clears_pending_merge() {
2177 let tmp = TempDir::new().unwrap();
2178 let run_id = "01jxsnap000000000000000000";
2179 let paths = seed_run_with_node(&tmp, run_id);
2180 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2181 assert!(read_n0001(&paths).pending_merge.is_some());
2182
2183 let mut status = event(run_id);
2184 status.seq = 4;
2185 status.kind = "node.status".into();
2186 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2187 status.data = serde_json::json!({ "status": "failed" });
2188 apply_event(&paths, &status).unwrap();
2189 let n = read_n0001(&paths);
2190 assert_eq!(n.status, Status::Failed);
2191 assert!(
2192 n.pending_merge.is_none(),
2193 "terminal status clears the transaction"
2194 );
2195 }
2196
2197 #[test]
2201 fn supervisor_attached_sets_supervisor_pid() {
2202 let tmp = TempDir::new().unwrap();
2203 let run_id = "01jxsnap000000000000000000";
2204 let paths = seed_run_with_node(&tmp, run_id);
2205 assert_eq!(read_n0001(&paths).supervisor_pid, None);
2206
2207 let mut ev = event(run_id);
2208 ev.seq = 3;
2209 ev.kind = "supervisor.attached".into();
2210 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2211 ev.data = serde_json::json!({ "pid": 47820 });
2212 apply_event(&paths, &ev).expect("supervisor.attached applies");
2213 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2214 }
2215
2216 #[test]
2219 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2220 let tmp = TempDir::new().unwrap();
2221 let run_id = "01jxsnap000000000000000000";
2222 let paths = seed_run_with_node(&tmp, run_id);
2223
2224 let mut ev = event(run_id);
2225 ev.kind = "supervisor.attached".into();
2226 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2227
2228 ev.seq = 3;
2229 ev.data = serde_json::json!({ "pid": 100 });
2230 apply_event(&paths, &ev).expect("first attach applies");
2231 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2232
2233 ev.seq = 4;
2235 ev.data = serde_json::json!({ "pid": 200 });
2236 apply_event(&paths, &ev).expect("second attach applies");
2237 let after_second = read_n0001(&paths);
2238 assert_eq!(after_second.supervisor_pid, Some(200));
2239
2240 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2243 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2244 apply_event(&paths, &ev).expect("replay applies as no-op");
2245 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2246 }
2247
2248 #[test]
2251 fn supervisor_cursor_advanced_sets_report_cursor() {
2252 let tmp = TempDir::new().unwrap();
2253 let run_id = "01jxsnap000000000000000000";
2254 let paths = seed_run_with_node(&tmp, run_id);
2255 let child = "02jxsnap000000000000000000";
2256
2257 let mut ev = event(run_id);
2258 ev.seq = 3;
2259 ev.kind = "supervisor.cursor_advanced".into();
2260 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2261 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2262 apply_event(&paths, &ev).expect("cursor_advanced applies");
2263
2264 let n = read_n0001(&paths);
2265 assert_eq!(
2266 n.last_processed_report_seq_by_child.get(child),
2267 Some(&Value::from(7u64))
2268 );
2269 }
2270
2271 #[test]
2276 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2277 let tmp = TempDir::new().unwrap();
2278 let run_id = "01jxsnap000000000000000000";
2279 let paths = seed_run_with_node(&tmp, run_id);
2280 let child_a = "02jxsnap000000000000000000";
2281 let child_b = "03jxsnap000000000000000000";
2282
2283 let mut ev = event(run_id);
2284 ev.kind = "supervisor.cursor_advanced".into();
2285 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2286
2287 ev.seq = 3;
2288 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2289 apply_event(&paths, &ev).expect("seq 5 applies");
2290
2291 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2293 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2294
2295 ev.seq = 4;
2297 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2298 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2299 assert!(ops.is_empty(), "older seq must plan no ops");
2300 apply_event(&paths, &ev).expect("older seq applies as no-op");
2301 assert_eq!(
2302 read_n0001(&paths)
2303 .last_processed_report_seq_by_child
2304 .get(child_a),
2305 Some(&Value::from(5u64))
2306 );
2307
2308 ev.seq = 5;
2310 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2311 apply_event(&paths, &ev).expect("higher seq applies");
2312 ev.seq = 6;
2313 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2314 apply_event(&paths, &ev).expect("second child applies");
2315
2316 let n = read_n0001(&paths);
2317 assert_eq!(
2318 n.last_processed_report_seq_by_child.get(child_a),
2319 Some(&Value::from(9u64))
2320 );
2321 assert_eq!(
2322 n.last_processed_report_seq_by_child.get(child_b),
2323 Some(&Value::from(1u64))
2324 );
2325 }
2326
2327 #[test]
2330 fn supervisor_state_events_reject_malformed_payloads() {
2331 let tmp = TempDir::new().unwrap();
2332 let run_id = "01jxsnap000000000000000000";
2333 let paths = seed_run_with_node(&tmp, run_id);
2334 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2335
2336 let mut ev = event(run_id);
2338 ev.seq = 3;
2339 ev.kind = "supervisor.attached".into();
2340 ev.node_id = nid.clone();
2341 ev.data = serde_json::json!({});
2342 assert!(matches!(
2343 reduce_event_to_ops(&paths, &ev),
2344 Err(Error::CorruptEventLog { .. })
2345 ));
2346
2347 ev.node_id = None;
2349 ev.data = serde_json::json!({ "pid": 1 });
2350 assert!(matches!(
2351 reduce_event_to_ops(&paths, &ev),
2352 Err(Error::CorruptEventLog { .. })
2353 ));
2354
2355 let mut ev2 = event(run_id);
2357 ev2.seq = 4;
2358 ev2.kind = "supervisor.cursor_advanced".into();
2359 ev2.node_id = nid.clone();
2360 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
2361 assert!(matches!(
2362 reduce_event_to_ops(&paths, &ev2),
2363 Err(Error::CorruptEventLog { .. })
2364 ));
2365
2366 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2368 assert!(matches!(
2369 reduce_event_to_ops(&paths, &ev2),
2370 Err(Error::CorruptEventLog { .. })
2371 ));
2372 }
2373
2374 #[test]
2381 fn removed_or_garbage_kind_in_created_events_is_rejected() {
2382 let tmp = TempDir::new().unwrap();
2383 let run_id = "01jxsnap000000000000000000";
2384 let rid = RunId::parse_str(run_id).unwrap();
2385 let dir = crate::run_dir(tmp.path(), &rid);
2386 std::fs::create_dir_all(&dir).unwrap();
2387 let paths = RunPaths::new(dir, run_id).unwrap();
2388
2389 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2390 let mut ev = event(run_id);
2391 ev.kind = "run.created".into();
2392 ev.node_id = None;
2393 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2394 assert!(
2395 matches!(
2396 reduce_event_to_ops(&paths, &ev),
2397 Err(Error::CorruptEventLog { .. })
2398 ),
2399 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2400 );
2401 }
2402
2403 let mut ok = event(run_id);
2406 ok.kind = "run.created".into();
2407 ok.node_id = None;
2408 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2409 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2410 }
2411
2412 #[cfg(unix)]
2421 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2422 use std::os::unix::fs::MetadataExt;
2423 let mut consider = vec![paths.manifest()];
2424 for dir in [paths.nodes_dir()] {
2425 if let Ok(rd) = std::fs::read_dir(&dir) {
2426 for ent in rd.flatten() {
2427 let p = ent.path();
2428 if p.extension().and_then(|s| s.to_str()) == Some("json") {
2429 consider.push(p);
2430 }
2431 }
2432 }
2433 }
2434 let mut map = std::collections::BTreeMap::new();
2435 for p in consider {
2436 if let Ok(md) = std::fs::symlink_metadata(&p) {
2437 if md.file_type().is_file() {
2438 map.insert(p, md.ino());
2439 }
2440 }
2441 }
2442 map
2443 }
2444
2445 #[cfg(unix)]
2454 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2455 use std::collections::BTreeSet;
2456 let before = projection_inodes(paths);
2457 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2458 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2459 .into_iter()
2460 .collect();
2461 apply_event(paths, ev)
2462 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2463 let after = projection_inodes(paths);
2464 let touched: BTreeSet<PathBuf> = after
2465 .iter()
2466 .filter(|(p, ino)| before.get(*p) != Some(*ino))
2467 .map(|(p, _)| p.clone())
2468 .collect();
2469 assert_eq!(
2470 planned, touched,
2471 "kind={}: plan_projections must name exactly the files apply_event writes",
2472 ev.kind
2473 );
2474 if expect_writes {
2475 assert!(
2476 !touched.is_empty(),
2477 "kind={}: expected this event to write at least one projection",
2478 ev.kind
2479 );
2480 }
2481 }
2482
2483 #[cfg(unix)]
2490 #[test]
2491 fn plan_projections_matches_apply_for_every_kind() {
2492 let tmp = TempDir::new().unwrap();
2493 let run_id = "01jxsnap000000000000000000";
2494 let rid = RunId::parse_str(run_id).unwrap();
2495 let dir = crate::run_dir(tmp.path(), &rid);
2496 std::fs::create_dir_all(&dir).unwrap();
2497 let paths = RunPaths::new(dir, run_id).unwrap();
2498 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2499 let child = "02jxsnap000000000000000000";
2500
2501 let mut next_seq = 0u64;
2503 let mut at = |kind: &str, node_id, data| {
2504 next_seq += 1;
2505 Event {
2506 ts: Utc::now(),
2507 seq: next_seq,
2508 kind: kind.into(),
2509 run_id: rid.clone(),
2510 node_id,
2511 idempotency_key: None,
2512 data,
2513 }
2514 };
2515
2516 assert_plan_matches_apply(
2518 &paths,
2519 &at(
2520 "run.created",
2521 None,
2522 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2523 ),
2524 true,
2525 );
2526 assert_plan_matches_apply(
2528 &paths,
2529 &at(
2530 "run.status",
2531 None,
2532 serde_json::json!({ "status": "running" }),
2533 ),
2534 true,
2535 );
2536 assert_plan_matches_apply(
2538 &paths,
2539 &at(
2540 "node.created",
2541 nid(),
2542 serde_json::json!({ "kind": "spinoff" }),
2543 ),
2544 true,
2545 );
2546 assert_plan_matches_apply(
2548 &paths,
2549 &at(
2550 "node.status",
2551 nid(),
2552 serde_json::json!({ "status": "running" }),
2553 ),
2554 true,
2555 );
2556 assert_plan_matches_apply(
2558 &paths,
2559 &at(
2560 "supervisor.attached",
2561 nid(),
2562 serde_json::json!({ "pid": 4242 }),
2563 ),
2564 true,
2565 );
2566 assert_plan_matches_apply(
2568 &paths,
2569 &at(
2570 "supervisor.cursor_advanced",
2571 nid(),
2572 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2573 ),
2574 true,
2575 );
2576 assert_plan_matches_apply(
2578 &paths,
2579 &at(
2580 "child.spawned",
2581 nid(),
2582 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2583 ),
2584 true,
2585 );
2586 assert_plan_matches_apply(
2588 &paths,
2589 &at("node.report", nid(), serde_json::json!({ "success": true })),
2590 true,
2591 );
2592 assert_plan_matches_apply(
2595 &paths,
2596 &at(
2597 "node.status",
2598 nid(),
2599 serde_json::json!({ "status": "failed" }),
2600 ),
2601 false,
2602 );
2603 for kind in [
2605 "supervisor.exited",
2606 "orchestrator.decision",
2607 "discuss.critical",
2608 "cleanup.window_missing",
2609 ] {
2610 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
2611 }
2612 }
2613
2614 #[test]
2618 fn worker_exited_records_clean_exit_without_transitioning_status() {
2619 let tmp = TempDir::new().unwrap();
2620 let run_id = "01jxsnap000000000000000000";
2621 let paths = bootstrap_retry_node(&tmp, run_id);
2622
2623 let mut ev = event(run_id);
2624 ev.seq = 3;
2625 ev.kind = "worker.exited".into();
2626 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2627 ev.data = serde_json::json!({ "exit_code": 0 });
2628 apply_event(&paths, &ev).expect("worker.exited applies");
2629
2630 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2631 .unwrap()
2632 .unwrap();
2633 let exit = n.worker_exit.expect("worker_exit recorded");
2634 assert_eq!(exit.code, Some(0));
2635 assert_eq!(exit.signal, None);
2636 assert!(exit.is_clean());
2637 assert_eq!(
2638 n.status,
2639 Status::Pending,
2640 "the exit fact never transitions status"
2641 );
2642 }
2643
2644 #[test]
2648 fn worker_exited_records_signal_and_is_first_write_wins() {
2649 let tmp = TempDir::new().unwrap();
2650 let run_id = "01jxsnap000000000000000000";
2651 let paths = bootstrap_retry_node(&tmp, run_id);
2652 let nid = NodeId::parse_str("n-0001").unwrap();
2653
2654 let mut ev = event(run_id);
2655 ev.seq = 3;
2656 ev.kind = "worker.exited".into();
2657 ev.node_id = Some(nid.clone());
2658 ev.data = serde_json::json!({ "signal": 9 });
2659 apply_event(&paths, &ev).expect("worker.exited applies");
2660
2661 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2662 let exit = n.worker_exit.expect("worker_exit recorded");
2663 assert_eq!(exit.signal, Some(9));
2664 assert!(exit.is_failure());
2665
2666 let mut dup = event(run_id);
2669 dup.seq = 4;
2670 dup.kind = "worker.exited".into();
2671 dup.node_id = Some(nid.clone());
2672 dup.data = serde_json::json!({ "exit_code": 0 });
2673 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
2674 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
2675 assert_eq!(
2676 n2.worker_exit.unwrap().signal,
2677 Some(9),
2678 "first-write-wins: the replayed exit must not overwrite the recorded fact"
2679 );
2680 }
2681
2682 #[test]
2686 fn worker_exited_without_code_or_signal_is_corrupt() {
2687 let tmp = TempDir::new().unwrap();
2688 let run_id = "01jxsnap000000000000000000";
2689 let paths = bootstrap_retry_node(&tmp, run_id);
2690
2691 let mut ev = event(run_id);
2692 ev.seq = 3;
2693 ev.kind = "worker.exited".into();
2694 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2695 ev.data = serde_json::json!({});
2696 match reduce_event_to_ops(&paths, &ev) {
2697 Err(Error::CorruptEventLog { .. }) => {}
2698 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
2699 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2700 }
2701
2702 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
2705 match reduce_event_to_ops(&paths, &ev) {
2706 Err(Error::CorruptEventLog { .. }) => {}
2707 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
2708 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2709 }
2710 }
2711
2712 #[test]
2718 fn node_death_observed_records_first_death_first_write_wins() {
2719 let tmp = TempDir::new().unwrap();
2720 let run_id = "01jxsnap000000000000000000";
2721 let paths = bootstrap_retry_node(&tmp, run_id);
2722 let nid = NodeId::parse_str("n-0001").unwrap();
2723
2724 let mut ev = event(run_id);
2725 ev.seq = 3;
2726 ev.kind = "node.death_observed".into();
2727 ev.node_id = Some(nid.clone());
2728 ev.data = serde_json::json!({});
2729 apply_event(&paths, &ev).expect("node.death_observed applies");
2730 let first = read_node_opt(&paths, &nid)
2731 .unwrap()
2732 .unwrap()
2733 .first_death_at
2734 .expect("first_death_at recorded");
2735 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
2736
2737 let mut later = event(run_id);
2739 later.seq = 4;
2740 later.kind = "node.death_observed".into();
2741 later.node_id = Some(nid.clone());
2742 later.ts = ev.ts + chrono::Duration::seconds(30);
2743 later.data = serde_json::json!({});
2744 apply_event(&paths, &later).expect("re-observation applies as no-op");
2745 assert_eq!(
2746 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2747 Some(first),
2748 "first-write-wins: a re-observation must not reset the anchor"
2749 );
2750 }
2751
2752 #[test]
2757 fn node_death_observed_noop_when_worker_exit_present() {
2758 let tmp = TempDir::new().unwrap();
2759 let run_id = "01jxsnap000000000000000000";
2760 let paths = bootstrap_retry_node(&tmp, run_id);
2761 let nid = NodeId::parse_str("n-0001").unwrap();
2762
2763 let mut exit = event(run_id);
2765 exit.seq = 3;
2766 exit.kind = "worker.exited".into();
2767 exit.node_id = Some(nid.clone());
2768 exit.data = serde_json::json!({ "exit_code": 0 });
2769 apply_event(&paths, &exit).unwrap();
2770
2771 let mut death = event(run_id);
2773 death.seq = 4;
2774 death.kind = "node.death_observed".into();
2775 death.node_id = Some(nid.clone());
2776 death.data = serde_json::json!({});
2777 apply_event(&paths, &death).expect("applies as no-op");
2778 assert_eq!(
2779 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2780 None,
2781 "a told worker.exited makes the crash backstop moot; no anchor recorded"
2782 );
2783 }
2784
2785 #[test]
2790 fn node_retry_clears_worker_exit() {
2791 let tmp = TempDir::new().unwrap();
2792 let run_id = "01jxsnap000000000000000000";
2793 let paths = bootstrap_retry_node(&tmp, run_id);
2794 let nid = NodeId::parse_str("n-0001").unwrap();
2795
2796 let mut exit = event(run_id);
2798 exit.seq = 3;
2799 exit.kind = "worker.exited".into();
2800 exit.node_id = Some(nid.clone());
2801 exit.data = serde_json::json!({ "exit_code": 7 });
2802 apply_event(&paths, &exit).unwrap();
2803 assert!(read_node_opt(&paths, &nid)
2804 .unwrap()
2805 .unwrap()
2806 .worker_exit
2807 .is_some());
2808
2809 let mut retry = event(run_id);
2811 retry.seq = 4;
2812 retry.kind = "node.retry".into();
2813 retry.node_id = Some(nid.clone());
2814 retry.data = serde_json::json!({
2815 "attempt": 1,
2816 "reason": "agent-died",
2817 "branch": "wt/foo",
2818 "worktree_path": "/tmp/new-wt",
2819 "agent_pid": 222,
2820 });
2821 apply_event(&paths, &retry).unwrap();
2822
2823 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2824 assert!(
2825 n.worker_exit.is_none(),
2826 "node.retry must clear the previous attempt's worker_exit"
2827 );
2828 assert_eq!(
2829 n.status,
2830 Status::Pending,
2831 "retry returns the node to Pending"
2832 );
2833 }
2834}