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
425pub fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
431 reduce_event_to_ops(paths, ev).map(|_| ())
432}
433
434fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
438 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
439 path: events_path.to_path_buf(),
440 reason: format!(
441 "event seq={} kind={} missing top-level `node_id`",
442 ev.seq, ev.kind
443 ),
444 })
445}
446
447fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
448 if let Some(existing) = read_manifest_opt(paths)? {
452 if existing.run_id != ev.run_id {
453 return Err(Error::CorruptEventLog {
454 path: paths.manifest(),
455 reason: format!(
456 "run.created run_id={} conflicts with existing manifest run_id={}",
457 ev.run_id, existing.run_id
458 ),
459 });
460 }
461 return Ok(vec![]);
462 }
463 let events_path = paths.events();
464 let d = &ev.data;
465 let kind =
466 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
467 path: events_path.clone(),
468 reason: "run.created missing/invalid `kind`".into(),
469 })?;
470 let lifecycle: Lifecycle = serde_json::from_value(
471 d.get("lifecycle").cloned().unwrap_or(Value::Null),
472 )
473 .map_err(|_| Error::CorruptEventLog {
474 path: events_path.clone(),
475 reason: "run.created missing/invalid `lifecycle`".into(),
476 })?;
477 let title = want_str(&events_path, ev, d, "title")?.to_string();
478 let agent_selection: Option<crate::schema::AgentSelection> = d
479 .get("agent_selection")
480 .cloned()
481 .map(serde_json::from_value)
482 .transpose()
483 .map_err(|e| Error::CorruptEventLog {
484 path: events_path.clone(),
485 reason: format!("run.created invalid `agent_selection`: {e}"),
486 })?;
487 if let Some(selection) = &agent_selection {
488 selection
489 .validate()
490 .map_err(|reason| Error::CorruptEventLog {
491 path: events_path.clone(),
492 reason: format!("run.created invalid `agent_selection`: {reason}"),
493 })?;
494 }
495 let m = Manifest {
496 schema_version: STATE_SCHEMA_VERSION,
497 applied_seq: 0,
500 run_id: paths.run_id.clone(),
502 kind,
503 lifecycle,
504 title,
505 status: Status::Pending,
506 created_at: ev.ts,
507 updated_at: ev.ts,
508 source_repo: d
509 .get("source_repo")
510 .and_then(Value::as_str)
511 .map(str::to_string),
512 source_branch: d
513 .get("source_branch")
514 .and_then(Value::as_str)
515 .map(str::to_string),
516 worktree_root: d
517 .get("worktree_root")
518 .and_then(Value::as_str)
519 .map(str::to_string),
520 managed_tmux_session: d
521 .get("managed_tmux_session")
522 .and_then(Value::as_str)
523 .map(str::to_string),
524 notify_cmd: d
525 .get("notify_cmd")
526 .and_then(Value::as_str)
527 .map(str::to_string),
528 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
529 agent_selection,
530 node_count: 0,
531 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
532 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
533 };
534 Ok(vec![ProjectionOp::Manifest(m)])
535}
536
537fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
538 let mut m = match read_manifest_opt(paths)? {
539 Some(m) => m,
540 None => return Ok(vec![]),
541 };
542 let new_status = require_status(ev, paths.events())?;
543 if m.status.is_terminal() {
546 trace_terminal_noop(ev, m.status, new_status);
547 return Ok(vec![]);
548 }
549 if m.status == new_status {
550 return Ok(vec![]);
551 }
552 m.status = new_status;
553 m.updated_at = ev.ts;
554 Ok(vec![ProjectionOp::Manifest(m)])
555}
556
557fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
568 let nonempty = |key| {
569 d.get(key)
570 .and_then(Value::as_str)
571 .map(str::trim)
572 .filter(|s| !s.is_empty())
573 .map(str::to_string)
574 };
575 let session = nonempty("tmux_session")?;
576 let window_id = nonempty("tmux_window_id")?;
577 Some(TmuxIdentity {
578 socket: nonempty("tmux_socket"),
579 session,
580 window_id,
581 pane_id: nonempty("tmux_pane_id"),
584 })
585}
586
587fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
588 let events_path = paths.events();
589 let node_id = require_envelope_node_id(&events_path, ev)?;
592 let is_default_node = node_id.as_str() == "n-0001";
593 if read_node_opt(paths, &node_id)?.is_some() {
595 return Ok(vec![]);
596 }
597 let d = &ev.data;
598 let kind =
599 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
600 path: events_path.clone(),
601 reason: format!(
602 "event seq={} kind=node.created missing/invalid `kind`",
603 ev.seq
604 ),
605 })?;
606 let n = Node {
607 schema_version: STATE_SCHEMA_VERSION,
608 node_id,
609 run_id: paths.run_id.clone(),
611 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
612 kind,
613 status: Status::Pending,
614 task: d.get("task").and_then(Value::as_str).map(str::to_string),
615 worktree_path: d
616 .get("worktree_path")
617 .and_then(Value::as_str)
618 .map(str::to_string),
619 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
620 base_sha: d
621 .get("base_sha")
622 .and_then(Value::as_str)
623 .filter(|s| !s.is_empty())
624 .map(str::to_string),
625 tmux_window: d
626 .get("tmux_window")
627 .and_then(Value::as_str)
628 .map(str::to_string),
629 tmux_identity: tmux_identity_from_data(d),
630 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
631 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
632 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
633 children: Vec::new(),
634 started_at: Some(ev.ts),
635 updated_at: ev.ts,
636 last_report: None,
637 last_processed_report_seq_by_child: serde_json::Map::default(),
638 retry_attempts: 0,
639 worker_exit: None,
640 pending_merge: None,
641 first_death_at: None,
642 awaiting_input: None,
643 };
644 let mut ops = vec![ProjectionOp::Node(n)];
645 if let Some(mut m) = read_manifest_opt(paths)? {
646 if is_default_node && m.source_branch.is_none() {
651 m.source_branch = d
652 .get("source_branch")
653 .and_then(Value::as_str)
654 .filter(|branch| !branch.is_empty())
655 .map(str::to_string);
656 }
657 m.updated_at = ev.ts;
662 ops.push(ProjectionOp::Manifest(m));
663 }
664 Ok(ops)
665}
666
667fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
686 let events_path = paths.events();
687 let node_id = require_envelope_node_id(&events_path, ev)?;
688 let mut n = match read_node_opt(paths, &node_id)? {
689 Some(n) => n,
690 None => return Ok(vec![]),
691 };
692 if n.status.is_terminal() {
695 tracing::debug!(
696 target: "taskfleet_core::reducer",
697 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
698 "no-op: node.retry against terminal node"
699 );
700 return Ok(vec![]);
701 }
702 let d = &ev.data;
703 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
706 n.base_sha = d
707 .get("base_sha")
708 .and_then(Value::as_str)
709 .filter(|s| !s.is_empty())
710 .map(str::to_string);
711 n.worktree_path = d
712 .get("worktree_path")
713 .and_then(Value::as_str)
714 .map(str::to_string);
715 n.tmux_window = d
716 .get("tmux_window")
717 .and_then(Value::as_str)
718 .map(str::to_string);
719 n.tmux_identity = tmux_identity_from_data(d);
720 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
721 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
722 n.status = Status::Pending;
723 n.started_at = Some(ev.ts);
724 n.updated_at = ev.ts;
725 n.last_report = None;
726 n.pending_merge = None;
732 n.worker_exit = None;
737 n.first_death_at = None;
742 n.awaiting_input = None;
745 n.retry_attempts = d
753 .get("attempt")
754 .and_then(Value::as_u64)
755 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
756 Ok(vec![ProjectionOp::Node(n)])
757}
758
759fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
760 let events_path = paths.events();
761 let node_id = require_envelope_node_id(&events_path, ev)?;
762 let mut n = match read_node_opt(paths, &node_id)? {
763 Some(n) => n,
764 None => return Ok(vec![]),
765 };
766 let new_status = require_status(ev, events_path)?;
767 if n.status.is_terminal() {
770 trace_terminal_noop(ev, n.status, new_status);
771 return Ok(vec![]);
772 }
773 if n.status == new_status {
774 return Ok(vec![]);
775 }
776 n.status = new_status;
777 if new_status.is_terminal() {
783 n.pending_merge = None;
784 n.awaiting_input = None;
785 }
786 n.updated_at = ev.ts;
787 Ok(vec![ProjectionOp::Node(n)])
788}
789
790fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
791 let events_path = paths.events();
792 let node_id = require_envelope_node_id(&events_path, ev)?;
793 let mut n = match read_node_opt(paths, &node_id)? {
794 Some(n) => n,
795 None => return Ok(vec![]),
796 };
797 if n.status.is_terminal() {
809 if matches!(n.status, Status::Failed | Status::Done)
848 && report_is_confirmed_explicit_merge(&ev.data)
849 {
850 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
851 return Ok(vec![]);
852 }
853 tracing::info!(
854 target: "taskfleet_core::reducer",
855 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
856 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
857 );
858 n.last_report = Some(ev.data.clone());
859 n.status = Status::Done;
863 n.pending_merge = None;
867 n.awaiting_input = None;
868 n.updated_at = ev.ts;
869 return Ok(vec![ProjectionOp::Node(n)]);
870 }
871 tracing::debug!(
872 target: "taskfleet_core::reducer",
873 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
874 "no-op: node.report against terminal node"
875 );
876 return Ok(vec![]);
877 }
878 let new_status = report_terminal_status(&events_path, ev)?;
886 n.last_report = Some(ev.data.clone());
887 n.status = new_status;
888 n.awaiting_input = None;
892 n.pending_merge = None;
897 n.updated_at = ev.ts;
898 Ok(vec![ProjectionOp::Node(n)])
899}
900
901fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
916 let events_path = paths.events();
917 let node_id = require_envelope_node_id(&events_path, ev)?;
918 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
919 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
920 match (code, signal) {
925 (Some(_), None) | (None, Some(_)) => {}
926 _ => {
927 return Err(Error::CorruptEventLog {
928 path: events_path,
929 reason: format!(
930 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
931 ev.seq
932 ),
933 });
934 }
935 }
936 let mut n = match read_node_opt(paths, &node_id)? {
937 Some(n) => n,
938 None => return Ok(vec![]),
945 };
946 if n.worker_exit.is_some() {
950 return Ok(vec![]);
951 }
952 n.worker_exit = Some(WorkerExit {
953 code,
954 signal,
955 at: ev.ts,
956 });
957 n.awaiting_input = None;
961 n.updated_at = ev.ts;
962 Ok(vec![ProjectionOp::Node(n)])
963}
964
965fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
978 let events_path = paths.events();
979 let node_id = require_envelope_node_id(&events_path, ev)?;
980 let mut n = match read_node_opt(paths, &node_id)? {
981 Some(n) => n,
982 None => return Ok(vec![]),
983 };
984 if n.first_death_at.is_some()
990 || n.status.is_terminal()
991 || n.worker_exit.is_some()
992 || n.last_report.is_some()
993 || n.pending_merge.is_some()
994 {
995 return Ok(vec![]);
996 }
997 n.first_death_at = Some(ev.ts);
998 n.updated_at = ev.ts;
999 Ok(vec![ProjectionOp::Node(n)])
1000}
1001
1002fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1011 const MAX_ITEMS: usize = 8;
1012 const MAX_TOPIC_CHARS: usize = 512;
1013 const MAX_OPTIONS: usize = 16;
1014 const MAX_OPTION_CHARS: usize = 256;
1015
1016 let events_path = paths.events();
1017 let node_id = require_envelope_node_id(&events_path, ev)?;
1018 let items = ev
1021 .data
1022 .get("discussion_items")
1023 .and_then(Value::as_array)
1024 .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1025 .ok_or_else(|| Error::CorruptEventLog {
1026 path: events_path.clone(),
1027 reason: format!(
1028 "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1029 ev.seq
1030 ),
1031 })?;
1032 for (index, item) in items.iter().enumerate() {
1033 let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1034 path: events_path.clone(),
1035 reason: format!(
1036 "event seq={} discussion_items[{index}] must be an object",
1037 ev.seq
1038 ),
1039 })?;
1040 let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1041 let default = obj
1042 .get("recommended_default")
1043 .and_then(Value::as_str)
1044 .unwrap_or("");
1045 let options = obj.get("options").and_then(Value::as_array);
1046 let options_valid = options.is_some_and(|values| {
1047 !values.is_empty()
1048 && values.len() <= MAX_OPTIONS
1049 && values.iter().all(|v| {
1050 v.as_str().is_some_and(|s| {
1051 !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1052 })
1053 })
1054 && values.iter().any(|v| v.as_str() == Some(default))
1055 });
1056 if topic.trim().is_empty()
1057 || topic.chars().count() > MAX_TOPIC_CHARS
1058 || default.trim().is_empty()
1059 || !options_valid
1060 {
1061 return Err(Error::CorruptEventLog {
1062 path: events_path.clone(),
1063 reason: format!(
1064 "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1065 ev.seq
1066 ),
1067 });
1068 }
1069 }
1070
1071 let mut n = match read_node_opt(paths, &node_id)? {
1072 Some(n) => n,
1073 None => return Ok(vec![]),
1074 };
1075 if n.status.is_terminal() || n.worker_exit.is_some() {
1076 return Ok(vec![]);
1077 }
1078 if let Some(open) = n.awaiting_input.as_mut() {
1079 if open.discussion_items.len() + items.len() > MAX_ITEMS {
1082 return Err(Error::CorruptEventLog {
1083 path: events_path,
1084 reason: format!(
1085 "event seq={} would exceed {MAX_ITEMS} open discussion items",
1086 ev.seq
1087 ),
1088 });
1089 }
1090 open.discussion_items.extend(items.iter().cloned());
1091 } else {
1092 n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1093 opened_at: ev.ts,
1094 event_seq: ev.seq,
1095 discussion_items: items.clone(),
1096 }));
1097 }
1098 n.updated_at = ev.ts;
1099 Ok(vec![ProjectionOp::Node(n)])
1100}
1101
1102fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1106 let events_path = paths.events();
1107 let node_id = require_envelope_node_id(&events_path, ev)?;
1108 let seq = ev
1111 .data
1112 .get("event_seq")
1113 .and_then(Value::as_u64)
1114 .ok_or_else(|| Error::CorruptEventLog {
1115 path: events_path.clone(),
1116 reason: format!(
1117 "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1118 ev.seq
1119 ),
1120 })?;
1121 let mut n = match read_node_opt(paths, &node_id)? {
1122 Some(n) => n,
1123 None => return Ok(vec![]),
1124 };
1125 let Some(open) = n.awaiting_input.as_ref() else {
1126 return Ok(vec![]);
1127 };
1128 if seq != open.event_seq {
1129 return Ok(vec![]);
1130 }
1131 n.awaiting_input = None;
1132 n.updated_at = ev.ts;
1133 Ok(vec![ProjectionOp::Node(n)])
1134}
1135
1136pub const KIND_MERGE_STARTED: &str = "merge.started";
1140
1141pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1146
1147fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1161 let events_path = paths.events();
1162 let node_id = require_envelope_node_id(&events_path, ev)?;
1163 let txn: MergeTxn =
1164 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1165 path: events_path.clone(),
1166 reason: format!(
1167 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1168 ev.seq
1169 ),
1170 })?;
1171 let mut n = match read_node_opt(paths, &node_id)? {
1172 Some(n) => n,
1173 None => return Ok(vec![]),
1174 };
1175 if n.status.is_terminal() {
1179 return Ok(vec![]);
1180 }
1181 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1184 return Ok(vec![]);
1185 }
1186 n.pending_merge = Some(Box::new(txn));
1187 n.updated_at = ev.ts;
1188 Ok(vec![ProjectionOp::Node(n)])
1189}
1190
1191fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1203 let events_path = paths.events();
1204 let node_id = require_envelope_node_id(&events_path, ev)?;
1205 let op_id = ev
1206 .data
1207 .get("op_id")
1208 .and_then(Value::as_str)
1209 .ok_or_else(|| Error::CorruptEventLog {
1210 path: events_path.clone(),
1211 reason: format!(
1212 "event seq={} kind=merge.aborted is missing string `op_id`",
1213 ev.seq
1214 ),
1215 })?;
1216 let mut n = match read_node_opt(paths, &node_id)? {
1217 Some(n) => n,
1218 None => return Ok(vec![]),
1219 };
1220 match n.pending_merge.as_ref() {
1223 Some(t) if t.op_id == op_id => {}
1224 _ => return Ok(vec![]),
1225 }
1226 n.pending_merge = None;
1227 n.updated_at = ev.ts;
1228 Ok(vec![ProjectionOp::Node(n)])
1229}
1230
1231fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1239 if current == incoming {
1240 tracing::debug!(
1241 target: "taskfleet_core::reducer",
1242 seq = ev.seq, kind = %ev.kind, status = ?current,
1243 "no-op: status re-applied to terminal target"
1244 );
1245 } else {
1246 tracing::warn!(
1247 target: "taskfleet_core::reducer",
1248 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1249 "no-op: ignored conflicting transition from terminal target"
1250 );
1251 }
1252}
1253
1254fn report_is_confirmed_explicit_merge(data: &Value) -> bool {
1271 ReportOrigin::report_is_confirmed_merge(data)
1272}
1273
1274fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1283 let corrupt = |reason: String| Error::CorruptEventLog {
1284 path: events_path.to_path_buf(),
1285 reason,
1286 };
1287 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1288 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1289 if cancelled {
1290 if success == Some(true) {
1291 return Err(corrupt(format!(
1292 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1293 ev.seq
1294 )));
1295 }
1296 Ok(Status::Cancelled)
1297 } else {
1298 match success {
1299 Some(true) => Ok(Status::Done),
1300 Some(false) => Ok(Status::Failed),
1301 None => Err(corrupt(format!(
1302 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1303 ev.seq
1304 ))),
1305 }
1306 }
1307}
1308
1309fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1310 let events_path = paths.events();
1313 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1314 path: events_path.clone(),
1315 reason: format!(
1316 "event seq={} kind=child.spawned missing parent `node_id`",
1317 ev.seq
1318 ),
1319 })?;
1320 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1321 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1322 let child_node_id = NodeId::parse_str(
1323 ev.data
1324 .get("child_node_id")
1325 .and_then(Value::as_str)
1326 .unwrap_or("n-0001"),
1327 )
1328 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1329 let mut n = match read_node_opt(paths, &parent_node_id)? {
1330 Some(n) => n,
1331 None => return Ok(vec![]),
1332 };
1333 let new_ref = ChildRef {
1334 run_id: child_run_id,
1335 node_id: child_node_id,
1336 };
1337 if n.children.iter().any(|c| c == &new_ref) {
1338 return Ok(vec![]);
1341 }
1342 n.children.push(new_ref);
1343 n.updated_at = ev.ts;
1344 Ok(vec![ProjectionOp::Node(n)])
1345}
1346
1347fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1358 let events_path = paths.events();
1359 let node_id = require_envelope_node_id(&events_path, ev)?;
1360 let raw = ev
1361 .data
1362 .get("pid")
1363 .and_then(Value::as_i64)
1364 .ok_or_else(|| Error::CorruptEventLog {
1365 path: events_path.clone(),
1366 reason: format!(
1367 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1368 ev.seq
1369 ),
1370 })?;
1371 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1372 path: events_path.clone(),
1373 reason: format!(
1374 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1375 ev.seq
1376 ),
1377 })?;
1378 let mut n = match read_node_opt(paths, &node_id)? {
1379 Some(n) => n,
1380 None => return Ok(vec![]),
1381 };
1382 if n.supervisor_pid == Some(pid) {
1383 return Ok(vec![]);
1384 }
1385 n.supervisor_pid = Some(pid);
1386 n.updated_at = ev.ts;
1387 Ok(vec![ProjectionOp::Node(n)])
1388}
1389
1390fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1401 let events_path = paths.events();
1402 let node_id = require_envelope_node_id(&events_path, ev)?;
1403 let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1404 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1408 let report_seq = ev
1409 .data
1410 .get("report_seq")
1411 .and_then(Value::as_u64)
1412 .ok_or_else(|| Error::CorruptEventLog {
1413 path: events_path.clone(),
1414 reason: format!(
1415 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1416 ev.seq
1417 ),
1418 })?;
1419 let mut n = match read_node_opt(paths, &node_id)? {
1420 Some(n) => n,
1421 None => return Ok(vec![]),
1422 };
1423 if let Some(prev) = n
1424 .last_processed_report_seq_by_child
1425 .get(child_run_id)
1426 .and_then(Value::as_u64)
1427 {
1428 if report_seq <= prev {
1429 return Ok(vec![]);
1430 }
1431 }
1432 n.last_processed_report_seq_by_child
1433 .insert(child_run_id.to_string(), Value::from(report_seq));
1434 n.updated_at = ev.ts;
1435 Ok(vec![ProjectionOp::Node(n)])
1436}
1437
1438#[cfg(test)]
1439mod tests {
1440 use super::*;
1441 use crate::schema::Event;
1442 use chrono::Utc;
1443 use tempfile::TempDir;
1444
1445 fn event(run_id: &str) -> Event {
1446 Event {
1447 ts: Utc::now(),
1448 seq: 1,
1449 kind: "run.status".into(),
1450 run_id: RunId::parse_str(run_id).unwrap(),
1451 node_id: None,
1452 idempotency_key: None,
1453 data: serde_json::json!({ "status": "running" }),
1454 }
1455 }
1456
1457 #[test]
1458 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1459 let tmp = TempDir::new().unwrap();
1463 let run_id = "01jxsnap000000000000000000";
1464 let rid = RunId::parse_str(run_id).unwrap();
1465 let dir = crate::run_dir(tmp.path(), &rid);
1466 std::fs::create_dir_all(&dir).unwrap();
1467 let paths = RunPaths::new(dir, run_id).unwrap();
1468
1469 let mut created = event(run_id);
1472 created.kind = "run.created".into();
1473 created.data = serde_json::json!({
1474 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1475 });
1476 apply_event(&paths, &created).expect("run.created applies");
1477 let manifest_before = std::fs::read(paths.manifest()).unwrap();
1478
1479 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1480 let mut ev = event(run_id);
1481 ev.seq = seq;
1482 ev.kind = kind.into();
1483 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1485 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1486 assert!(ops.is_empty(), "{kind} must plan no projection ops");
1487 apply_event(&paths, &ev).expect("audit kind applies as no-op");
1489 }
1490
1491 assert_eq!(
1493 std::fs::read(paths.manifest()).unwrap(),
1494 manifest_before,
1495 "audit events must not mutate the manifest"
1496 );
1497 assert!(!paths.nodes_dir().exists(), "no node projection created");
1498 }
1499
1500 #[test]
1501 fn run_created_folds_harness_when_present_and_defaults_none() {
1502 let tmp = TempDir::new().unwrap();
1503
1504 let run_id = "01jxhrnsaa0000000000000001";
1506 let rid = RunId::parse_str(run_id).unwrap();
1507 let dir = crate::run_dir(tmp.path(), &rid);
1508 std::fs::create_dir_all(&dir).unwrap();
1509 let paths = RunPaths::new(dir, run_id).unwrap();
1510 let mut created = event(run_id);
1511 created.kind = "run.created".into();
1512 created.data = serde_json::json!({
1513 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1514 "harness": "pi", "harness_source": "flag",
1515 });
1516 apply_event(&paths, &created).expect("run.created applies");
1517 let m = read_manifest_opt(&paths).unwrap().unwrap();
1518 assert_eq!(m.harness.as_deref(), Some("pi"));
1519
1520 let run_id2 = "01jxhrnsaa0000000000000002";
1522 let rid2 = RunId::parse_str(run_id2).unwrap();
1523 let dir2 = crate::run_dir(tmp.path(), &rid2);
1524 std::fs::create_dir_all(&dir2).unwrap();
1525 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1526 let mut created2 = event(run_id2);
1527 created2.kind = "run.created".into();
1528 created2.data =
1529 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1530 apply_event(&paths2, &created2).expect("run.created applies");
1531 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1532 assert_eq!(m2.harness, None);
1533 }
1534
1535 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1538 let rid = RunId::parse_str(run_id).unwrap();
1539 let dir = crate::run_dir(tmp.path(), &rid);
1540 std::fs::create_dir_all(&dir).unwrap();
1541 let paths = RunPaths::new(dir, run_id).unwrap();
1542 let mut created = event(run_id);
1543 created.kind = "run.created".into();
1544 created.data =
1545 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1546 apply_event(&paths, &created).expect("run.created applies");
1547 let mut node = event(run_id);
1548 node.seq = 2;
1549 node.kind = "node.created".into();
1550 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1551 node.data = serde_json::json!({
1552 "kind": "spinoff",
1553 "branch": "wt/foo",
1554 "worktree_path": "/tmp/old-wt",
1555 "agent_pid": 111,
1556 });
1557 apply_event(&paths, &node).expect("node.created applies");
1558 paths
1559 }
1560
1561 #[test]
1565 fn node_retry_rewires_node_and_increments_attempts() {
1566 let tmp = TempDir::new().unwrap();
1567 let run_id = "01jxsnap000000000000000000";
1568 let paths = bootstrap_retry_node(&tmp, run_id);
1569
1570 let mut retry = event(run_id);
1571 retry.seq = 3;
1572 retry.kind = "node.retry".into();
1573 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1574 retry.data = serde_json::json!({
1575 "attempt": 1,
1576 "reason": "agent-died",
1577 "branch": "wt/foo-r1",
1578 "base_sha": "a".repeat(40),
1579 "worktree_path": "/tmp/new-wt",
1580 "agent_pid": 222,
1581 "tmux_session": "s",
1582 "tmux_window_id": "@9",
1583 });
1584 apply_event(&paths, &retry).expect("node.retry applies");
1585
1586 let n = read_n0001(&paths);
1587 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1588 assert_eq!(
1589 n.branch.as_deref(),
1590 Some("wt/foo-r1"),
1591 "rewired to new branch"
1592 );
1593 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1594 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1595 assert_eq!(n.status, Status::Pending, "node returns to pending");
1596 assert!(n.last_report.is_none());
1597 assert_eq!(
1598 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1599 Some("@9"),
1600 "rewired tmux identity"
1601 );
1602
1603 let mut retry2 = event(run_id);
1605 retry2.seq = 4;
1606 retry2.kind = "node.retry".into();
1607 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1608 retry2.data = serde_json::json!({
1609 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1610 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1611 });
1612 apply_event(&paths, &retry2).expect("node.retry applies");
1613 assert_eq!(read_n0001(&paths).retry_attempts, 2);
1614 }
1615
1616 #[test]
1620 fn node_retry_against_terminal_node_is_noop() {
1621 let tmp = TempDir::new().unwrap();
1622 let run_id = "01jxsnap000000000000000000";
1623 let paths = bootstrap_retry_node(&tmp, run_id);
1624
1625 let mut report = event(run_id);
1627 report.seq = 3;
1628 report.kind = "node.report".into();
1629 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1630 report.data = serde_json::json!({ "success": true });
1631 apply_event(&paths, &report).expect("node.report applies");
1632 assert_eq!(read_n0001(&paths).status, Status::Done);
1633
1634 let mut retry = event(run_id);
1635 retry.seq = 4;
1636 retry.kind = "node.retry".into();
1637 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1638 retry.data = serde_json::json!({
1639 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1640 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1641 });
1642 apply_event(&paths, &retry).expect("node.retry applies as no-op");
1643
1644 let n = read_n0001(&paths);
1645 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
1646 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
1647 assert_eq!(n.agent_pid, Some(111), "not rewired");
1648 }
1649
1650 #[test]
1651 fn apply_event_rejects_event_from_a_different_run() {
1652 let tmp = TempDir::new().unwrap();
1653 let run_id = "01jxsnap000000000000000000";
1654 let rid = RunId::parse_str(run_id).unwrap();
1655 let dir = crate::run_dir(tmp.path(), &rid);
1656 std::fs::create_dir_all(&dir).unwrap();
1657 let paths = RunPaths::new(dir, run_id).unwrap();
1658
1659 let foreign = event("02jxsnap000000000000000000");
1661 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
1662 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
1663
1664 let mine = event(run_id);
1667 apply_event(&paths, &mine).expect("matching run_id must be accepted");
1668 }
1669
1670 #[test]
1671 fn tmux_identity_from_data_reads_qualified_fields() {
1672 let d = serde_json::json!({
1673 "tmux_socket": "/private/tmp/tmux-501/default",
1674 "tmux_session": "octl",
1675 "tmux_window_id": "@42",
1676 });
1677 let id = tmux_identity_from_data(&d).expect("qualified identity");
1678 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
1679 assert_eq!(id.session, "octl");
1680 assert_eq!(id.window_id, "@42");
1681 assert_eq!(id.pane_id, None);
1683
1684 let d2 = serde_json::json!({
1686 "tmux_socket": null,
1687 "tmux_session": "octl",
1688 "tmux_window_id": "@7",
1689 });
1690 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
1691 assert_eq!(id2.socket, None);
1692 assert_eq!(id2.window_id, "@7");
1693
1694 let d3 = serde_json::json!({
1696 "tmux_session": "octl",
1697 "tmux_window_id": "@42",
1698 "tmux_pane_id": "%7",
1699 });
1700 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
1701 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
1702 assert_eq!(id3.capture_target(), "%7");
1703
1704 let d4 = serde_json::json!({
1707 "tmux_session": "octl",
1708 "tmux_window_id": "@42",
1709 "tmux_pane_id": null,
1710 });
1711 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
1712 assert_eq!(id4.pane_id, None);
1713 assert_eq!(id4.capture_target(), "@42");
1714 }
1715
1716 #[test]
1717 fn tmux_identity_from_data_back_compat_is_none() {
1718 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
1720 assert!(tmux_identity_from_data(&legacy).is_none());
1721 let partial = serde_json::json!({ "tmux_window_id": "@42" });
1723 assert!(tmux_identity_from_data(&partial).is_none());
1724 }
1725
1726 #[test]
1729 fn node_created_populates_tmux_identity() {
1730 let tmp = TempDir::new().unwrap();
1731 let run_id = "01jxsnap000000000000000000";
1732 let rid = RunId::parse_str(run_id).unwrap();
1733 let dir = crate::run_dir(tmp.path(), &rid);
1734 std::fs::create_dir_all(&dir).unwrap();
1735 let paths = RunPaths::new(dir, run_id).unwrap();
1736
1737 let mut ev = event(run_id);
1738 ev.seq = 2;
1739 ev.kind = "node.created".into();
1740 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1741 ev.data = serde_json::json!({
1742 "kind": "spinoff",
1743 "tmux_window": "🚀 wt/x",
1744 "tmux_socket": "/private/tmp/tmux-501/default",
1745 "tmux_session": "octl",
1746 "tmux_window_id": "@42",
1747 });
1748 apply_event(&paths, &ev).expect("node.created applies");
1749 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
1750 .unwrap()
1751 .unwrap();
1752 let id = n.tmux_identity.expect("qualified identity recorded");
1753 assert_eq!(id.session, "octl");
1754 assert_eq!(id.window_id, "@42");
1755 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
1756
1757 let run2 = "02jxsnap000000000000000000";
1759 let rid2 = RunId::parse_str(run2).unwrap();
1760 let dir2 = crate::run_dir(tmp.path(), &rid2);
1761 std::fs::create_dir_all(&dir2).unwrap();
1762 let paths2 = RunPaths::new(dir2, run2).unwrap();
1763 let mut ev2 = event(run2);
1764 ev2.seq = 2;
1765 ev2.kind = "node.created".into();
1766 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1767 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
1768 apply_event(&paths2, &ev2).expect("legacy node.created applies");
1769 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
1770 .unwrap()
1771 .unwrap();
1772 assert!(n2.tmux_identity.is_none());
1773 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
1774 }
1775
1776 #[test]
1777 fn node_materialization_populates_missing_manifest_source_branch() {
1778 let tmp = TempDir::new().unwrap();
1779 let run_id = "01jxsnap000000000000000001";
1780 let rid = RunId::parse_str(run_id).unwrap();
1781 let dir = crate::run_dir(tmp.path(), &rid);
1782 std::fs::create_dir_all(&dir).unwrap();
1783 let paths = RunPaths::new(dir, run_id).unwrap();
1784
1785 let mut created = event(run_id);
1786 created.kind = "run.created".into();
1787 created.data = serde_json::json!({
1788 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1789 });
1790 apply_event(&paths, &created).unwrap();
1791 assert!(read_manifest_opt(&paths)
1792 .unwrap()
1793 .unwrap()
1794 .source_branch
1795 .is_none());
1796
1797 let mut node = event(run_id);
1798 node.seq = 2;
1799 node.kind = "node.created".into();
1800 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1801 node.data = serde_json::json!({
1802 "kind": "spinoff",
1803 "source_branch": "main",
1804 "worktree_path": "/tmp/wt/pending"
1805 });
1806 apply_event(&paths, &node).unwrap();
1807
1808 let manifest = read_manifest_opt(&paths).unwrap().unwrap();
1809 assert_eq!(manifest.status, Status::Pending);
1810 assert_eq!(manifest.source_branch.as_deref(), Some("main"));
1811 let node = read_n0001(&paths);
1812 assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
1813 }
1814
1815 #[test]
1816 fn node_materialization_preserves_explicit_manifest_source_branch() {
1817 let tmp = TempDir::new().unwrap();
1818 let run_id = "01jxsnap000000000000000002";
1819 let rid = RunId::parse_str(run_id).unwrap();
1820 let dir = crate::run_dir(tmp.path(), &rid);
1821 std::fs::create_dir_all(&dir).unwrap();
1822 let paths = RunPaths::new(dir, run_id).unwrap();
1823
1824 let mut created = event(run_id);
1825 created.kind = "run.created".into();
1826 created.data = serde_json::json!({
1827 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1828 "source_branch": "release"
1829 });
1830 apply_event(&paths, &created).unwrap();
1831
1832 let mut node = event(run_id);
1833 node.seq = 2;
1834 node.kind = "node.created".into();
1835 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1836 node.data = serde_json::json!({
1837 "kind": "spinoff", "source_branch": "main"
1838 });
1839 apply_event(&paths, &node).unwrap();
1840
1841 assert_eq!(
1842 read_manifest_opt(&paths)
1843 .unwrap()
1844 .unwrap()
1845 .source_branch
1846 .as_deref(),
1847 Some("release")
1848 );
1849 }
1850
1851 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1854 let rid = RunId::parse_str(run_id).unwrap();
1855 let dir = crate::run_dir(tmp.path(), &rid);
1856 std::fs::create_dir_all(&dir).unwrap();
1857 let paths = RunPaths::new(dir, run_id).unwrap();
1858
1859 let mut created = event(run_id);
1860 created.kind = "run.created".into();
1861 created.data = serde_json::json!({
1862 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1863 });
1864 apply_event(&paths, &created).expect("run.created applies");
1865
1866 let mut node = event(run_id);
1867 node.seq = 2;
1868 node.kind = "node.created".into();
1869 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1870 node.data = serde_json::json!({ "kind": "spinoff" });
1871 apply_event(&paths, &node).expect("node.created applies");
1872 paths
1873 }
1874
1875 fn read_n0001(paths: &RunPaths) -> Node {
1876 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
1877 .unwrap()
1878 .unwrap()
1879 }
1880
1881 #[test]
1882 fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
1883 let tmp = TempDir::new().unwrap();
1884 let run_id = "01jxwd0000000000000000000w";
1885 let paths = seed_run_with_node(&tmp, run_id);
1886 let nid = Some(NodeId::parse_str("n-0001").unwrap());
1887 let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
1888 let mut open = event(run_id);
1889 open.seq = 3;
1890 open.ts = opened_at;
1891 open.kind = "node.awaiting_input".into();
1892 open.node_id = nid.clone();
1893 open.data = serde_json::json!({ "discussion_items": [{
1894 "topic": "Which scope?",
1895 "options": ["small", "large"],
1896 "recommended_default": "small"
1897 }] });
1898 apply_event(&paths, &open).unwrap();
1899 let first = read_n0001(&paths).awaiting_input.unwrap();
1900 assert_eq!(first.opened_at, opened_at);
1901 assert_eq!(first.event_seq, 3);
1902
1903 let mut duplicate = open.clone();
1905 duplicate.seq = 4;
1906 duplicate.ts = opened_at + chrono::Duration::hours(1);
1907 apply_event(&paths, &duplicate).unwrap();
1908 let still_first = read_n0001(&paths).awaiting_input.unwrap();
1909 assert_eq!(still_first.opened_at, opened_at);
1910 assert_eq!(still_first.event_seq, 3);
1911
1912 let mut stale = event(run_id);
1914 stale.seq = 5;
1915 stale.kind = "node.input_resolved".into();
1916 stale.node_id = nid.clone();
1917 stale.data = serde_json::json!({ "event_seq": 2 });
1918 apply_event(&paths, &stale).unwrap();
1919 assert!(read_n0001(&paths).awaiting_input.is_some());
1920
1921 let mut resolved = stale;
1922 resolved.seq = 6;
1923 resolved.data = serde_json::json!({ "event_seq": 3 });
1924 apply_event(&paths, &resolved).unwrap();
1925 assert!(read_n0001(&paths).awaiting_input.is_none());
1926 }
1927
1928 #[test]
1929 fn awaiting_input_rejects_missing_default_without_mutating_projection() {
1930 let tmp = TempDir::new().unwrap();
1931 let run_id = "01jxwd0000000000000000000x";
1932 let paths = seed_run_with_node(&tmp, run_id);
1933 let mut open = event(run_id);
1934 open.seq = 3;
1935 open.kind = "node.awaiting_input".into();
1936 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1937 open.data = serde_json::json!({ "discussion_items": [{
1938 "topic": "Which scope?", "options": ["small", "large"]
1939 }] });
1940 assert!(reduce_event_to_ops(&paths, &open).is_err());
1941 assert!(read_n0001(&paths).awaiting_input.is_none());
1942 }
1943
1944 #[test]
1945 fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
1946 let tmp = TempDir::new().unwrap();
1947 let run_id = "01jxwd0000000000000000000y";
1948 let paths = seed_run_with_node(&tmp, run_id);
1949 let nid = Some(NodeId::parse_str("n-0001").unwrap());
1950
1951 let mut open = event(run_id);
1952 open.seq = 3;
1953 open.kind = "node.awaiting_input".into();
1954 open.node_id = nid.clone();
1955 open.data = serde_json::json!({ "discussion_items": [{
1956 "topic": "Which scope?", "options": ["small", "large"],
1957 "recommended_default": "small"
1958 }] });
1959 apply_event(&paths, &open).unwrap();
1960
1961 let mut malformed_duplicate = open.clone();
1962 malformed_duplicate.seq = 4;
1963 malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
1964 assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
1965
1966 let mut exited = event(run_id);
1967 exited.seq = 5;
1968 exited.kind = "worker.exited".into();
1969 exited.node_id = nid;
1970 exited.data = serde_json::json!({ "exit_code": 0 });
1971 apply_event(&paths, &exited).unwrap();
1972 let node = read_n0001(&paths);
1973 assert!(node.awaiting_input.is_none());
1974 assert!(node.worker_exit.is_some());
1975
1976 let mut delayed_open = open;
1977 delayed_open.seq = 6;
1978 assert!(reduce_event_to_ops(&paths, &delayed_open)
1979 .unwrap()
1980 .is_empty());
1981 }
1982
1983 #[test]
1984 fn input_resolved_requires_generation_even_when_nothing_is_open() {
1985 let tmp = TempDir::new().unwrap();
1986 let run_id = "01jxwd0000000000000000000z";
1987 let paths = seed_run_with_node(&tmp, run_id);
1988 let mut resolved = event(run_id);
1989 resolved.seq = 3;
1990 resolved.kind = "node.input_resolved".into();
1991 resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1992 resolved.data = serde_json::json!({});
1993 assert!(reduce_event_to_ops(&paths, &resolved).is_err());
1994 }
1995
1996 #[test]
1997 fn awaiting_input_default_must_be_one_of_options() {
1998 let tmp = TempDir::new().unwrap();
1999 let run_id = "01jxwd00000000000000000010";
2000 let paths = seed_run_with_node(&tmp, run_id);
2001 let mut open = event(run_id);
2002 open.seq = 3;
2003 open.kind = "node.awaiting_input".into();
2004 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2005 open.data = serde_json::json!({ "discussion_items": [{
2006 "topic": "Which scope?", "options": ["small", "large"],
2007 "recommended_default": "other"
2008 }] });
2009 assert!(reduce_event_to_ops(&paths, &open).is_err());
2010 }
2011
2012 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
2013 let mut ev = event(run_id);
2014 ev.seq = seq;
2015 ev.kind = KIND_MERGE_STARTED.into();
2016 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2017 ev.data = serde_json::json!({
2018 "op_id": op_id,
2019 "source_branch": "main",
2020 "worker_branch": "wt/worker",
2021 "expected_source_oid": expected,
2022 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2023 "base_sha": null,
2024 "driver_pid": 4242,
2025 "driver_pid_start_secs": null,
2026 "started_at": "2026-08-15T00:00:00Z",
2027 });
2028 ev
2029 }
2030
2031 #[test]
2034 fn merge_started_records_pending_transaction() {
2035 let tmp = TempDir::new().unwrap();
2036 let run_id = "01jxsnap000000000000000000";
2037 let paths = seed_run_with_node(&tmp, run_id);
2038
2039 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2040 let n = read_n0001(&paths);
2041 assert_eq!(
2042 n.status,
2043 Status::Pending,
2044 "recording a merge is not terminal"
2045 );
2046 let txn = n.pending_merge.expect("transaction recorded");
2047 assert_eq!(txn.op_id, "op-1");
2048 assert_eq!(txn.expected_source_oid, "aaa");
2049 }
2050
2051 #[test]
2054 fn merge_aborted_clears_matching_transaction_only() {
2055 let tmp = TempDir::new().unwrap();
2056 let run_id = "01jxsnap000000000000000000";
2057 let paths = seed_run_with_node(&tmp, run_id);
2058 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2059
2060 let mut stale = event(run_id);
2062 stale.seq = 4;
2063 stale.kind = KIND_MERGE_ABORTED.into();
2064 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2065 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2066 apply_event(&paths, &stale).unwrap();
2067 assert!(
2068 read_n0001(&paths).pending_merge.is_some(),
2069 "stale abort is a no-op"
2070 );
2071
2072 let mut abort = event(run_id);
2074 abort.seq = 5;
2075 abort.kind = KIND_MERGE_ABORTED.into();
2076 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2077 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2078 apply_event(&paths, &abort).unwrap();
2079 let n = read_n0001(&paths);
2080 assert!(
2081 n.pending_merge.is_none(),
2082 "matching abort clears the transaction"
2083 );
2084 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2085 }
2086
2087 #[test]
2090 fn terminal_report_clears_pending_merge() {
2091 let tmp = TempDir::new().unwrap();
2092 let run_id = "01jxsnap000000000000000000";
2093 let paths = seed_run_with_node(&tmp, run_id);
2094 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2095
2096 let mut report = event(run_id);
2097 report.seq = 4;
2098 report.kind = "node.report".into();
2099 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2100 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2101 apply_event(&paths, &report).unwrap();
2102 let n = read_n0001(&paths);
2103 assert_eq!(n.status, Status::Done);
2104 assert!(
2105 n.pending_merge.is_none(),
2106 "completed merge clears the transaction"
2107 );
2108 }
2109
2110 #[test]
2119 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2120 let tmp = TempDir::new().unwrap();
2121
2122 let drive = |run_id: &str, report_data: Value| -> Status {
2125 let paths = seed_run_with_node(&tmp, run_id);
2126 let mut fail = event(run_id);
2128 fail.seq = 3;
2129 fail.kind = "node.status".into();
2130 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2131 fail.data = serde_json::json!({ "status": "failed" });
2132 apply_event(&paths, &fail).unwrap();
2133 assert_eq!(read_n0001(&paths).status, Status::Failed);
2134 let mut report = event(run_id);
2136 report.seq = 4;
2137 report.kind = "node.report".into();
2138 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2139 report.data = report_data;
2140 apply_event(&paths, &report).unwrap();
2141 read_n0001(&paths).status
2142 };
2143
2144 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2146 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2147 assert_eq!(
2148 drive("01jxsnap000000000000000001", agent_forged),
2149 Status::Failed,
2150 "an Agent-origin report with a forged via must not be adopted"
2151 );
2152
2153 let malformed = serde_json::json!({
2155 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2156 });
2157 assert_eq!(
2158 drive("01jxsnap000000000000000002", malformed),
2159 Status::Failed,
2160 "a malformed origin must not re-unlock the legacy via adoption path"
2161 );
2162
2163 let mut run_merge = serde_json::json!({ "success": true });
2165 crate::ReportOrigin::RunMerge {
2166 op_id: Some("op-1".into()),
2167 worker_oid: Some("cafebabe".into()),
2168 }
2169 .stamp(&mut run_merge);
2170 assert_eq!(
2171 drive("01jxsnap000000000000000003", run_merge),
2172 Status::Done,
2173 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2174 );
2175
2176 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2179 assert_eq!(
2180 drive("01jxsnap000000000000000004", legacy),
2181 Status::Done,
2182 "a legacy via-only report (no origin field) is still adopted"
2183 );
2184 }
2185
2186 #[test]
2190 fn terminal_node_status_clears_pending_merge() {
2191 let tmp = TempDir::new().unwrap();
2192 let run_id = "01jxsnap000000000000000000";
2193 let paths = seed_run_with_node(&tmp, run_id);
2194 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2195 assert!(read_n0001(&paths).pending_merge.is_some());
2196
2197 let mut status = event(run_id);
2198 status.seq = 4;
2199 status.kind = "node.status".into();
2200 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2201 status.data = serde_json::json!({ "status": "failed" });
2202 apply_event(&paths, &status).unwrap();
2203 let n = read_n0001(&paths);
2204 assert_eq!(n.status, Status::Failed);
2205 assert!(
2206 n.pending_merge.is_none(),
2207 "terminal status clears the transaction"
2208 );
2209 }
2210
2211 #[test]
2215 fn supervisor_attached_sets_supervisor_pid() {
2216 let tmp = TempDir::new().unwrap();
2217 let run_id = "01jxsnap000000000000000000";
2218 let paths = seed_run_with_node(&tmp, run_id);
2219 assert_eq!(read_n0001(&paths).supervisor_pid, None);
2220
2221 let mut ev = event(run_id);
2222 ev.seq = 3;
2223 ev.kind = "supervisor.attached".into();
2224 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2225 ev.data = serde_json::json!({ "pid": 47820 });
2226 apply_event(&paths, &ev).expect("supervisor.attached applies");
2227 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2228 }
2229
2230 #[test]
2233 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2234 let tmp = TempDir::new().unwrap();
2235 let run_id = "01jxsnap000000000000000000";
2236 let paths = seed_run_with_node(&tmp, run_id);
2237
2238 let mut ev = event(run_id);
2239 ev.kind = "supervisor.attached".into();
2240 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2241
2242 ev.seq = 3;
2243 ev.data = serde_json::json!({ "pid": 100 });
2244 apply_event(&paths, &ev).expect("first attach applies");
2245 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2246
2247 ev.seq = 4;
2249 ev.data = serde_json::json!({ "pid": 200 });
2250 apply_event(&paths, &ev).expect("second attach applies");
2251 let after_second = read_n0001(&paths);
2252 assert_eq!(after_second.supervisor_pid, Some(200));
2253
2254 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2257 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2258 apply_event(&paths, &ev).expect("replay applies as no-op");
2259 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2260 }
2261
2262 #[test]
2265 fn supervisor_cursor_advanced_sets_report_cursor() {
2266 let tmp = TempDir::new().unwrap();
2267 let run_id = "01jxsnap000000000000000000";
2268 let paths = seed_run_with_node(&tmp, run_id);
2269 let child = "02jxsnap000000000000000000";
2270
2271 let mut ev = event(run_id);
2272 ev.seq = 3;
2273 ev.kind = "supervisor.cursor_advanced".into();
2274 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2275 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2276 apply_event(&paths, &ev).expect("cursor_advanced applies");
2277
2278 let n = read_n0001(&paths);
2279 assert_eq!(
2280 n.last_processed_report_seq_by_child.get(child),
2281 Some(&Value::from(7u64))
2282 );
2283 }
2284
2285 #[test]
2290 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2291 let tmp = TempDir::new().unwrap();
2292 let run_id = "01jxsnap000000000000000000";
2293 let paths = seed_run_with_node(&tmp, run_id);
2294 let child_a = "02jxsnap000000000000000000";
2295 let child_b = "03jxsnap000000000000000000";
2296
2297 let mut ev = event(run_id);
2298 ev.kind = "supervisor.cursor_advanced".into();
2299 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2300
2301 ev.seq = 3;
2302 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2303 apply_event(&paths, &ev).expect("seq 5 applies");
2304
2305 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2307 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2308
2309 ev.seq = 4;
2311 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2312 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2313 assert!(ops.is_empty(), "older seq must plan no ops");
2314 apply_event(&paths, &ev).expect("older seq applies as no-op");
2315 assert_eq!(
2316 read_n0001(&paths)
2317 .last_processed_report_seq_by_child
2318 .get(child_a),
2319 Some(&Value::from(5u64))
2320 );
2321
2322 ev.seq = 5;
2324 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2325 apply_event(&paths, &ev).expect("higher seq applies");
2326 ev.seq = 6;
2327 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2328 apply_event(&paths, &ev).expect("second child applies");
2329
2330 let n = read_n0001(&paths);
2331 assert_eq!(
2332 n.last_processed_report_seq_by_child.get(child_a),
2333 Some(&Value::from(9u64))
2334 );
2335 assert_eq!(
2336 n.last_processed_report_seq_by_child.get(child_b),
2337 Some(&Value::from(1u64))
2338 );
2339 }
2340
2341 #[test]
2344 fn supervisor_state_events_reject_malformed_payloads() {
2345 let tmp = TempDir::new().unwrap();
2346 let run_id = "01jxsnap000000000000000000";
2347 let paths = seed_run_with_node(&tmp, run_id);
2348 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2349
2350 let mut ev = event(run_id);
2352 ev.seq = 3;
2353 ev.kind = "supervisor.attached".into();
2354 ev.node_id = nid.clone();
2355 ev.data = serde_json::json!({});
2356 assert!(matches!(
2357 reduce_event_to_ops(&paths, &ev),
2358 Err(Error::CorruptEventLog { .. })
2359 ));
2360
2361 ev.node_id = None;
2363 ev.data = serde_json::json!({ "pid": 1 });
2364 assert!(matches!(
2365 reduce_event_to_ops(&paths, &ev),
2366 Err(Error::CorruptEventLog { .. })
2367 ));
2368
2369 let mut ev2 = event(run_id);
2371 ev2.seq = 4;
2372 ev2.kind = "supervisor.cursor_advanced".into();
2373 ev2.node_id = nid.clone();
2374 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
2375 assert!(matches!(
2376 reduce_event_to_ops(&paths, &ev2),
2377 Err(Error::CorruptEventLog { .. })
2378 ));
2379
2380 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2382 assert!(matches!(
2383 reduce_event_to_ops(&paths, &ev2),
2384 Err(Error::CorruptEventLog { .. })
2385 ));
2386 }
2387
2388 #[test]
2395 fn removed_or_garbage_kind_in_created_events_is_rejected() {
2396 let tmp = TempDir::new().unwrap();
2397 let run_id = "01jxsnap000000000000000000";
2398 let rid = RunId::parse_str(run_id).unwrap();
2399 let dir = crate::run_dir(tmp.path(), &rid);
2400 std::fs::create_dir_all(&dir).unwrap();
2401 let paths = RunPaths::new(dir, run_id).unwrap();
2402
2403 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2404 let mut ev = event(run_id);
2405 ev.kind = "run.created".into();
2406 ev.node_id = None;
2407 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2408 assert!(
2409 matches!(
2410 reduce_event_to_ops(&paths, &ev),
2411 Err(Error::CorruptEventLog { .. })
2412 ),
2413 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2414 );
2415 }
2416
2417 let mut ok = event(run_id);
2420 ok.kind = "run.created".into();
2421 ok.node_id = None;
2422 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2423 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2424 }
2425
2426 #[cfg(unix)]
2435 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2436 use std::os::unix::fs::MetadataExt;
2437 let mut consider = vec![paths.manifest()];
2438 for dir in [paths.nodes_dir()] {
2439 if let Ok(rd) = std::fs::read_dir(&dir) {
2440 for ent in rd.flatten() {
2441 let p = ent.path();
2442 if p.extension().and_then(|s| s.to_str()) == Some("json") {
2443 consider.push(p);
2444 }
2445 }
2446 }
2447 }
2448 let mut map = std::collections::BTreeMap::new();
2449 for p in consider {
2450 if let Ok(md) = std::fs::symlink_metadata(&p) {
2451 if md.file_type().is_file() {
2452 map.insert(p, md.ino());
2453 }
2454 }
2455 }
2456 map
2457 }
2458
2459 #[cfg(unix)]
2468 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2469 use std::collections::BTreeSet;
2470 let before = projection_inodes(paths);
2471 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2472 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2473 .into_iter()
2474 .collect();
2475 apply_event(paths, ev)
2476 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2477 let after = projection_inodes(paths);
2478 let touched: BTreeSet<PathBuf> = after
2479 .iter()
2480 .filter(|(p, ino)| before.get(*p) != Some(*ino))
2481 .map(|(p, _)| p.clone())
2482 .collect();
2483 assert_eq!(
2484 planned, touched,
2485 "kind={}: plan_projections must name exactly the files apply_event writes",
2486 ev.kind
2487 );
2488 if expect_writes {
2489 assert!(
2490 !touched.is_empty(),
2491 "kind={}: expected this event to write at least one projection",
2492 ev.kind
2493 );
2494 }
2495 }
2496
2497 #[cfg(unix)]
2504 #[test]
2505 fn plan_projections_matches_apply_for_every_kind() {
2506 let tmp = TempDir::new().unwrap();
2507 let run_id = "01jxsnap000000000000000000";
2508 let rid = RunId::parse_str(run_id).unwrap();
2509 let dir = crate::run_dir(tmp.path(), &rid);
2510 std::fs::create_dir_all(&dir).unwrap();
2511 let paths = RunPaths::new(dir, run_id).unwrap();
2512 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2513 let child = "02jxsnap000000000000000000";
2514
2515 let mut next_seq = 0u64;
2517 let mut at = |kind: &str, node_id, data| {
2518 next_seq += 1;
2519 Event {
2520 ts: Utc::now(),
2521 seq: next_seq,
2522 kind: kind.into(),
2523 run_id: rid.clone(),
2524 node_id,
2525 idempotency_key: None,
2526 data,
2527 }
2528 };
2529
2530 assert_plan_matches_apply(
2532 &paths,
2533 &at(
2534 "run.created",
2535 None,
2536 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2537 ),
2538 true,
2539 );
2540 assert_plan_matches_apply(
2542 &paths,
2543 &at(
2544 "run.status",
2545 None,
2546 serde_json::json!({ "status": "running" }),
2547 ),
2548 true,
2549 );
2550 assert_plan_matches_apply(
2552 &paths,
2553 &at(
2554 "node.created",
2555 nid(),
2556 serde_json::json!({ "kind": "spinoff" }),
2557 ),
2558 true,
2559 );
2560 assert_plan_matches_apply(
2562 &paths,
2563 &at(
2564 "node.status",
2565 nid(),
2566 serde_json::json!({ "status": "running" }),
2567 ),
2568 true,
2569 );
2570 assert_plan_matches_apply(
2572 &paths,
2573 &at(
2574 "supervisor.attached",
2575 nid(),
2576 serde_json::json!({ "pid": 4242 }),
2577 ),
2578 true,
2579 );
2580 assert_plan_matches_apply(
2582 &paths,
2583 &at(
2584 "supervisor.cursor_advanced",
2585 nid(),
2586 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2587 ),
2588 true,
2589 );
2590 assert_plan_matches_apply(
2592 &paths,
2593 &at(
2594 "child.spawned",
2595 nid(),
2596 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2597 ),
2598 true,
2599 );
2600 assert_plan_matches_apply(
2602 &paths,
2603 &at("node.report", nid(), serde_json::json!({ "success": true })),
2604 true,
2605 );
2606 assert_plan_matches_apply(
2609 &paths,
2610 &at(
2611 "node.status",
2612 nid(),
2613 serde_json::json!({ "status": "failed" }),
2614 ),
2615 false,
2616 );
2617 for kind in [
2619 "supervisor.exited",
2620 "orchestrator.decision",
2621 "discuss.critical",
2622 "cleanup.window_missing",
2623 ] {
2624 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
2625 }
2626 }
2627
2628 #[test]
2632 fn worker_exited_records_clean_exit_without_transitioning_status() {
2633 let tmp = TempDir::new().unwrap();
2634 let run_id = "01jxsnap000000000000000000";
2635 let paths = bootstrap_retry_node(&tmp, run_id);
2636
2637 let mut ev = event(run_id);
2638 ev.seq = 3;
2639 ev.kind = "worker.exited".into();
2640 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2641 ev.data = serde_json::json!({ "exit_code": 0 });
2642 apply_event(&paths, &ev).expect("worker.exited applies");
2643
2644 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2645 .unwrap()
2646 .unwrap();
2647 let exit = n.worker_exit.expect("worker_exit recorded");
2648 assert_eq!(exit.code, Some(0));
2649 assert_eq!(exit.signal, None);
2650 assert!(exit.is_clean());
2651 assert_eq!(
2652 n.status,
2653 Status::Pending,
2654 "the exit fact never transitions status"
2655 );
2656 }
2657
2658 #[test]
2662 fn worker_exited_records_signal_and_is_first_write_wins() {
2663 let tmp = TempDir::new().unwrap();
2664 let run_id = "01jxsnap000000000000000000";
2665 let paths = bootstrap_retry_node(&tmp, run_id);
2666 let nid = NodeId::parse_str("n-0001").unwrap();
2667
2668 let mut ev = event(run_id);
2669 ev.seq = 3;
2670 ev.kind = "worker.exited".into();
2671 ev.node_id = Some(nid.clone());
2672 ev.data = serde_json::json!({ "signal": 9 });
2673 apply_event(&paths, &ev).expect("worker.exited applies");
2674
2675 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2676 let exit = n.worker_exit.expect("worker_exit recorded");
2677 assert_eq!(exit.signal, Some(9));
2678 assert!(exit.is_failure());
2679
2680 let mut dup = event(run_id);
2683 dup.seq = 4;
2684 dup.kind = "worker.exited".into();
2685 dup.node_id = Some(nid.clone());
2686 dup.data = serde_json::json!({ "exit_code": 0 });
2687 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
2688 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
2689 assert_eq!(
2690 n2.worker_exit.unwrap().signal,
2691 Some(9),
2692 "first-write-wins: the replayed exit must not overwrite the recorded fact"
2693 );
2694 }
2695
2696 #[test]
2700 fn worker_exited_without_code_or_signal_is_corrupt() {
2701 let tmp = TempDir::new().unwrap();
2702 let run_id = "01jxsnap000000000000000000";
2703 let paths = bootstrap_retry_node(&tmp, run_id);
2704
2705 let mut ev = event(run_id);
2706 ev.seq = 3;
2707 ev.kind = "worker.exited".into();
2708 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2709 ev.data = serde_json::json!({});
2710 match reduce_event_to_ops(&paths, &ev) {
2711 Err(Error::CorruptEventLog { .. }) => {}
2712 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
2713 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2714 }
2715
2716 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
2719 match reduce_event_to_ops(&paths, &ev) {
2720 Err(Error::CorruptEventLog { .. }) => {}
2721 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
2722 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
2723 }
2724 }
2725
2726 #[test]
2732 fn node_death_observed_records_first_death_first_write_wins() {
2733 let tmp = TempDir::new().unwrap();
2734 let run_id = "01jxsnap000000000000000000";
2735 let paths = bootstrap_retry_node(&tmp, run_id);
2736 let nid = NodeId::parse_str("n-0001").unwrap();
2737
2738 let mut ev = event(run_id);
2739 ev.seq = 3;
2740 ev.kind = "node.death_observed".into();
2741 ev.node_id = Some(nid.clone());
2742 ev.data = serde_json::json!({});
2743 apply_event(&paths, &ev).expect("node.death_observed applies");
2744 let first = read_node_opt(&paths, &nid)
2745 .unwrap()
2746 .unwrap()
2747 .first_death_at
2748 .expect("first_death_at recorded");
2749 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
2750
2751 let mut later = event(run_id);
2753 later.seq = 4;
2754 later.kind = "node.death_observed".into();
2755 later.node_id = Some(nid.clone());
2756 later.ts = ev.ts + chrono::Duration::seconds(30);
2757 later.data = serde_json::json!({});
2758 apply_event(&paths, &later).expect("re-observation applies as no-op");
2759 assert_eq!(
2760 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2761 Some(first),
2762 "first-write-wins: a re-observation must not reset the anchor"
2763 );
2764 }
2765
2766 #[test]
2771 fn node_death_observed_noop_when_worker_exit_present() {
2772 let tmp = TempDir::new().unwrap();
2773 let run_id = "01jxsnap000000000000000000";
2774 let paths = bootstrap_retry_node(&tmp, run_id);
2775 let nid = NodeId::parse_str("n-0001").unwrap();
2776
2777 let mut exit = event(run_id);
2779 exit.seq = 3;
2780 exit.kind = "worker.exited".into();
2781 exit.node_id = Some(nid.clone());
2782 exit.data = serde_json::json!({ "exit_code": 0 });
2783 apply_event(&paths, &exit).unwrap();
2784
2785 let mut death = event(run_id);
2787 death.seq = 4;
2788 death.kind = "node.death_observed".into();
2789 death.node_id = Some(nid.clone());
2790 death.data = serde_json::json!({});
2791 apply_event(&paths, &death).expect("applies as no-op");
2792 assert_eq!(
2793 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
2794 None,
2795 "a told worker.exited makes the crash backstop moot; no anchor recorded"
2796 );
2797 }
2798
2799 #[test]
2804 fn node_retry_clears_worker_exit() {
2805 let tmp = TempDir::new().unwrap();
2806 let run_id = "01jxsnap000000000000000000";
2807 let paths = bootstrap_retry_node(&tmp, run_id);
2808 let nid = NodeId::parse_str("n-0001").unwrap();
2809
2810 let mut exit = event(run_id);
2812 exit.seq = 3;
2813 exit.kind = "worker.exited".into();
2814 exit.node_id = Some(nid.clone());
2815 exit.data = serde_json::json!({ "exit_code": 7 });
2816 apply_event(&paths, &exit).unwrap();
2817 assert!(read_node_opt(&paths, &nid)
2818 .unwrap()
2819 .unwrap()
2820 .worker_exit
2821 .is_some());
2822
2823 let mut retry = event(run_id);
2825 retry.seq = 4;
2826 retry.kind = "node.retry".into();
2827 retry.node_id = Some(nid.clone());
2828 retry.data = serde_json::json!({
2829 "attempt": 1,
2830 "reason": "agent-died",
2831 "branch": "wt/foo",
2832 "worktree_path": "/tmp/new-wt",
2833 "agent_pid": 222,
2834 });
2835 apply_event(&paths, &retry).unwrap();
2836
2837 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2838 assert!(
2839 n.worker_exit.is_none(),
2840 "node.retry must clear the previous attempt's worker_exit"
2841 );
2842 assert_eq!(
2843 n.status,
2844 Status::Pending,
2845 "retry returns the node to Pending"
2846 );
2847 }
2848}