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 CallerPiLifecycle, CallerPiSession, CallerPiState, CallerSettlementIntent, ChildRef, Event,
59 EvidenceStatus, IdValidationError, Kind, Lifecycle, Manifest, MergeTxn, Node, NodeId, RunId,
60 Status, TmuxIdentity, WorkerEvidence, WorkerExit, STATE_SCHEMA_VERSION,
61};
62
63fn corrupt_id(events_path: &Path, ev: &Event, e: &IdValidationError) -> Error {
69 Error::CorruptEventLog {
70 path: events_path.to_path_buf(),
71 reason: format!("event seq={} kind={}: {e}", ev.seq, ev.kind),
72 }
73}
74
75fn opt_run_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<RunId>> {
79 match d.get(field) {
80 None | Some(Value::Null) => Ok(None),
81 Some(Value::String(s)) => RunId::parse_str(s)
82 .map(Some)
83 .map_err(|e| corrupt_id(events_path, ev, &e)),
84 Some(_) => Err(Error::CorruptEventLog {
85 path: events_path.to_path_buf(),
86 reason: format!(
87 "event seq={} kind={} `{field}` must be a JSON string or null",
88 ev.seq, ev.kind
89 ),
90 }),
91 }
92}
93
94fn opt_node_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<NodeId>> {
96 match d.get(field) {
97 None | Some(Value::Null) => Ok(None),
98 Some(Value::String(s)) => NodeId::parse_str(s)
99 .map(Some)
100 .map_err(|e| corrupt_id(events_path, ev, &e)),
101 Some(_) => Err(Error::CorruptEventLog {
102 path: events_path.to_path_buf(),
103 reason: format!(
104 "event seq={} kind={} `{field}` must be a JSON string or null",
105 ev.seq, ev.kind
106 ),
107 }),
108 }
109}
110
111fn data_kind(v: &Value) -> Option<Kind> {
125 match serde_json::from_value::<Kind>(v.clone()) {
126 Ok(Kind::Unknown) | Err(_) => None,
127 Ok(k) => Some(k),
128 }
129}
130
131fn data_status(v: &Value) -> Option<Status> {
132 serde_json::from_value(v.clone()).ok()
133}
134
135fn require_status(ev: &Event, path: PathBuf) -> Result<Status> {
136 data_status(ev.data.get("status").unwrap_or(&Value::Null)).ok_or_else(|| {
137 Error::CorruptEventLog {
138 path,
139 reason: format!("{} missing/invalid `status`", ev.kind),
140 }
141 })
142}
143
144fn want_str<'a>(events_path: &Path, ev: &Event, d: &'a Value, field: &str) -> Result<&'a str> {
145 d.get(field)
146 .and_then(Value::as_str)
147 .ok_or_else(|| Error::CorruptEventLog {
148 path: events_path.to_path_buf(),
149 reason: format!(
150 "event seq={} kind={} missing `{field}` string field",
151 ev.seq, ev.kind
152 ),
153 })
154}
155
156fn optional_bool(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<bool>> {
162 match d.get(field) {
163 None | Some(Value::Null) => Ok(None),
164 Some(Value::Bool(b)) => Ok(Some(*b)),
165 Some(_) => Err(Error::CorruptEventLog {
166 path: events_path.to_path_buf(),
167 reason: format!(
168 "event seq={} kind={} `{field}` must be a JSON boolean or null",
169 ev.seq, ev.kind
170 ),
171 }),
172 }
173}
174
175fn optional_i32(d: &Value, field: &str, events_path: &Path, ev: &Event) -> Result<Option<i32>> {
176 match d.get(field) {
177 None | Some(Value::Null) => Ok(None),
178 Some(v) => {
179 let raw = v.as_i64().ok_or_else(|| Error::CorruptEventLog {
180 path: events_path.to_path_buf(),
181 reason: format!(
182 "event seq={} kind={} `{field}` must be integer",
183 ev.seq, ev.kind
184 ),
185 })?;
186 i32::try_from(raw)
187 .map(Some)
188 .map_err(|_| Error::CorruptEventLog {
189 path: events_path.to_path_buf(),
190 reason: format!(
191 "event seq={} kind={} `{field}` out of i32 range: {raw}",
192 ev.seq, ev.kind
193 ),
194 })
195 }
196 }
197}
198
199fn optional_ts(
200 d: &Value,
201 field: &str,
202 events_path: &Path,
203 ev: &Event,
204) -> Result<Option<DateTime<Utc>>> {
205 match d.get(field) {
206 None | Some(Value::Null) => Ok(None),
207 Some(Value::String(s)) => DateTime::parse_from_rfc3339(s)
208 .map(|dt| Some(dt.with_timezone(&Utc)))
209 .map_err(|_| Error::CorruptEventLog {
210 path: events_path.to_path_buf(),
211 reason: format!(
212 "event seq={} kind={} `{field}` not RFC3339",
213 ev.seq, ev.kind
214 ),
215 }),
216 Some(_) => Err(Error::CorruptEventLog {
217 path: events_path.to_path_buf(),
218 reason: format!(
219 "event seq={} kind={} `{field}` must be RFC3339 string or null",
220 ev.seq, ev.kind
221 ),
222 }),
223 }
224}
225
226#[allow(clippy::large_enum_variant)]
238pub(crate) enum ProjectionOp {
239 Manifest(Manifest),
241 Node(Node),
243}
244
245pub(crate) fn commit_ops(paths: &RunPaths, ops: Vec<ProjectionOp>) -> Result<()> {
253 for op in ops {
254 match op {
255 ProjectionOp::Manifest(m) => write_manifest(paths, &m)?,
256 ProjectionOp::Node(n) => write_node(paths, &n)?,
257 }
258 }
259 Ok(())
260}
261
262pub(crate) fn reduce_event_to_ops(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
278 if ev.run_id != paths.run_id {
282 return Err(Error::CorruptEventLog {
283 path: paths.events(),
284 reason: format!(
285 "event seq={} envelope run_id {:?} does not match run {:?}",
286 ev.seq,
287 ev.run_id.as_str(),
288 paths.run_id.as_str()
289 ),
290 });
291 }
292 #[allow(clippy::match_same_arms)]
295 match ev.kind.as_str() {
296 "run.created" => reduce_run_created(paths, ev),
297 "run.status" => reduce_run_status(paths, ev),
298 "node.created" => reduce_node_created(paths, ev),
299 "node.status" => reduce_node_status(paths, ev),
300 "node.report" => reduce_node_report(paths, ev),
301 "node.retry" => reduce_node_retry(paths, ev),
302 "caller.pi.session_bound" => reduce_caller_pi_session_bound(paths, ev),
303 "caller.pi.lifecycle" => reduce_caller_pi_lifecycle(paths, ev),
304 "caller.settlement_intent" => reduce_caller_settlement_intent(paths, ev),
305 "worker.exited" => reduce_worker_exited(paths, ev),
306 "worker.evidence.archived" => reduce_worker_evidence_archived(paths, ev),
307 "worker.evidence.failed" => reduce_worker_evidence_failed(paths, ev),
308 "worker.display.retained" => reduce_worker_display_retained(paths, ev),
309 "worker.display.expired" => reduce_worker_display_expired(paths, ev),
310 "worker.display.unavailable" => reduce_worker_display_unavailable(paths, ev),
311 "node.death_observed" => reduce_node_death_observed(paths, ev),
312 "node.awaiting_input" => reduce_node_awaiting_input(paths, ev),
313 "node.input_resolved" => reduce_node_input_resolved(paths, ev),
314 KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
315 KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
316 "child.spawned" => reduce_child_spawned(paths, ev),
317 "supervisor.attached" => reduce_supervisor_attached(paths, ev),
318 "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
319 "supervisor.exited" => Ok(vec![]),
320 "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
328 "run.notified" | "run.awaiting_input_notified" => Ok(vec![]),
335 "cleanup.window_missing"
360 | "cleanup.worktree_missing"
361 | "cleanup.branch_remove_failed"
362 | "cleanup.branch_preserved"
363 | "cleanup.discard_authorized"
364 | "cleanup.session_killed"
365 | "cleanup.session_retained" => Ok(vec![]),
366 "supervisor.child_id_quarantined" => Ok(vec![]),
376 _ => Ok(vec![]),
377 }
378}
379
380fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
386 match op {
387 ProjectionOp::Manifest(_) => paths.manifest(),
388 ProjectionOp::Node(n) => paths.node(&n.node_id),
389 }
390}
391
392pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
414 let ops = reduce_event_to_ops(paths, event)?;
415 Ok(ops.iter().map(|op| op_path(paths, op)).collect())
416}
417
418pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
434 let ops = reduce_event_to_ops(paths, ev)?;
435 commit_ops(paths, ops)
436}
437
438pub fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
444 reduce_event_to_ops(paths, ev).map(|_| ())
445}
446
447fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
451 ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
452 path: events_path.to_path_buf(),
453 reason: format!(
454 "event seq={} kind={} missing top-level `node_id`",
455 ev.seq, ev.kind
456 ),
457 })
458}
459
460fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
461 if let Some(existing) = read_manifest_opt(paths)? {
465 if existing.run_id != ev.run_id {
466 return Err(Error::CorruptEventLog {
467 path: paths.manifest(),
468 reason: format!(
469 "run.created run_id={} conflicts with existing manifest run_id={}",
470 ev.run_id, existing.run_id
471 ),
472 });
473 }
474 return Ok(vec![]);
475 }
476 let events_path = paths.events();
477 let d = &ev.data;
478 let kind =
479 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
480 path: events_path.clone(),
481 reason: "run.created missing/invalid `kind`".into(),
482 })?;
483 let lifecycle: Lifecycle = serde_json::from_value(
484 d.get("lifecycle").cloned().unwrap_or(Value::Null),
485 )
486 .map_err(|_| Error::CorruptEventLog {
487 path: events_path.clone(),
488 reason: "run.created missing/invalid `lifecycle`".into(),
489 })?;
490 let title = want_str(&events_path, ev, d, "title")?.to_string();
491 let agent_selection: Option<crate::schema::AgentSelection> = d
492 .get("agent_selection")
493 .cloned()
494 .map(serde_json::from_value)
495 .transpose()
496 .map_err(|e| Error::CorruptEventLog {
497 path: events_path.clone(),
498 reason: format!("run.created invalid `agent_selection`: {e}"),
499 })?;
500 if let Some(selection) = &agent_selection {
501 selection
502 .validate()
503 .map_err(|reason| Error::CorruptEventLog {
504 path: events_path.clone(),
505 reason: format!("run.created invalid `agent_selection`: {reason}"),
506 })?;
507 }
508 let m = Manifest {
509 schema_version: STATE_SCHEMA_VERSION,
510 applied_seq: 0,
513 run_id: paths.run_id.clone(),
515 kind,
516 lifecycle,
517 caller_settlement_intent: None,
518 agent_owner: d
519 .get("agent_owner")
520 .cloned()
521 .map(serde_json::from_value)
522 .transpose()
523 .map_err(|e| Error::CorruptEventLog {
524 path: events_path.clone(),
525 reason: format!("run.created invalid agent_owner: {e}"),
526 })?
527 .unwrap_or_default(),
528 title,
529 status: Status::Pending,
530 created_at: ev.ts,
531 updated_at: ev.ts,
532 source_repo: d
533 .get("source_repo")
534 .and_then(Value::as_str)
535 .map(str::to_string),
536 source_branch: d
537 .get("source_branch")
538 .and_then(Value::as_str)
539 .map(str::to_string),
540 worktree_root: d
541 .get("worktree_root")
542 .and_then(Value::as_str)
543 .map(str::to_string),
544 managed_tmux_session: d
545 .get("managed_tmux_session")
546 .and_then(Value::as_str)
547 .map(str::to_string),
548 tmux_retention: retention_policy_from_data(&events_path, d)?,
549 notify_cmd: d
550 .get("notify_cmd")
551 .and_then(Value::as_str)
552 .map(str::to_string),
553 harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
554 agent_selection,
555 node_count: 0,
556 parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
557 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
558 };
559 Ok(vec![ProjectionOp::Manifest(m)])
560}
561
562fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
563 let mut m = match read_manifest_opt(paths)? {
564 Some(m) => m,
565 None => return Ok(vec![]),
566 };
567 let new_status = require_status(ev, paths.events())?;
568 if m.status.is_terminal() {
574 let recovery_seq = ev
575 .data
576 .get("recovery_merge_report_seq")
577 .and_then(Value::as_u64);
578 let children_successful =
579 ev.data.get("children_successful").and_then(Value::as_bool) == Some(true);
580 let recovery_allowed = m.status == Status::Failed
581 && new_status == Status::Done
582 && children_successful
583 && recovery_seq.is_some_and(|wanted| {
584 crate::cancel::read_node_status_facts(paths, Some(ev.seq))
585 .ok()
586 .filter(|facts| {
587 crate::aggregate_terminal_status(facts.iter().map(|fact| fact.status))
588 == Some(Status::Done)
589 })
590 .is_some_and(|facts| {
591 facts
592 .iter()
593 .any(|fact| fact.confirmed_merge_seq == Some(wanted))
594 })
595 });
596 if !recovery_allowed {
597 trace_terminal_noop(ev, m.status, new_status);
598 return Ok(vec![]);
599 }
600 }
601 if m.status == new_status {
602 return Ok(vec![]);
603 }
604 m.status = new_status;
605 m.updated_at = ev.ts;
606 Ok(vec![ProjectionOp::Manifest(m)])
607}
608
609fn retention_policy_from_data(
620 events_path: &Path,
621 data: &Value,
622) -> Result<Option<Box<crate::schema::TmuxRetentionPolicy>>> {
623 let Some(value) = data.get("tmux_retention") else {
624 return Ok(None);
625 };
626 let policy: crate::schema::TmuxRetentionPolicy = serde_json::from_value(value.clone())
627 .map_err(|e| Error::CorruptEventLog {
628 path: events_path.to_path_buf(),
629 reason: format!("run.created invalid `tmux_retention`: {e}"),
630 })?;
631 if !policy.persistent
632 || policy.completed_window_ttl_secs == 0
633 || policy.completed_window_max == 0
634 {
635 return Err(Error::CorruptEventLog {
636 path: events_path.to_path_buf(),
637 reason: "run.created tmux_retention must be persistent with positive ttl/max".into(),
638 });
639 }
640 Ok(Some(Box::new(policy)))
641}
642
643fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
644 let nonempty = |key| {
645 d.get(key)
646 .and_then(Value::as_str)
647 .map(str::trim)
648 .filter(|s| !s.is_empty())
649 .map(str::to_string)
650 };
651 let session = nonempty("tmux_session")?;
652 let window_id = nonempty("tmux_window_id")?;
653 Some(TmuxIdentity {
654 socket: nonempty("tmux_socket"),
655 session,
656 window_id,
657 pane_id: nonempty("tmux_pane_id"),
660 server_pid: d
661 .get("tmux_server_pid")
662 .and_then(Value::as_u64)
663 .and_then(|value| u32::try_from(value).ok()),
664 server_pid_start_secs: d.get("tmux_server_pid_start_secs").and_then(Value::as_u64),
665 server_marker: nonempty("tmux_server_marker"),
666 })
667}
668
669fn worker_evidence_from_spawn_data(
670 events_path: &Path,
671 ev: &Event,
672 data: &Value,
673) -> Result<Option<WorkerEvidence>> {
674 let Some(session_id) = data.get("pi_session_id") else {
675 if data.get("pi_session_path").is_some() || data.get("pi_session_cwd").is_some() {
676 return Err(Error::CorruptEventLog {
677 path: events_path.to_path_buf(),
678 reason: format!(
679 "event seq={} Pi evidence fields must be all-or-none",
680 ev.seq
681 ),
682 });
683 }
684 return Ok(None);
685 };
686 if session_id.is_null() {
687 if data
688 .get("pi_session_path")
689 .is_some_and(|value| !value.is_null())
690 || data
691 .get("pi_session_cwd")
692 .is_some_and(|value| !value.is_null())
693 {
694 return Err(Error::CorruptEventLog {
695 path: events_path.to_path_buf(),
696 reason: format!(
697 "event seq={} Pi evidence fields must be all-or-none",
698 ev.seq
699 ),
700 });
701 }
702 return Ok(None);
703 }
704 let session_id = session_id
705 .as_str()
706 .filter(|v| !v.is_empty())
707 .ok_or_else(|| Error::CorruptEventLog {
708 path: events_path.to_path_buf(),
709 reason: format!(
710 "event seq={} pi_session_id must be a non-empty string",
711 ev.seq
712 ),
713 })?;
714 if session_id.len() != 36
715 || !session_id.bytes().enumerate().all(|(index, byte)| {
716 if matches!(index, 8 | 13 | 18 | 23) {
717 byte == b'-'
718 } else {
719 byte.is_ascii_hexdigit()
720 }
721 })
722 {
723 return Err(Error::CorruptEventLog {
724 path: events_path.to_path_buf(),
725 reason: format!("event seq={} pi_session_id must be a UUID", ev.seq),
726 });
727 }
728 let original_cwd = want_str(events_path, ev, data, "pi_session_cwd")?;
729 let live_session_path = want_str(events_path, ev, data, "pi_session_path")?;
730 let expected_live_path = format!(
731 ".creating/pi-sessions/{}/pi-session-{session_id}.jsonl",
732 ev.run_id.as_str()
733 );
734 if live_session_path != expected_live_path {
735 return Err(Error::CorruptEventLog {
736 path: events_path.to_path_buf(),
737 reason: format!(
738 "event seq={} pi_session_path is not the canonical state-relative path",
739 ev.seq
740 ),
741 });
742 }
743 let attempt = match data.get("attempt") {
744 None | Some(Value::Null) => 0,
745 Some(value) => value
746 .as_u64()
747 .and_then(|raw| u32::try_from(raw).ok())
748 .ok_or_else(|| Error::CorruptEventLog {
749 path: events_path.to_path_buf(),
750 reason: format!("event seq={} attempt must be a u32", ev.seq),
751 })?,
752 };
753 Ok(Some(WorkerEvidence {
754 attempt,
755 session_id: session_id.to_string(),
756 original_cwd: original_cwd.to_string(),
757 live_session_path: live_session_path.to_string(),
758 status: EvidenceStatus::Pending,
759 transcript_path: None,
760 resume_path: None,
761 pane_path: None,
762 report_path: None,
763 transcript_sha256: None,
764 error: None,
765 }))
766}
767
768fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
769 let events_path = paths.events();
770 let node_id = require_envelope_node_id(&events_path, ev)?;
773 let is_default_node = node_id.as_str() == "n-0001";
774 if read_node_opt(paths, &node_id)?.is_some() {
776 return Ok(vec![]);
777 }
778 let d = &ev.data;
779 let kind =
780 data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
781 path: events_path.clone(),
782 reason: format!(
783 "event seq={} kind=node.created missing/invalid `kind`",
784 ev.seq
785 ),
786 })?;
787 let n = Node {
788 schema_version: STATE_SCHEMA_VERSION,
789 node_id,
790 run_id: paths.run_id.clone(),
792 parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
793 kind,
794 status: Status::Pending,
795 task: d.get("task").and_then(Value::as_str).map(str::to_string),
796 worktree_path: d
797 .get("worktree_path")
798 .and_then(Value::as_str)
799 .map(str::to_string),
800 branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
801 base_sha: d
802 .get("base_sha")
803 .and_then(Value::as_str)
804 .filter(|s| !s.is_empty())
805 .map(str::to_string),
806 tmux_window: d
807 .get("tmux_window")
808 .and_then(Value::as_str)
809 .map(str::to_string),
810 tmux_identity: tmux_identity_from_data(d).map(Box::new),
811 evidence: worker_evidence_from_spawn_data(&events_path, ev, d)?,
812 caller_pi_session: None,
813 caller_pi_lifecycle: None,
814 retained_display: None,
815 retention_unavailable: None,
816 agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
817 agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
818 supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
819 children: Vec::new(),
820 started_at: Some(ev.ts),
821 updated_at: ev.ts,
822 last_report: None,
823 last_processed_report_seq_by_child: serde_json::Map::default(),
824 retry_attempts: 0,
825 worker_exit: None,
826 pending_merge: None,
827 first_death_at: None,
828 awaiting_input: None,
829 };
830 let mut ops = vec![ProjectionOp::Node(n)];
831 if let Some(mut m) = read_manifest_opt(paths)? {
832 if is_default_node && m.source_branch.is_none() {
837 m.source_branch = d
838 .get("source_branch")
839 .and_then(Value::as_str)
840 .filter(|branch| !branch.is_empty())
841 .map(str::to_string);
842 }
843 m.updated_at = ev.ts;
848 ops.push(ProjectionOp::Manifest(m));
849 }
850 Ok(ops)
851}
852
853fn reduce_caller_settlement_intent(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
856 let bad = |reason: &str| Error::CorruptEventLog {
857 path: paths.events(),
858 reason: format!("event seq={} caller.settlement_intent: {reason}", ev.seq),
859 };
860 let id = require_envelope_node_id(&paths.events(), ev)?;
861 let mut manifest = read_manifest_opt(paths)?.ok_or_else(|| bad("missing run"))?;
862 let node = read_node_opt(paths, &id)?.ok_or_else(|| bad("missing node"))?;
863 let mut intent: CallerSettlementIntent =
864 serde_json::from_value(ev.data.clone()).map_err(|_| bad("invalid typed intent"))?;
865 if manifest.agent_owner != crate::schema::AgentOwner::Caller
866 || intent.run_id != ev.run_id
867 || intent.node_id != id
868 || node.run_id != ev.run_id
869 || ev.idempotency_key.as_deref() != Some(intent.key.as_str())
870 || intent.key.trim().is_empty()
871 || intent.actor.trim().is_empty()
872 || intent.writer_dev == 0
873 || intent.writer_ino == 0
874 || intent.gate_dev == 0
875 || intent.gate_ino == 0
876 || intent.seq != 0
877 || intent.generation
878 != node
879 .caller_pi_lifecycle
880 .as_ref()
881 .map_or(0, |v| v.generation)
882 {
883 return Err(bad("invalid identity, generation or audit fields"));
884 }
885 intent.seq = ev.seq;
886 if let Some(existing) = &manifest.caller_settlement_intent {
887 return if existing == &intent {
888 Ok(vec![])
889 } else {
890 Err(bad("intent already recorded"))
891 };
892 }
893 if manifest.status.is_terminal()
894 && !(manifest.status == Status::Failed
895 && node.status == Status::Failed
896 && manifest.node_count == 1
897 && id.as_str() == "n-0001"
898 && matches!(
899 intent.operation,
900 crate::schema::SettlementOperation::Merge
901 | crate::schema::SettlementOperation::Discard
902 ))
903 {
904 return Err(bad("terminal run"));
905 }
906 manifest.caller_settlement_intent = Some(intent);
907 manifest.updated_at = ev.ts;
908 Ok(vec![ProjectionOp::Manifest(manifest)])
909}
910
911fn reduce_caller_pi_session_bound(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
914 let bad = |reason: &str| Error::CorruptEventLog {
915 path: paths.events(),
916 reason: format!("event seq={} caller.pi.session_bound: {reason}", ev.seq),
917 };
918 let id = require_envelope_node_id(&paths.events(), ev)?;
919 let manifest = read_manifest_opt(paths)?.ok_or_else(|| bad("missing run"))?;
920 if manifest.agent_owner != crate::schema::AgentOwner::Caller {
921 return Err(bad("not a caller-owned run"));
922 }
923 let mut node = read_node_opt(paths, &id)?.ok_or_else(|| bad("missing node"))?;
924 let binding: CallerPiSession =
925 serde_json::from_value(ev.data.clone()).map_err(|_| bad("invalid binding"))?;
926 if binding.original_cwd != node.worktree_path.as_deref().unwrap_or("")
927 || binding.original_cwd.is_empty()
928 || binding.pi_session_id.is_empty()
929 || binding.session_path.is_empty()
930 {
931 return Err(bad("binding identity does not match node"));
932 }
933 match &node.caller_pi_session {
934 Some(existing) if existing == &binding => return Ok(vec![]),
935 Some(_) => return Err(bad("binding already exists with different identity")),
936 None => {}
937 }
938 if let Some(reserved) = &node.caller_pi_lifecycle {
939 if reserved.pi_session_id != binding.pi_session_id
940 || Some(reserved.generation) != binding.generation
941 || reserved
942 .session_path
943 .as_deref()
944 .is_some_and(|p| p != binding.session_path)
945 || !matches!(
946 reserved.state,
947 CallerPiState::Reserved | CallerPiState::Started
948 )
949 {
950 return Err(bad("binding does not match the current launch reservation"));
951 }
952 }
953 if binding.generation.is_some() && node.caller_pi_lifecycle.is_none() {
954 return Err(bad("binding generation has no reservation"));
955 }
956 if let Some(current) = &mut node.caller_pi_lifecycle {
957 current.session_path = Some(binding.session_path.clone());
958 }
959 node.caller_pi_session = Some(binding);
960 node.updated_at = ev.ts;
961 Ok(vec![ProjectionOp::Node(node)])
962}
963
964fn reduce_caller_pi_lifecycle(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
968 let bad = |reason: &str| Error::CorruptEventLog {
969 path: paths.events(),
970 reason: format!("event seq={} caller.pi.lifecycle: {reason}", ev.seq),
971 };
972 let id = require_envelope_node_id(&paths.events(), ev)?;
973 let manifest = read_manifest_opt(paths)?.ok_or_else(|| bad("missing run"))?;
974 if manifest.agent_owner != crate::schema::AgentOwner::Caller {
975 return Err(bad("not a caller-owned run"));
976 }
977 let mut node = read_node_opt(paths, &id)?.ok_or_else(|| bad("missing node"))?;
978 let fact: CallerPiLifecycle =
979 serde_json::from_value(ev.data.clone()).map_err(|_| bad("invalid lifecycle fact"))?;
980 let binding = node.caller_pi_session.as_ref();
981 if fact.generation == 0
982 || binding.is_some_and(|b| {
983 fact.pi_session_id != b.pi_session_id
984 || fact.session_path.as_deref() != Some(&b.session_path)
985 })
986 || (fact.state == CallerPiState::Reserved && fact.session_path.is_some())
987 || (fact.state == CallerPiState::Started && fact.session_path.is_none())
988 || (binding.is_none()
989 && fact.state != CallerPiState::Reserved
990 && !node
991 .caller_pi_lifecycle
992 .as_ref()
993 .is_some_and(|old| old.state == CallerPiState::Reserved))
994 || fact
995 .reason
996 .as_ref()
997 .is_some_and(|r| r.trim().is_empty() || r.len() > 1024)
998 || matches!(fact.state, CallerPiState::Started | CallerPiState::Reserved)
999 != fact.reason.is_none()
1000 {
1001 return Err(bad("invalid generation, session identity or reason"));
1002 }
1003 if node.caller_pi_lifecycle.as_ref() == Some(&fact) {
1004 return Ok(vec![]);
1005 }
1006 if manifest.status.is_terminal() {
1007 return Err(bad("terminal run cannot accept a new Pi transition"));
1008 }
1009 match &node.caller_pi_lifecycle {
1010 None if fact.generation == 1
1011 && fact.state == CallerPiState::Started
1012 && binding.is_some() => {}
1013 None if fact.generation == 1
1014 && fact.state == CallerPiState::Reserved
1015 && binding.is_none() => {}
1016 Some(old) if *old == fact => return Ok(vec![]),
1017 Some(old)
1018 if old.generation == fact.generation
1019 && old.state == CallerPiState::Started
1020 && matches!(
1021 fact.state,
1022 CallerPiState::Exited | CallerPiState::ControlUncertain
1023 ) => {}
1024 Some(old)
1025 if old.generation == fact.generation
1026 && old.state == CallerPiState::Reserved
1027 && old.pi_session_id == fact.pi_session_id
1028 && (old.session_path == fact.session_path
1029 || (old.session_path.is_none()
1030 && fact.state == CallerPiState::Started
1031 && binding.is_some()))
1032 && (matches!(
1033 fact.state,
1034 CallerPiState::LaunchFailed | CallerPiState::ControlUncertain
1035 ) || (fact.state == CallerPiState::Started && binding.is_some())) => {}
1036 Some(old)
1037 if old.generation.checked_add(1) == Some(fact.generation)
1038 && matches!(
1039 old.state,
1040 CallerPiState::Exited | CallerPiState::LaunchFailed
1041 )
1042 && fact.state == CallerPiState::Started => {}
1043 _ => return Err(bad("stale or conflicting generation/transition")),
1044 }
1045 node.caller_pi_lifecycle = Some(fact);
1046 node.updated_at = ev.ts;
1047 Ok(vec![ProjectionOp::Node(node)])
1048}
1049
1050fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1069 let events_path = paths.events();
1070 let node_id = require_envelope_node_id(&events_path, ev)?;
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() {
1078 tracing::debug!(
1079 target: "taskfleet_core::reducer",
1080 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
1081 "no-op: node.retry against terminal node"
1082 );
1083 return Ok(vec![]);
1084 }
1085 let d = &ev.data;
1086 n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
1089 n.base_sha = d
1090 .get("base_sha")
1091 .and_then(Value::as_str)
1092 .filter(|s| !s.is_empty())
1093 .map(str::to_string);
1094 n.worktree_path = d
1095 .get("worktree_path")
1096 .and_then(Value::as_str)
1097 .map(str::to_string);
1098 n.tmux_window = d
1099 .get("tmux_window")
1100 .and_then(Value::as_str)
1101 .map(str::to_string);
1102 n.tmux_identity = tmux_identity_from_data(d).map(Box::new);
1103 n.evidence = worker_evidence_from_spawn_data(&events_path, ev, d)?;
1104 n.retained_display = None;
1105 n.retention_unavailable = None;
1106 n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
1107 n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
1108 n.status = Status::Pending;
1109 n.started_at = Some(ev.ts);
1110 n.updated_at = ev.ts;
1111 n.last_report = None;
1112 n.pending_merge = None;
1118 n.worker_exit = None;
1123 n.first_death_at = None;
1128 n.awaiting_input = None;
1131 n.retry_attempts = d
1139 .get("attempt")
1140 .and_then(Value::as_u64)
1141 .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
1142 Ok(vec![ProjectionOp::Node(n)])
1143}
1144
1145fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1146 let events_path = paths.events();
1147 let node_id = require_envelope_node_id(&events_path, ev)?;
1148 let mut n = match read_node_opt(paths, &node_id)? {
1149 Some(n) => n,
1150 None => return Ok(vec![]),
1151 };
1152 let new_status = require_status(ev, events_path)?;
1153 if n.status.is_terminal() {
1156 trace_terminal_noop(ev, n.status, new_status);
1157 return Ok(vec![]);
1158 }
1159 if n.status == new_status {
1160 return Ok(vec![]);
1161 }
1162 n.status = new_status;
1163 if new_status.is_terminal() {
1169 n.pending_merge = None;
1170 n.awaiting_input = None;
1171 }
1172 n.updated_at = ev.ts;
1173 Ok(vec![ProjectionOp::Node(n)])
1174}
1175
1176fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1177 let events_path = paths.events();
1178 let node_id = require_envelope_node_id(&events_path, ev)?;
1179 let mut n = match read_node_opt(paths, &node_id)? {
1180 Some(n) => n,
1181 None => return Ok(vec![]),
1182 };
1183 if n.status.is_terminal() {
1195 if ReportOrigin::permits_terminal_merge_recovery(n.status, &ev.data) {
1228 if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
1229 return Ok(vec![]);
1230 }
1231 tracing::info!(
1232 target: "taskfleet_core::reducer",
1233 seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
1234 "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
1235 );
1236 n.last_report = Some(ev.data.clone());
1237 n.status = Status::Done;
1241 n.pending_merge = None;
1245 n.awaiting_input = None;
1246 n.updated_at = ev.ts;
1247 return Ok(vec![ProjectionOp::Node(n)]);
1248 }
1249 tracing::debug!(
1250 target: "taskfleet_core::reducer",
1251 seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
1252 "no-op: node.report against terminal node"
1253 );
1254 return Ok(vec![]);
1255 }
1256 let new_status = report_terminal_status(&events_path, ev)?;
1264 n.last_report = Some(ev.data.clone());
1265 n.status = new_status;
1266 n.awaiting_input = None;
1270 n.pending_merge = None;
1275 n.updated_at = ev.ts;
1276 Ok(vec![ProjectionOp::Node(n)])
1277}
1278
1279fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1294 let events_path = paths.events();
1295 let node_id = require_envelope_node_id(&events_path, ev)?;
1296 let attempt = match ev.data.get("attempt") {
1297 None | Some(Value::Null) => 0,
1298 Some(_) => evidence_attempt(&events_path, ev)?,
1299 };
1300 let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
1301 let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
1302 match (code, signal) {
1307 (Some(_), None) | (None, Some(_)) => {}
1308 _ => {
1309 return Err(Error::CorruptEventLog {
1310 path: events_path,
1311 reason: format!(
1312 "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
1313 ev.seq
1314 ),
1315 });
1316 }
1317 }
1318 let mut n = match read_node_opt(paths, &node_id)? {
1319 Some(n) => n,
1320 None => return Ok(vec![]),
1327 };
1328 if attempt != n.retry_attempts || n.worker_exit.is_some() {
1332 return Ok(vec![]);
1333 }
1334 n.worker_exit = Some(WorkerExit {
1335 code,
1336 signal,
1337 at: ev.ts,
1338 });
1339 n.awaiting_input = None;
1343 n.updated_at = ev.ts;
1344 Ok(vec![ProjectionOp::Node(n)])
1345}
1346
1347fn evidence_attempt(events_path: &Path, ev: &Event) -> Result<u32> {
1348 ev.data
1349 .get("attempt")
1350 .and_then(Value::as_u64)
1351 .and_then(|raw| u32::try_from(raw).ok())
1352 .ok_or_else(|| Error::CorruptEventLog {
1353 path: events_path.to_path_buf(),
1354 reason: format!("event seq={} attempt must be a u32", ev.seq),
1355 })
1356}
1357
1358fn reduce_worker_evidence_archived(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1359 let events_path = paths.events();
1360 let node_id = require_envelope_node_id(&events_path, ev)?;
1361 let attempt = evidence_attempt(&events_path, ev)?;
1362 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1363 let mut values = Vec::new();
1364 for field in [
1365 "transcript_path",
1366 "resume_path",
1367 "pane_path",
1368 "report_path",
1369 "transcript_sha256",
1370 ] {
1371 values.push(want_str(&events_path, ev, &ev.data, field)?.to_string());
1372 }
1373 let expected_prefix = format!("evidence/{}/", node_id.as_str());
1374 for (index, suffix) in [
1375 "pi-session.original.jsonl",
1376 "pi-session.resume.jsonl",
1377 "final-pane.log",
1378 "terminal-report.json",
1379 ]
1380 .iter()
1381 .enumerate()
1382 {
1383 if values[index] != format!("{expected_prefix}{suffix}") {
1384 return Err(Error::CorruptEventLog {
1385 path: events_path,
1386 reason: format!(
1387 "event seq={} evidence artifact path is not canonical",
1388 ev.seq
1389 ),
1390 });
1391 }
1392 }
1393 if values[4].len() != 64
1394 || !values[4]
1395 .bytes()
1396 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1397 {
1398 return Err(Error::CorruptEventLog {
1399 path: events_path,
1400 reason: format!(
1401 "event seq={} transcript_sha256 is not lowercase SHA-256",
1402 ev.seq
1403 ),
1404 });
1405 }
1406 let mut node = match read_node_opt(paths, &node_id)? {
1407 Some(node) => node,
1408 None => return Ok(vec![]),
1409 };
1410 let Some(evidence) = node.evidence.as_mut() else {
1411 return Err(Error::CorruptEventLog {
1412 path: events_path,
1413 reason: format!(
1414 "event seq={} evidence archive has no recorded Pi session",
1415 ev.seq
1416 ),
1417 });
1418 };
1419 if evidence.attempt != attempt || evidence.session_id != session_id {
1420 return Ok(vec![]);
1421 }
1422 if evidence.status == EvidenceStatus::Complete {
1423 return Ok(vec![]);
1424 }
1425 evidence.transcript_path = Some(values.remove(0));
1426 evidence.resume_path = Some(values.remove(0));
1427 evidence.pane_path = Some(values.remove(0));
1428 evidence.report_path = Some(values.remove(0));
1429 evidence.transcript_sha256 = Some(values.remove(0));
1430 evidence.status = EvidenceStatus::Complete;
1431 evidence.error = None;
1432 node.updated_at = ev.ts;
1433 Ok(vec![ProjectionOp::Node(node)])
1434}
1435
1436fn reduce_worker_display_retained(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1437 let events_path = paths.events();
1438 let node_id = require_envelope_node_id(&events_path, ev)?;
1439 let display: crate::schema::RetainedDisplay =
1440 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1441 path: events_path.clone(),
1442 reason: format!("event seq={} invalid retained display: {e}", ev.seq),
1443 })?;
1444 let mut node = match read_node_opt(paths, &node_id)? {
1445 Some(node) => node,
1446 None => return Ok(vec![]),
1447 };
1448 if display.attempt != node.retry_attempts || display.ownership_marker.is_empty() {
1449 return Ok(vec![]);
1450 }
1451 if node.retained_display.is_some() {
1452 return Ok(vec![]);
1453 }
1454 node.retained_display = Some(Box::new(display));
1455 node.updated_at = ev.ts;
1456 Ok(vec![ProjectionOp::Node(node)])
1457}
1458
1459fn reduce_worker_display_expired(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1460 let events_path = paths.events();
1461 let node_id = require_envelope_node_id(&events_path, ev)?;
1462 let marker = want_str(&events_path, ev, &ev.data, "ownership_marker")?;
1463 let attempt = evidence_attempt(&events_path, ev)?;
1464 let mut node = match read_node_opt(paths, &node_id)? {
1465 Some(node) => node,
1466 None => return Ok(vec![]),
1467 };
1468 let Some(display) = node.retained_display.as_mut() else {
1469 return Ok(vec![]);
1470 };
1471 if display.attempt != attempt
1472 || display.ownership_marker != marker
1473 || display.expired_at.is_some()
1474 {
1475 return Ok(vec![]);
1476 }
1477 display.expired_at = Some(ev.ts);
1478 node.updated_at = ev.ts;
1479 Ok(vec![ProjectionOp::Node(node)])
1480}
1481
1482fn reduce_worker_display_unavailable(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1483 let events_path = paths.events();
1484 let node_id = require_envelope_node_id(&events_path, ev)?;
1485 let attempt = evidence_attempt(&events_path, ev)?;
1486 let reason = want_str(&events_path, ev, &ev.data, "reason")?;
1487 let mut node = match read_node_opt(paths, &node_id)? {
1488 Some(node) => node,
1489 None => return Ok(vec![]),
1490 };
1491 if attempt != node.retry_attempts || node.retained_display.is_some() {
1492 return Ok(vec![]);
1493 }
1494 if node.retention_unavailable.is_none() {
1495 node.retention_unavailable = Some(reason.to_string());
1496 node.updated_at = ev.ts;
1497 }
1498 Ok(vec![ProjectionOp::Node(node)])
1499}
1500
1501fn reduce_worker_evidence_failed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1502 let events_path = paths.events();
1503 let node_id = require_envelope_node_id(&events_path, ev)?;
1504 let attempt = evidence_attempt(&events_path, ev)?;
1505 let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1506 let detail = want_str(&events_path, ev, &ev.data, "error")?.to_string();
1507 let mut node = match read_node_opt(paths, &node_id)? {
1508 Some(node) => node,
1509 None => return Ok(vec![]),
1510 };
1511 let Some(evidence) = node.evidence.as_mut() else {
1512 return Err(Error::CorruptEventLog {
1513 path: events_path,
1514 reason: format!(
1515 "event seq={} evidence failure has no recorded Pi session",
1516 ev.seq
1517 ),
1518 });
1519 };
1520 if evidence.attempt != attempt || evidence.session_id != session_id {
1521 return Ok(vec![]);
1522 }
1523 if evidence.status == EvidenceStatus::Complete {
1524 return Ok(vec![]);
1525 }
1526 evidence.status = EvidenceStatus::Failed;
1527 evidence.error = Some(detail);
1528 node.updated_at = ev.ts;
1529 Ok(vec![ProjectionOp::Node(node)])
1530}
1531
1532fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1545 let events_path = paths.events();
1546 let node_id = require_envelope_node_id(&events_path, ev)?;
1547 let mut n = match read_node_opt(paths, &node_id)? {
1548 Some(n) => n,
1549 None => return Ok(vec![]),
1550 };
1551 if n.first_death_at.is_some()
1557 || n.status.is_terminal()
1558 || n.worker_exit.is_some()
1559 || n.last_report.is_some()
1560 || n.pending_merge.is_some()
1561 {
1562 return Ok(vec![]);
1563 }
1564 n.first_death_at = Some(ev.ts);
1565 n.updated_at = ev.ts;
1566 Ok(vec![ProjectionOp::Node(n)])
1567}
1568
1569fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1578 const MAX_ITEMS: usize = 8;
1579 const MAX_TOPIC_CHARS: usize = 512;
1580 const MAX_OPTIONS: usize = 16;
1581 const MAX_OPTION_CHARS: usize = 256;
1582
1583 let events_path = paths.events();
1584 let node_id = require_envelope_node_id(&events_path, ev)?;
1585 let items = ev
1588 .data
1589 .get("discussion_items")
1590 .and_then(Value::as_array)
1591 .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1592 .ok_or_else(|| Error::CorruptEventLog {
1593 path: events_path.clone(),
1594 reason: format!(
1595 "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1596 ev.seq
1597 ),
1598 })?;
1599 for (index, item) in items.iter().enumerate() {
1600 let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1601 path: events_path.clone(),
1602 reason: format!(
1603 "event seq={} discussion_items[{index}] must be an object",
1604 ev.seq
1605 ),
1606 })?;
1607 let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1608 let default = obj
1609 .get("recommended_default")
1610 .and_then(Value::as_str)
1611 .unwrap_or("");
1612 let options = obj.get("options").and_then(Value::as_array);
1613 let options_valid = options.is_some_and(|values| {
1614 !values.is_empty()
1615 && values.len() <= MAX_OPTIONS
1616 && values.iter().all(|v| {
1617 v.as_str().is_some_and(|s| {
1618 !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1619 })
1620 })
1621 && values.iter().any(|v| v.as_str() == Some(default))
1622 });
1623 if topic.trim().is_empty()
1624 || topic.chars().count() > MAX_TOPIC_CHARS
1625 || default.trim().is_empty()
1626 || !options_valid
1627 {
1628 return Err(Error::CorruptEventLog {
1629 path: events_path.clone(),
1630 reason: format!(
1631 "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1632 ev.seq
1633 ),
1634 });
1635 }
1636 }
1637
1638 let mut n = match read_node_opt(paths, &node_id)? {
1639 Some(n) => n,
1640 None => return Ok(vec![]),
1641 };
1642 if n.status.is_terminal() || n.worker_exit.is_some() {
1643 return Ok(vec![]);
1644 }
1645 if let Some(open) = n.awaiting_input.as_mut() {
1646 if open.discussion_items.len() + items.len() > MAX_ITEMS {
1649 return Err(Error::CorruptEventLog {
1650 path: events_path,
1651 reason: format!(
1652 "event seq={} would exceed {MAX_ITEMS} open discussion items",
1653 ev.seq
1654 ),
1655 });
1656 }
1657 open.discussion_items.extend(items.iter().cloned());
1658 } else {
1659 n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1660 opened_at: ev.ts,
1661 event_seq: ev.seq,
1662 discussion_items: items.clone(),
1663 }));
1664 }
1665 n.updated_at = ev.ts;
1666 Ok(vec![ProjectionOp::Node(n)])
1667}
1668
1669fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1673 let events_path = paths.events();
1674 let node_id = require_envelope_node_id(&events_path, ev)?;
1675 let seq = ev
1678 .data
1679 .get("event_seq")
1680 .and_then(Value::as_u64)
1681 .ok_or_else(|| Error::CorruptEventLog {
1682 path: events_path.clone(),
1683 reason: format!(
1684 "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1685 ev.seq
1686 ),
1687 })?;
1688 let mut n = match read_node_opt(paths, &node_id)? {
1689 Some(n) => n,
1690 None => return Ok(vec![]),
1691 };
1692 let Some(open) = n.awaiting_input.as_ref() else {
1693 return Ok(vec![]);
1694 };
1695 if seq != open.event_seq {
1696 return Ok(vec![]);
1697 }
1698 n.awaiting_input = None;
1699 n.updated_at = ev.ts;
1700 Ok(vec![ProjectionOp::Node(n)])
1701}
1702
1703pub const KIND_MERGE_STARTED: &str = "merge.started";
1707
1708pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1713
1714fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1728 let events_path = paths.events();
1729 let node_id = require_envelope_node_id(&events_path, ev)?;
1730 let txn: MergeTxn =
1731 serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1732 path: events_path.clone(),
1733 reason: format!(
1734 "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1735 ev.seq
1736 ),
1737 })?;
1738 let manifest = read_manifest_opt(paths)?.ok_or_else(|| Error::CorruptEventLog {
1739 path: events_path.clone(),
1740 reason: "merge.started without a run".into(),
1741 })?;
1742 let bad_link = || Error::CorruptEventLog {
1743 path: events_path.clone(),
1744 reason: format!(
1745 "event seq={} merge.started caller authority does not match intent",
1746 ev.seq
1747 ),
1748 };
1749 match manifest.agent_owner {
1750 crate::schema::AgentOwner::Caller => {
1751 let intent = manifest
1752 .caller_settlement_intent
1753 .as_ref()
1754 .ok_or_else(bad_link)?;
1755 let link = txn.caller_authority.as_ref().ok_or_else(bad_link)?;
1756 if intent.operation != crate::schema::SettlementOperation::Merge
1757 || intent.run_id != ev.run_id
1758 || intent.node_id != node_id
1759 || link.intent_seq != intent.seq
1760 || link.intent_key != intent.key
1761 || (link.writer_dev, link.writer_ino) != (intent.writer_dev, intent.writer_ino)
1762 || txn.source_branch != manifest.source_branch.as_deref().unwrap_or("")
1763 {
1764 return Err(bad_link());
1765 }
1766 }
1767 crate::schema::AgentOwner::Taskfleet if txn.caller_authority.is_some() => {
1768 return Err(bad_link())
1769 }
1770 crate::schema::AgentOwner::Taskfleet => {}
1771 }
1772 let mut n = match read_node_opt(paths, &node_id)? {
1773 Some(n) => n,
1774 None => return Ok(vec![]),
1775 };
1776 if manifest.agent_owner == crate::schema::AgentOwner::Caller
1777 && (n.branch.as_deref() != Some(&txn.worker_branch)
1778 || n.caller_pi_lifecycle.as_ref().map_or(0, |p| p.generation)
1779 != manifest
1780 .caller_settlement_intent
1781 .as_ref()
1782 .unwrap()
1783 .generation)
1784 {
1785 return Err(bad_link());
1786 }
1787 if n.status.is_terminal()
1791 && !(manifest.agent_owner == crate::schema::AgentOwner::Caller
1792 && n.status == Status::Failed
1793 && manifest.status == Status::Failed
1794 && manifest.node_count == 1
1795 && manifest.caller_settlement_intent.as_ref().is_some_and(|i| {
1796 i.operation == crate::schema::SettlementOperation::Merge && i.node_id == node_id
1797 }))
1798 {
1799 return Ok(vec![]);
1800 }
1801 if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1804 return Ok(vec![]);
1805 }
1806 n.pending_merge = Some(Box::new(txn));
1807 n.updated_at = ev.ts;
1808 Ok(vec![ProjectionOp::Node(n)])
1809}
1810
1811fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1823 let events_path = paths.events();
1824 let node_id = require_envelope_node_id(&events_path, ev)?;
1825 let op_id = ev
1826 .data
1827 .get("op_id")
1828 .and_then(Value::as_str)
1829 .ok_or_else(|| Error::CorruptEventLog {
1830 path: events_path.clone(),
1831 reason: format!(
1832 "event seq={} kind=merge.aborted is missing string `op_id`",
1833 ev.seq
1834 ),
1835 })?;
1836 let mut n = match read_node_opt(paths, &node_id)? {
1837 Some(n) => n,
1838 None => return Ok(vec![]),
1839 };
1840 match n.pending_merge.as_ref() {
1843 Some(t) if t.op_id == op_id => {}
1844 _ => return Ok(vec![]),
1845 }
1846 n.pending_merge = None;
1847 n.updated_at = ev.ts;
1848 Ok(vec![ProjectionOp::Node(n)])
1849}
1850
1851fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1859 if current == incoming {
1860 tracing::debug!(
1861 target: "taskfleet_core::reducer",
1862 seq = ev.seq, kind = %ev.kind, status = ?current,
1863 "no-op: status re-applied to terminal target"
1864 );
1865 } else {
1866 tracing::warn!(
1867 target: "taskfleet_core::reducer",
1868 seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1869 "no-op: ignored conflicting transition from terminal target"
1870 );
1871 }
1872}
1873
1874fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1883 let corrupt = |reason: String| Error::CorruptEventLog {
1884 path: events_path.to_path_buf(),
1885 reason,
1886 };
1887 let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1888 let success = optional_bool(events_path, ev, &ev.data, "success")?;
1889 if cancelled {
1890 if success == Some(true) {
1891 return Err(corrupt(format!(
1892 "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1893 ev.seq
1894 )));
1895 }
1896 Ok(Status::Cancelled)
1897 } else {
1898 match success {
1899 Some(true) => Ok(Status::Done),
1900 Some(false) => Ok(Status::Failed),
1901 None => Err(corrupt(format!(
1902 "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1903 ev.seq
1904 ))),
1905 }
1906 }
1907}
1908
1909fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1910 let events_path = paths.events();
1913 let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1914 path: events_path.clone(),
1915 reason: format!(
1916 "event seq={} kind=child.spawned missing parent `node_id`",
1917 ev.seq
1918 ),
1919 })?;
1920 let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1921 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1922 let child_node_id = NodeId::parse_str(
1923 ev.data
1924 .get("child_node_id")
1925 .and_then(Value::as_str)
1926 .unwrap_or("n-0001"),
1927 )
1928 .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1929 let mut n = match read_node_opt(paths, &parent_node_id)? {
1930 Some(n) => n,
1931 None => return Ok(vec![]),
1932 };
1933 let new_ref = ChildRef {
1934 run_id: child_run_id,
1935 node_id: child_node_id,
1936 };
1937 if n.children.iter().any(|c| c == &new_ref) {
1938 return Ok(vec![]);
1941 }
1942 n.children.push(new_ref);
1943 n.updated_at = ev.ts;
1944 Ok(vec![ProjectionOp::Node(n)])
1945}
1946
1947fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1958 let events_path = paths.events();
1959 let node_id = require_envelope_node_id(&events_path, ev)?;
1960 let raw = ev
1961 .data
1962 .get("pid")
1963 .and_then(Value::as_i64)
1964 .ok_or_else(|| Error::CorruptEventLog {
1965 path: events_path.clone(),
1966 reason: format!(
1967 "event seq={} kind=supervisor.attached missing/invalid `pid`",
1968 ev.seq
1969 ),
1970 })?;
1971 let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1972 path: events_path.clone(),
1973 reason: format!(
1974 "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1975 ev.seq
1976 ),
1977 })?;
1978 let mut n = match read_node_opt(paths, &node_id)? {
1979 Some(n) => n,
1980 None => return Ok(vec![]),
1981 };
1982 if n.supervisor_pid == Some(pid) {
1983 return Ok(vec![]);
1984 }
1985 n.supervisor_pid = Some(pid);
1986 n.updated_at = ev.ts;
1987 Ok(vec![ProjectionOp::Node(n)])
1988}
1989
1990fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
2001 let events_path = paths.events();
2002 let node_id = require_envelope_node_id(&events_path, ev)?;
2003 let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
2004 RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
2008 let report_seq = ev
2009 .data
2010 .get("report_seq")
2011 .and_then(Value::as_u64)
2012 .ok_or_else(|| Error::CorruptEventLog {
2013 path: events_path.clone(),
2014 reason: format!(
2015 "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
2016 ev.seq
2017 ),
2018 })?;
2019 let mut n = match read_node_opt(paths, &node_id)? {
2020 Some(n) => n,
2021 None => return Ok(vec![]),
2022 };
2023 if let Some(prev) = n
2024 .last_processed_report_seq_by_child
2025 .get(child_run_id)
2026 .and_then(Value::as_u64)
2027 {
2028 if report_seq <= prev {
2029 return Ok(vec![]);
2030 }
2031 }
2032 n.last_processed_report_seq_by_child
2033 .insert(child_run_id.to_string(), Value::from(report_seq));
2034 n.updated_at = ev.ts;
2035 Ok(vec![ProjectionOp::Node(n)])
2036}
2037
2038#[cfg(test)]
2039mod tests {
2040 use super::*;
2041 use crate::schema::Event;
2042 use chrono::Utc;
2043 use tempfile::TempDir;
2044
2045 fn event(run_id: &str) -> Event {
2046 Event {
2047 ts: Utc::now(),
2048 seq: 1,
2049 kind: "run.status".into(),
2050 run_id: RunId::parse_str(run_id).unwrap(),
2051 node_id: None,
2052 idempotency_key: None,
2053 data: serde_json::json!({ "status": "running" }),
2054 }
2055 }
2056
2057 #[test]
2058 fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
2059 let tmp = TempDir::new().unwrap();
2063 let run_id = "01jxsnap000000000000000000";
2064 let rid = RunId::parse_str(run_id).unwrap();
2065 let dir = crate::run_dir(tmp.path(), &rid);
2066 std::fs::create_dir_all(&dir).unwrap();
2067 let paths = RunPaths::new(dir, run_id).unwrap();
2068
2069 let mut created = event(run_id);
2072 created.kind = "run.created".into();
2073 created.data = serde_json::json!({
2074 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2075 });
2076 apply_event(&paths, &created).expect("run.created applies");
2077 let manifest_before = std::fs::read(paths.manifest()).unwrap();
2078
2079 for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
2080 let mut ev = event(run_id);
2081 ev.seq = seq;
2082 ev.kind = kind.into();
2083 ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
2085 let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
2086 assert!(ops.is_empty(), "{kind} must plan no projection ops");
2087 apply_event(&paths, &ev).expect("audit kind applies as no-op");
2089 }
2090
2091 assert_eq!(
2093 std::fs::read(paths.manifest()).unwrap(),
2094 manifest_before,
2095 "audit events must not mutate the manifest"
2096 );
2097 assert!(!paths.nodes_dir().exists(), "no node projection created");
2098 }
2099
2100 #[test]
2101 fn run_created_folds_harness_when_present_and_defaults_none() {
2102 let tmp = TempDir::new().unwrap();
2103
2104 let run_id = "01jxhrnsaa0000000000000001";
2106 let rid = RunId::parse_str(run_id).unwrap();
2107 let dir = crate::run_dir(tmp.path(), &rid);
2108 std::fs::create_dir_all(&dir).unwrap();
2109 let paths = RunPaths::new(dir, run_id).unwrap();
2110 let mut created = event(run_id);
2111 created.kind = "run.created".into();
2112 created.data = serde_json::json!({
2113 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
2114 "harness": "pi", "harness_source": "flag",
2115 });
2116 apply_event(&paths, &created).expect("run.created applies");
2117 let m = read_manifest_opt(&paths).unwrap().unwrap();
2118 assert_eq!(m.harness.as_deref(), Some("pi"));
2119
2120 let run_id2 = "01jxhrnsaa0000000000000002";
2122 let rid2 = RunId::parse_str(run_id2).unwrap();
2123 let dir2 = crate::run_dir(tmp.path(), &rid2);
2124 std::fs::create_dir_all(&dir2).unwrap();
2125 let paths2 = RunPaths::new(dir2, run_id2).unwrap();
2126 let mut created2 = event(run_id2);
2127 created2.kind = "run.created".into();
2128 created2.data =
2129 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2130 apply_event(&paths2, &created2).expect("run.created applies");
2131 let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
2132 assert_eq!(m2.harness, None);
2133 }
2134
2135 fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
2138 let rid = RunId::parse_str(run_id).unwrap();
2139 let dir = crate::run_dir(tmp.path(), &rid);
2140 std::fs::create_dir_all(&dir).unwrap();
2141 let paths = RunPaths::new(dir, run_id).unwrap();
2142 let mut created = event(run_id);
2143 created.kind = "run.created".into();
2144 created.data =
2145 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2146 apply_event(&paths, &created).expect("run.created applies");
2147 let mut node = event(run_id);
2148 node.seq = 2;
2149 node.kind = "node.created".into();
2150 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2151 node.data = serde_json::json!({
2152 "kind": "spinoff",
2153 "branch": "wt/foo",
2154 "worktree_path": "/tmp/old-wt",
2155 "agent_pid": 111,
2156 });
2157 apply_event(&paths, &node).expect("node.created applies");
2158 paths
2159 }
2160
2161 #[test]
2165 fn node_retry_rewires_node_and_increments_attempts() {
2166 let tmp = TempDir::new().unwrap();
2167 let run_id = "01jxsnap000000000000000000";
2168 let paths = bootstrap_retry_node(&tmp, run_id);
2169
2170 let mut retry = event(run_id);
2171 retry.seq = 3;
2172 retry.kind = "node.retry".into();
2173 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2174 retry.data = serde_json::json!({
2175 "attempt": 1,
2176 "reason": "agent-died",
2177 "branch": "wt/foo-r1",
2178 "base_sha": "a".repeat(40),
2179 "worktree_path": "/tmp/new-wt",
2180 "agent_pid": 222,
2181 "tmux_session": "s",
2182 "tmux_window_id": "@9",
2183 });
2184 apply_event(&paths, &retry).expect("node.retry applies");
2185
2186 let n = read_n0001(&paths);
2187 assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
2188 assert_eq!(
2189 n.branch.as_deref(),
2190 Some("wt/foo-r1"),
2191 "rewired to new branch"
2192 );
2193 assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
2194 assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
2195 assert_eq!(n.status, Status::Pending, "node returns to pending");
2196 assert!(n.last_report.is_none());
2197 assert_eq!(
2198 n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
2199 Some("@9"),
2200 "rewired tmux identity"
2201 );
2202
2203 let mut retry2 = event(run_id);
2205 retry2.seq = 4;
2206 retry2.kind = "node.retry".into();
2207 retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2208 retry2.data = serde_json::json!({
2209 "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
2210 "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
2211 });
2212 apply_event(&paths, &retry2).expect("node.retry applies");
2213 assert_eq!(read_n0001(&paths).retry_attempts, 2);
2214 }
2215
2216 #[test]
2217 fn stale_worker_evidence_and_exit_do_not_cross_retry_generation() {
2218 let tmp = TempDir::new().unwrap();
2219 let run_id = "01jxsnap000000000000000000";
2220 let paths = bootstrap_retry_node(&tmp, run_id);
2221 let nid = NodeId::parse_str("n-0001").unwrap();
2222 let current_session = "018f5f64-b137-7d44-b2b4-4f02c3f646e8";
2223
2224 let mut retry = event(run_id);
2225 retry.seq = 3;
2226 retry.kind = "node.retry".into();
2227 retry.node_id = Some(nid.clone());
2228 retry.data = serde_json::json!({
2229 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
2230 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
2231 "pi_session_id": current_session,
2232 "pi_session_path": format!(".creating/pi-sessions/{run_id}/pi-session-{current_session}.jsonl"),
2233 "pi_session_cwd": "/tmp/new-wt"
2234 });
2235 apply_event(&paths, &retry).unwrap();
2236
2237 let mut stale_failure = event(run_id);
2238 stale_failure.seq = 4;
2239 stale_failure.kind = "worker.evidence.failed".into();
2240 stale_failure.node_id = Some(nid.clone());
2241 stale_failure.data = serde_json::json!({
2242 "attempt": 0,
2243 "session_id": "118f5f64-b137-7d44-b2b4-4f02c3f646e8",
2244 "error": "old attempt"
2245 });
2246 apply_event(&paths, &stale_failure).unwrap();
2247 assert_eq!(
2248 read_n0001(&paths).evidence.unwrap().status,
2249 EvidenceStatus::Pending
2250 );
2251
2252 let mut stale_exit = event(run_id);
2253 stale_exit.seq = 5;
2254 stale_exit.kind = "worker.exited".into();
2255 stale_exit.node_id = Some(nid.clone());
2256 stale_exit.data = serde_json::json!({"attempt":0,"exit_code":9});
2257 apply_event(&paths, &stale_exit).unwrap();
2258 assert!(read_n0001(&paths).worker_exit.is_none());
2259
2260 let mut archived = event(run_id);
2261 archived.seq = 6;
2262 archived.kind = "worker.evidence.archived".into();
2263 archived.node_id = Some(nid);
2264 archived.data = serde_json::json!({
2265 "attempt":1, "session_id":current_session,
2266 "transcript_path":"evidence/n-0001/pi-session.original.jsonl",
2267 "resume_path":"evidence/n-0001/pi-session.resume.jsonl",
2268 "pane_path":"evidence/n-0001/final-pane.log",
2269 "report_path":"evidence/n-0001/terminal-report.json",
2270 "transcript_sha256":"0".repeat(64)
2271 });
2272 apply_event(&paths, &archived).unwrap();
2273 assert_eq!(
2274 read_n0001(&paths).evidence.unwrap().status,
2275 EvidenceStatus::Complete
2276 );
2277 }
2278
2279 #[test]
2283 fn node_retry_against_terminal_node_is_noop() {
2284 let tmp = TempDir::new().unwrap();
2285 let run_id = "01jxsnap000000000000000000";
2286 let paths = bootstrap_retry_node(&tmp, run_id);
2287
2288 let mut report = event(run_id);
2290 report.seq = 3;
2291 report.kind = "node.report".into();
2292 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2293 report.data = serde_json::json!({ "success": true });
2294 apply_event(&paths, &report).expect("node.report applies");
2295 assert_eq!(read_n0001(&paths).status, Status::Done);
2296
2297 let mut retry = event(run_id);
2298 retry.seq = 4;
2299 retry.kind = "node.retry".into();
2300 retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2301 retry.data = serde_json::json!({
2302 "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
2303 "worktree_path": "/tmp/new-wt", "agent_pid": 222,
2304 });
2305 apply_event(&paths, &retry).expect("node.retry applies as no-op");
2306
2307 let n = read_n0001(&paths);
2308 assert_eq!(n.status, Status::Done, "terminal node not resurrected");
2309 assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
2310 assert_eq!(n.agent_pid, Some(111), "not rewired");
2311 }
2312
2313 #[test]
2314 fn apply_event_rejects_event_from_a_different_run() {
2315 let tmp = TempDir::new().unwrap();
2316 let run_id = "01jxsnap000000000000000000";
2317 let rid = RunId::parse_str(run_id).unwrap();
2318 let dir = crate::run_dir(tmp.path(), &rid);
2319 std::fs::create_dir_all(&dir).unwrap();
2320 let paths = RunPaths::new(dir, run_id).unwrap();
2321
2322 let foreign = event("02jxsnap000000000000000000");
2324 let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
2325 assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
2326
2327 let mine = event(run_id);
2330 apply_event(&paths, &mine).expect("matching run_id must be accepted");
2331 }
2332
2333 #[test]
2334 fn tmux_identity_from_data_reads_qualified_fields() {
2335 let d = serde_json::json!({
2336 "tmux_socket": "/private/tmp/tmux-501/default",
2337 "tmux_session": "taskfleet",
2338 "tmux_window_id": "@42",
2339 });
2340 let id = tmux_identity_from_data(&d).expect("qualified identity");
2341 assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
2342 assert_eq!(id.session, "taskfleet");
2343 assert_eq!(id.window_id, "@42");
2344 assert_eq!(id.pane_id, None);
2346
2347 let d2 = serde_json::json!({
2349 "tmux_socket": null,
2350 "tmux_session": "taskfleet",
2351 "tmux_window_id": "@7",
2352 });
2353 let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
2354 assert_eq!(id2.socket, None);
2355 assert_eq!(id2.window_id, "@7");
2356
2357 let d3 = serde_json::json!({
2359 "tmux_session": "taskfleet",
2360 "tmux_window_id": "@42",
2361 "tmux_pane_id": "%7",
2362 });
2363 let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
2364 assert_eq!(id3.pane_id.as_deref(), Some("%7"));
2365 assert_eq!(id3.capture_target(), "%7");
2366
2367 let d4 = serde_json::json!({
2370 "tmux_session": "taskfleet",
2371 "tmux_window_id": "@42",
2372 "tmux_pane_id": null,
2373 });
2374 let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
2375 assert_eq!(id4.pane_id, None);
2376 assert_eq!(id4.capture_target(), "@42");
2377 }
2378
2379 #[test]
2380 fn tmux_identity_from_data_back_compat_is_none() {
2381 let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
2383 assert!(tmux_identity_from_data(&legacy).is_none());
2384 let partial = serde_json::json!({ "tmux_window_id": "@42" });
2386 assert!(tmux_identity_from_data(&partial).is_none());
2387 }
2388
2389 #[test]
2392 fn node_created_populates_tmux_identity() {
2393 let tmp = TempDir::new().unwrap();
2394 let run_id = "01jxsnap000000000000000000";
2395 let rid = RunId::parse_str(run_id).unwrap();
2396 let dir = crate::run_dir(tmp.path(), &rid);
2397 std::fs::create_dir_all(&dir).unwrap();
2398 let paths = RunPaths::new(dir, run_id).unwrap();
2399
2400 let mut ev = event(run_id);
2401 ev.seq = 2;
2402 ev.kind = "node.created".into();
2403 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2404 ev.data = serde_json::json!({
2405 "kind": "spinoff",
2406 "tmux_window": "🚀 wt/x",
2407 "tmux_socket": "/private/tmp/tmux-501/default",
2408 "tmux_session": "taskfleet",
2409 "tmux_window_id": "@42",
2410 });
2411 apply_event(&paths, &ev).expect("node.created applies");
2412 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2413 .unwrap()
2414 .unwrap();
2415 let id = n.tmux_identity.expect("qualified identity recorded");
2416 assert_eq!(id.session, "taskfleet");
2417 assert_eq!(id.window_id, "@42");
2418 assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
2419
2420 let run2 = "02jxsnap000000000000000000";
2422 let rid2 = RunId::parse_str(run2).unwrap();
2423 let dir2 = crate::run_dir(tmp.path(), &rid2);
2424 std::fs::create_dir_all(&dir2).unwrap();
2425 let paths2 = RunPaths::new(dir2, run2).unwrap();
2426 let mut ev2 = event(run2);
2427 ev2.seq = 2;
2428 ev2.kind = "node.created".into();
2429 ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2430 ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
2431 apply_event(&paths2, &ev2).expect("legacy node.created applies");
2432 let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
2433 .unwrap()
2434 .unwrap();
2435 assert!(n2.tmux_identity.is_none());
2436 assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
2437 }
2438
2439 #[test]
2440 fn node_materialization_populates_missing_manifest_source_branch() {
2441 let tmp = TempDir::new().unwrap();
2442 let run_id = "01jxsnap000000000000000001";
2443 let rid = RunId::parse_str(run_id).unwrap();
2444 let dir = crate::run_dir(tmp.path(), &rid);
2445 std::fs::create_dir_all(&dir).unwrap();
2446 let paths = RunPaths::new(dir, run_id).unwrap();
2447
2448 let mut created = event(run_id);
2449 created.kind = "run.created".into();
2450 created.data = serde_json::json!({
2451 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2452 });
2453 apply_event(&paths, &created).unwrap();
2454 assert!(read_manifest_opt(&paths)
2455 .unwrap()
2456 .unwrap()
2457 .source_branch
2458 .is_none());
2459
2460 let mut node = event(run_id);
2461 node.seq = 2;
2462 node.kind = "node.created".into();
2463 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2464 node.data = serde_json::json!({
2465 "kind": "spinoff",
2466 "source_branch": "main",
2467 "worktree_path": "/tmp/wt/pending"
2468 });
2469 apply_event(&paths, &node).unwrap();
2470
2471 let manifest = read_manifest_opt(&paths).unwrap().unwrap();
2472 assert_eq!(manifest.status, Status::Pending);
2473 assert_eq!(manifest.source_branch.as_deref(), Some("main"));
2474 let node = read_n0001(&paths);
2475 assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
2476 }
2477
2478 #[test]
2479 fn node_materialization_preserves_explicit_manifest_source_branch() {
2480 let tmp = TempDir::new().unwrap();
2481 let run_id = "01jxsnap000000000000000002";
2482 let rid = RunId::parse_str(run_id).unwrap();
2483 let dir = crate::run_dir(tmp.path(), &rid);
2484 std::fs::create_dir_all(&dir).unwrap();
2485 let paths = RunPaths::new(dir, run_id).unwrap();
2486
2487 let mut created = event(run_id);
2488 created.kind = "run.created".into();
2489 created.data = serde_json::json!({
2490 "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
2491 "source_branch": "release"
2492 });
2493 apply_event(&paths, &created).unwrap();
2494
2495 let mut node = event(run_id);
2496 node.seq = 2;
2497 node.kind = "node.created".into();
2498 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2499 node.data = serde_json::json!({
2500 "kind": "spinoff", "source_branch": "main"
2501 });
2502 apply_event(&paths, &node).unwrap();
2503
2504 assert_eq!(
2505 read_manifest_opt(&paths)
2506 .unwrap()
2507 .unwrap()
2508 .source_branch
2509 .as_deref(),
2510 Some("release")
2511 );
2512 }
2513
2514 fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
2517 let rid = RunId::parse_str(run_id).unwrap();
2518 let dir = crate::run_dir(tmp.path(), &rid);
2519 std::fs::create_dir_all(&dir).unwrap();
2520 let paths = RunPaths::new(dir, run_id).unwrap();
2521
2522 let mut created = event(run_id);
2523 created.kind = "run.created".into();
2524 created.data = serde_json::json!({
2525 "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2526 });
2527 apply_event(&paths, &created).expect("run.created applies");
2528
2529 let mut node = event(run_id);
2530 node.seq = 2;
2531 node.kind = "node.created".into();
2532 node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2533 node.data = serde_json::json!({ "kind": "spinoff" });
2534 apply_event(&paths, &node).expect("node.created applies");
2535 paths
2536 }
2537
2538 fn read_n0001(paths: &RunPaths) -> Node {
2539 read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
2540 .unwrap()
2541 .unwrap()
2542 }
2543
2544 #[test]
2545 fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
2546 let tmp = TempDir::new().unwrap();
2547 let run_id = "01jxwd0000000000000000000w";
2548 let paths = seed_run_with_node(&tmp, run_id);
2549 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2550 let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
2551 let mut open = event(run_id);
2552 open.seq = 3;
2553 open.ts = opened_at;
2554 open.kind = "node.awaiting_input".into();
2555 open.node_id = nid.clone();
2556 open.data = serde_json::json!({ "discussion_items": [{
2557 "topic": "Which scope?",
2558 "options": ["small", "large"],
2559 "recommended_default": "small"
2560 }] });
2561 apply_event(&paths, &open).unwrap();
2562 let first = read_n0001(&paths).awaiting_input.unwrap();
2563 assert_eq!(first.opened_at, opened_at);
2564 assert_eq!(first.event_seq, 3);
2565
2566 let mut duplicate = open.clone();
2568 duplicate.seq = 4;
2569 duplicate.ts = opened_at + chrono::Duration::hours(1);
2570 apply_event(&paths, &duplicate).unwrap();
2571 let still_first = read_n0001(&paths).awaiting_input.unwrap();
2572 assert_eq!(still_first.opened_at, opened_at);
2573 assert_eq!(still_first.event_seq, 3);
2574
2575 let mut stale = event(run_id);
2577 stale.seq = 5;
2578 stale.kind = "node.input_resolved".into();
2579 stale.node_id = nid.clone();
2580 stale.data = serde_json::json!({ "event_seq": 2 });
2581 apply_event(&paths, &stale).unwrap();
2582 assert!(read_n0001(&paths).awaiting_input.is_some());
2583
2584 let mut resolved = stale;
2585 resolved.seq = 6;
2586 resolved.data = serde_json::json!({ "event_seq": 3 });
2587 apply_event(&paths, &resolved).unwrap();
2588 assert!(read_n0001(&paths).awaiting_input.is_none());
2589 }
2590
2591 #[test]
2592 fn awaiting_input_rejects_missing_default_without_mutating_projection() {
2593 let tmp = TempDir::new().unwrap();
2594 let run_id = "01jxwd0000000000000000000x";
2595 let paths = seed_run_with_node(&tmp, run_id);
2596 let mut open = event(run_id);
2597 open.seq = 3;
2598 open.kind = "node.awaiting_input".into();
2599 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2600 open.data = serde_json::json!({ "discussion_items": [{
2601 "topic": "Which scope?", "options": ["small", "large"]
2602 }] });
2603 assert!(reduce_event_to_ops(&paths, &open).is_err());
2604 assert!(read_n0001(&paths).awaiting_input.is_none());
2605 }
2606
2607 #[test]
2608 fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
2609 let tmp = TempDir::new().unwrap();
2610 let run_id = "01jxwd0000000000000000000y";
2611 let paths = seed_run_with_node(&tmp, run_id);
2612 let nid = Some(NodeId::parse_str("n-0001").unwrap());
2613
2614 let mut open = event(run_id);
2615 open.seq = 3;
2616 open.kind = "node.awaiting_input".into();
2617 open.node_id = nid.clone();
2618 open.data = serde_json::json!({ "discussion_items": [{
2619 "topic": "Which scope?", "options": ["small", "large"],
2620 "recommended_default": "small"
2621 }] });
2622 apply_event(&paths, &open).unwrap();
2623
2624 let mut malformed_duplicate = open.clone();
2625 malformed_duplicate.seq = 4;
2626 malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
2627 assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
2628
2629 let mut exited = event(run_id);
2630 exited.seq = 5;
2631 exited.kind = "worker.exited".into();
2632 exited.node_id = nid;
2633 exited.data = serde_json::json!({ "exit_code": 0 });
2634 apply_event(&paths, &exited).unwrap();
2635 let node = read_n0001(&paths);
2636 assert!(node.awaiting_input.is_none());
2637 assert!(node.worker_exit.is_some());
2638
2639 let mut delayed_open = open;
2640 delayed_open.seq = 6;
2641 assert!(reduce_event_to_ops(&paths, &delayed_open)
2642 .unwrap()
2643 .is_empty());
2644 }
2645
2646 #[test]
2647 fn input_resolved_requires_generation_even_when_nothing_is_open() {
2648 let tmp = TempDir::new().unwrap();
2649 let run_id = "01jxwd0000000000000000000z";
2650 let paths = seed_run_with_node(&tmp, run_id);
2651 let mut resolved = event(run_id);
2652 resolved.seq = 3;
2653 resolved.kind = "node.input_resolved".into();
2654 resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2655 resolved.data = serde_json::json!({});
2656 assert!(reduce_event_to_ops(&paths, &resolved).is_err());
2657 }
2658
2659 #[test]
2660 fn awaiting_input_default_must_be_one_of_options() {
2661 let tmp = TempDir::new().unwrap();
2662 let run_id = "01jxwd00000000000000000010";
2663 let paths = seed_run_with_node(&tmp, run_id);
2664 let mut open = event(run_id);
2665 open.seq = 3;
2666 open.kind = "node.awaiting_input".into();
2667 open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2668 open.data = serde_json::json!({ "discussion_items": [{
2669 "topic": "Which scope?", "options": ["small", "large"],
2670 "recommended_default": "other"
2671 }] });
2672 assert!(reduce_event_to_ops(&paths, &open).is_err());
2673 }
2674
2675 fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
2676 let mut ev = event(run_id);
2677 ev.seq = seq;
2678 ev.kind = KIND_MERGE_STARTED.into();
2679 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2680 ev.data = serde_json::json!({
2681 "op_id": op_id,
2682 "source_branch": "main",
2683 "worker_branch": "wt/worker",
2684 "expected_source_oid": expected,
2685 "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2686 "base_sha": null,
2687 "driver_pid": 4242,
2688 "driver_pid_start_secs": null,
2689 "started_at": "2026-08-15T00:00:00Z",
2690 });
2691 ev
2692 }
2693
2694 #[test]
2697 fn merge_started_records_pending_transaction() {
2698 let tmp = TempDir::new().unwrap();
2699 let run_id = "01jxsnap000000000000000000";
2700 let paths = seed_run_with_node(&tmp, run_id);
2701
2702 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2703 let n = read_n0001(&paths);
2704 assert_eq!(
2705 n.status,
2706 Status::Pending,
2707 "recording a merge is not terminal"
2708 );
2709 let txn = n.pending_merge.expect("transaction recorded");
2710 assert_eq!(txn.op_id, "op-1");
2711 assert_eq!(txn.expected_source_oid, "aaa");
2712 }
2713
2714 #[test]
2717 fn merge_aborted_clears_matching_transaction_only() {
2718 let tmp = TempDir::new().unwrap();
2719 let run_id = "01jxsnap000000000000000000";
2720 let paths = seed_run_with_node(&tmp, run_id);
2721 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2722
2723 let mut stale = event(run_id);
2725 stale.seq = 4;
2726 stale.kind = KIND_MERGE_ABORTED.into();
2727 stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2728 stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2729 apply_event(&paths, &stale).unwrap();
2730 assert!(
2731 read_n0001(&paths).pending_merge.is_some(),
2732 "stale abort is a no-op"
2733 );
2734
2735 let mut abort = event(run_id);
2737 abort.seq = 5;
2738 abort.kind = KIND_MERGE_ABORTED.into();
2739 abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2740 abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2741 apply_event(&paths, &abort).unwrap();
2742 let n = read_n0001(&paths);
2743 assert!(
2744 n.pending_merge.is_none(),
2745 "matching abort clears the transaction"
2746 );
2747 assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2748 }
2749
2750 #[test]
2753 fn terminal_report_clears_pending_merge() {
2754 let tmp = TempDir::new().unwrap();
2755 let run_id = "01jxsnap000000000000000000";
2756 let paths = seed_run_with_node(&tmp, run_id);
2757 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2758
2759 let mut report = event(run_id);
2760 report.seq = 4;
2761 report.kind = "node.report".into();
2762 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2763 report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2764 apply_event(&paths, &report).unwrap();
2765 let n = read_n0001(&paths);
2766 assert_eq!(n.status, Status::Done);
2767 assert!(
2768 n.pending_merge.is_none(),
2769 "completed merge clears the transaction"
2770 );
2771 }
2772
2773 #[test]
2782 fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2783 let tmp = TempDir::new().unwrap();
2784
2785 let drive = |run_id: &str, report_data: Value| -> Status {
2788 let paths = seed_run_with_node(&tmp, run_id);
2789 let mut fail = event(run_id);
2791 fail.seq = 3;
2792 fail.kind = "node.status".into();
2793 fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2794 fail.data = serde_json::json!({ "status": "failed" });
2795 apply_event(&paths, &fail).unwrap();
2796 assert_eq!(read_n0001(&paths).status, Status::Failed);
2797 let mut report = event(run_id);
2799 report.seq = 4;
2800 report.kind = "node.report".into();
2801 report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2802 report.data = report_data;
2803 apply_event(&paths, &report).unwrap();
2804 read_n0001(&paths).status
2805 };
2806
2807 let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2809 crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2810 assert_eq!(
2811 drive("01jxsnap000000000000000001", agent_forged),
2812 Status::Failed,
2813 "an Agent-origin report with a forged via must not be adopted"
2814 );
2815
2816 let malformed = serde_json::json!({
2818 "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2819 });
2820 assert_eq!(
2821 drive("01jxsnap000000000000000002", malformed),
2822 Status::Failed,
2823 "a malformed origin must not re-unlock the legacy via adoption path"
2824 );
2825
2826 let mut run_merge = serde_json::json!({ "success": true });
2828 crate::ReportOrigin::RunMerge {
2829 op_id: Some("op-1".into()),
2830 worker_oid: Some("cafebabe".into()),
2831 }
2832 .stamp(&mut run_merge);
2833 assert_eq!(
2834 drive("01jxsnap000000000000000003", run_merge),
2835 Status::Done,
2836 "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2837 );
2838
2839 let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2842 assert_eq!(
2843 drive("01jxsnap000000000000000004", legacy),
2844 Status::Done,
2845 "a legacy via-only report (no origin field) is still adopted"
2846 );
2847 }
2848
2849 #[test]
2853 fn terminal_node_status_clears_pending_merge() {
2854 let tmp = TempDir::new().unwrap();
2855 let run_id = "01jxsnap000000000000000000";
2856 let paths = seed_run_with_node(&tmp, run_id);
2857 apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2858 assert!(read_n0001(&paths).pending_merge.is_some());
2859
2860 let mut status = event(run_id);
2861 status.seq = 4;
2862 status.kind = "node.status".into();
2863 status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2864 status.data = serde_json::json!({ "status": "failed" });
2865 apply_event(&paths, &status).unwrap();
2866 let n = read_n0001(&paths);
2867 assert_eq!(n.status, Status::Failed);
2868 assert!(
2869 n.pending_merge.is_none(),
2870 "terminal status clears the transaction"
2871 );
2872 }
2873
2874 #[test]
2878 fn supervisor_attached_sets_supervisor_pid() {
2879 let tmp = TempDir::new().unwrap();
2880 let run_id = "01jxsnap000000000000000000";
2881 let paths = seed_run_with_node(&tmp, run_id);
2882 assert_eq!(read_n0001(&paths).supervisor_pid, None);
2883
2884 let mut ev = event(run_id);
2885 ev.seq = 3;
2886 ev.kind = "supervisor.attached".into();
2887 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2888 ev.data = serde_json::json!({ "pid": 47820 });
2889 apply_event(&paths, &ev).expect("supervisor.attached applies");
2890 assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2891 }
2892
2893 #[test]
2896 fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2897 let tmp = TempDir::new().unwrap();
2898 let run_id = "01jxsnap000000000000000000";
2899 let paths = seed_run_with_node(&tmp, run_id);
2900
2901 let mut ev = event(run_id);
2902 ev.kind = "supervisor.attached".into();
2903 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2904
2905 ev.seq = 3;
2906 ev.data = serde_json::json!({ "pid": 100 });
2907 apply_event(&paths, &ev).expect("first attach applies");
2908 assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2909
2910 ev.seq = 4;
2912 ev.data = serde_json::json!({ "pid": 200 });
2913 apply_event(&paths, &ev).expect("second attach applies");
2914 let after_second = read_n0001(&paths);
2915 assert_eq!(after_second.supervisor_pid, Some(200));
2916
2917 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2920 assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2921 apply_event(&paths, &ev).expect("replay applies as no-op");
2922 assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2923 }
2924
2925 #[test]
2928 fn supervisor_cursor_advanced_sets_report_cursor() {
2929 let tmp = TempDir::new().unwrap();
2930 let run_id = "01jxsnap000000000000000000";
2931 let paths = seed_run_with_node(&tmp, run_id);
2932 let child = "02jxsnap000000000000000000";
2933
2934 let mut ev = event(run_id);
2935 ev.seq = 3;
2936 ev.kind = "supervisor.cursor_advanced".into();
2937 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2938 ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2939 apply_event(&paths, &ev).expect("cursor_advanced applies");
2940
2941 let n = read_n0001(&paths);
2942 assert_eq!(
2943 n.last_processed_report_seq_by_child.get(child),
2944 Some(&Value::from(7u64))
2945 );
2946 }
2947
2948 #[test]
2953 fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2954 let tmp = TempDir::new().unwrap();
2955 let run_id = "01jxsnap000000000000000000";
2956 let paths = seed_run_with_node(&tmp, run_id);
2957 let child_a = "02jxsnap000000000000000000";
2958 let child_b = "03jxsnap000000000000000000";
2959
2960 let mut ev = event(run_id);
2961 ev.kind = "supervisor.cursor_advanced".into();
2962 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2963
2964 ev.seq = 3;
2965 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2966 apply_event(&paths, &ev).expect("seq 5 applies");
2967
2968 let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2970 assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2971
2972 ev.seq = 4;
2974 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2975 let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2976 assert!(ops.is_empty(), "older seq must plan no ops");
2977 apply_event(&paths, &ev).expect("older seq applies as no-op");
2978 assert_eq!(
2979 read_n0001(&paths)
2980 .last_processed_report_seq_by_child
2981 .get(child_a),
2982 Some(&Value::from(5u64))
2983 );
2984
2985 ev.seq = 5;
2987 ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2988 apply_event(&paths, &ev).expect("higher seq applies");
2989 ev.seq = 6;
2990 ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2991 apply_event(&paths, &ev).expect("second child applies");
2992
2993 let n = read_n0001(&paths);
2994 assert_eq!(
2995 n.last_processed_report_seq_by_child.get(child_a),
2996 Some(&Value::from(9u64))
2997 );
2998 assert_eq!(
2999 n.last_processed_report_seq_by_child.get(child_b),
3000 Some(&Value::from(1u64))
3001 );
3002 }
3003
3004 #[test]
3007 fn supervisor_state_events_reject_malformed_payloads() {
3008 let tmp = TempDir::new().unwrap();
3009 let run_id = "01jxsnap000000000000000000";
3010 let paths = seed_run_with_node(&tmp, run_id);
3011 let nid = Some(NodeId::parse_str("n-0001").unwrap());
3012
3013 let mut ev = event(run_id);
3015 ev.seq = 3;
3016 ev.kind = "supervisor.attached".into();
3017 ev.node_id = nid.clone();
3018 ev.data = serde_json::json!({});
3019 assert!(matches!(
3020 reduce_event_to_ops(&paths, &ev),
3021 Err(Error::CorruptEventLog { .. })
3022 ));
3023
3024 ev.node_id = None;
3026 ev.data = serde_json::json!({ "pid": 1 });
3027 assert!(matches!(
3028 reduce_event_to_ops(&paths, &ev),
3029 Err(Error::CorruptEventLog { .. })
3030 ));
3031
3032 let mut ev2 = event(run_id);
3034 ev2.seq = 4;
3035 ev2.kind = "supervisor.cursor_advanced".into();
3036 ev2.node_id = nid.clone();
3037 ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
3038 assert!(matches!(
3039 reduce_event_to_ops(&paths, &ev2),
3040 Err(Error::CorruptEventLog { .. })
3041 ));
3042
3043 ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
3045 assert!(matches!(
3046 reduce_event_to_ops(&paths, &ev2),
3047 Err(Error::CorruptEventLog { .. })
3048 ));
3049 }
3050
3051 #[test]
3058 fn removed_or_garbage_kind_in_created_events_is_rejected() {
3059 let tmp = TempDir::new().unwrap();
3060 let run_id = "01jxsnap000000000000000000";
3061 let rid = RunId::parse_str(run_id).unwrap();
3062 let dir = crate::run_dir(tmp.path(), &rid);
3063 std::fs::create_dir_all(&dir).unwrap();
3064 let paths = RunPaths::new(dir, run_id).unwrap();
3065
3066 for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
3067 let mut ev = event(run_id);
3068 ev.kind = "run.created".into();
3069 ev.node_id = None;
3070 ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
3071 assert!(
3072 matches!(
3073 reduce_event_to_ops(&paths, &ev),
3074 Err(Error::CorruptEventLog { .. })
3075 ),
3076 "run.created with kind {bad:?} must be rejected, not folded to Unknown"
3077 );
3078 }
3079
3080 let mut ok = event(run_id);
3083 ok.kind = "run.created".into();
3084 ok.node_id = None;
3085 ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
3086 assert!(reduce_event_to_ops(&paths, &ok).is_ok());
3087 }
3088
3089 #[cfg(unix)]
3098 fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
3099 use std::os::unix::fs::MetadataExt;
3100 let mut consider = vec![paths.manifest()];
3101 for dir in [paths.nodes_dir()] {
3102 if let Ok(rd) = std::fs::read_dir(&dir) {
3103 for ent in rd.flatten() {
3104 let p = ent.path();
3105 if p.extension().and_then(|s| s.to_str()) == Some("json") {
3106 consider.push(p);
3107 }
3108 }
3109 }
3110 }
3111 let mut map = std::collections::BTreeMap::new();
3112 for p in consider {
3113 if let Ok(md) = std::fs::symlink_metadata(&p) {
3114 if md.file_type().is_file() {
3115 map.insert(p, md.ino());
3116 }
3117 }
3118 }
3119 map
3120 }
3121
3122 #[cfg(unix)]
3131 fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
3132 use std::collections::BTreeSet;
3133 let before = projection_inodes(paths);
3134 let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
3135 .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
3136 .into_iter()
3137 .collect();
3138 apply_event(paths, ev)
3139 .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
3140 let after = projection_inodes(paths);
3141 let touched: BTreeSet<PathBuf> = after
3142 .iter()
3143 .filter(|(p, ino)| before.get(*p) != Some(*ino))
3144 .map(|(p, _)| p.clone())
3145 .collect();
3146 assert_eq!(
3147 planned, touched,
3148 "kind={}: plan_projections must name exactly the files apply_event writes",
3149 ev.kind
3150 );
3151 if expect_writes {
3152 assert!(
3153 !touched.is_empty(),
3154 "kind={}: expected this event to write at least one projection",
3155 ev.kind
3156 );
3157 }
3158 }
3159
3160 #[cfg(unix)]
3167 #[test]
3168 fn plan_projections_matches_apply_for_every_kind() {
3169 let tmp = TempDir::new().unwrap();
3170 let run_id = "01jxsnap000000000000000000";
3171 let rid = RunId::parse_str(run_id).unwrap();
3172 let dir = crate::run_dir(tmp.path(), &rid);
3173 std::fs::create_dir_all(&dir).unwrap();
3174 let paths = RunPaths::new(dir, run_id).unwrap();
3175 let nid = || Some(NodeId::parse_str("n-0001").unwrap());
3176 let child = "02jxsnap000000000000000000";
3177
3178 let mut next_seq = 0u64;
3180 let mut at = |kind: &str, node_id, data| {
3181 next_seq += 1;
3182 Event {
3183 ts: Utc::now(),
3184 seq: next_seq,
3185 kind: kind.into(),
3186 run_id: rid.clone(),
3187 node_id,
3188 idempotency_key: None,
3189 data,
3190 }
3191 };
3192
3193 assert_plan_matches_apply(
3195 &paths,
3196 &at(
3197 "run.created",
3198 None,
3199 serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
3200 ),
3201 true,
3202 );
3203 assert_plan_matches_apply(
3205 &paths,
3206 &at(
3207 "run.status",
3208 None,
3209 serde_json::json!({ "status": "running" }),
3210 ),
3211 true,
3212 );
3213 assert_plan_matches_apply(
3215 &paths,
3216 &at(
3217 "node.created",
3218 nid(),
3219 serde_json::json!({ "kind": "spinoff" }),
3220 ),
3221 true,
3222 );
3223 assert_plan_matches_apply(
3225 &paths,
3226 &at(
3227 "node.status",
3228 nid(),
3229 serde_json::json!({ "status": "running" }),
3230 ),
3231 true,
3232 );
3233 assert_plan_matches_apply(
3235 &paths,
3236 &at(
3237 "supervisor.attached",
3238 nid(),
3239 serde_json::json!({ "pid": 4242 }),
3240 ),
3241 true,
3242 );
3243 assert_plan_matches_apply(
3245 &paths,
3246 &at(
3247 "supervisor.cursor_advanced",
3248 nid(),
3249 serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
3250 ),
3251 true,
3252 );
3253 assert_plan_matches_apply(
3255 &paths,
3256 &at(
3257 "child.spawned",
3258 nid(),
3259 serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
3260 ),
3261 true,
3262 );
3263 assert_plan_matches_apply(
3265 &paths,
3266 &at("node.report", nid(), serde_json::json!({ "success": true })),
3267 true,
3268 );
3269 assert_plan_matches_apply(
3272 &paths,
3273 &at(
3274 "node.status",
3275 nid(),
3276 serde_json::json!({ "status": "failed" }),
3277 ),
3278 false,
3279 );
3280 for kind in [
3282 "supervisor.exited",
3283 "orchestrator.decision",
3284 "discuss.critical",
3285 "cleanup.window_missing",
3286 ] {
3287 assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
3288 }
3289 }
3290
3291 #[test]
3295 fn worker_exited_records_clean_exit_without_transitioning_status() {
3296 let tmp = TempDir::new().unwrap();
3297 let run_id = "01jxsnap000000000000000000";
3298 let paths = bootstrap_retry_node(&tmp, run_id);
3299
3300 let mut ev = event(run_id);
3301 ev.seq = 3;
3302 ev.kind = "worker.exited".into();
3303 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
3304 ev.data = serde_json::json!({ "exit_code": 0 });
3305 apply_event(&paths, &ev).expect("worker.exited applies");
3306
3307 let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
3308 .unwrap()
3309 .unwrap();
3310 let exit = n.worker_exit.expect("worker_exit recorded");
3311 assert_eq!(exit.code, Some(0));
3312 assert_eq!(exit.signal, None);
3313 assert!(exit.is_clean());
3314 assert_eq!(
3315 n.status,
3316 Status::Pending,
3317 "the exit fact never transitions status"
3318 );
3319 }
3320
3321 #[test]
3325 fn worker_exited_records_signal_and_is_first_write_wins() {
3326 let tmp = TempDir::new().unwrap();
3327 let run_id = "01jxsnap000000000000000000";
3328 let paths = bootstrap_retry_node(&tmp, run_id);
3329 let nid = NodeId::parse_str("n-0001").unwrap();
3330
3331 let mut ev = event(run_id);
3332 ev.seq = 3;
3333 ev.kind = "worker.exited".into();
3334 ev.node_id = Some(nid.clone());
3335 ev.data = serde_json::json!({ "signal": 9 });
3336 apply_event(&paths, &ev).expect("worker.exited applies");
3337
3338 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3339 let exit = n.worker_exit.expect("worker_exit recorded");
3340 assert_eq!(exit.signal, Some(9));
3341 assert!(exit.is_failure());
3342
3343 let mut dup = event(run_id);
3346 dup.seq = 4;
3347 dup.kind = "worker.exited".into();
3348 dup.node_id = Some(nid.clone());
3349 dup.data = serde_json::json!({ "exit_code": 0 });
3350 apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
3351 let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
3352 assert_eq!(
3353 n2.worker_exit.unwrap().signal,
3354 Some(9),
3355 "first-write-wins: the replayed exit must not overwrite the recorded fact"
3356 );
3357 }
3358
3359 #[test]
3363 fn worker_exited_without_code_or_signal_is_corrupt() {
3364 let tmp = TempDir::new().unwrap();
3365 let run_id = "01jxsnap000000000000000000";
3366 let paths = bootstrap_retry_node(&tmp, run_id);
3367
3368 let mut ev = event(run_id);
3369 ev.seq = 3;
3370 ev.kind = "worker.exited".into();
3371 ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
3372 ev.data = serde_json::json!({});
3373 match reduce_event_to_ops(&paths, &ev) {
3374 Err(Error::CorruptEventLog { .. }) => {}
3375 Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
3376 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3377 }
3378
3379 ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
3382 match reduce_event_to_ops(&paths, &ev) {
3383 Err(Error::CorruptEventLog { .. }) => {}
3384 Ok(_) => panic!("a worker.exited with both fields must be rejected"),
3385 Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3386 }
3387 }
3388
3389 #[test]
3395 fn node_death_observed_records_first_death_first_write_wins() {
3396 let tmp = TempDir::new().unwrap();
3397 let run_id = "01jxsnap000000000000000000";
3398 let paths = bootstrap_retry_node(&tmp, run_id);
3399 let nid = NodeId::parse_str("n-0001").unwrap();
3400
3401 let mut ev = event(run_id);
3402 ev.seq = 3;
3403 ev.kind = "node.death_observed".into();
3404 ev.node_id = Some(nid.clone());
3405 ev.data = serde_json::json!({});
3406 apply_event(&paths, &ev).expect("node.death_observed applies");
3407 let first = read_node_opt(&paths, &nid)
3408 .unwrap()
3409 .unwrap()
3410 .first_death_at
3411 .expect("first_death_at recorded");
3412 assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
3413
3414 let mut later = event(run_id);
3416 later.seq = 4;
3417 later.kind = "node.death_observed".into();
3418 later.node_id = Some(nid.clone());
3419 later.ts = ev.ts + chrono::Duration::seconds(30);
3420 later.data = serde_json::json!({});
3421 apply_event(&paths, &later).expect("re-observation applies as no-op");
3422 assert_eq!(
3423 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3424 Some(first),
3425 "first-write-wins: a re-observation must not reset the anchor"
3426 );
3427 }
3428
3429 #[test]
3434 fn node_death_observed_noop_when_worker_exit_present() {
3435 let tmp = TempDir::new().unwrap();
3436 let run_id = "01jxsnap000000000000000000";
3437 let paths = bootstrap_retry_node(&tmp, run_id);
3438 let nid = NodeId::parse_str("n-0001").unwrap();
3439
3440 let mut exit = event(run_id);
3442 exit.seq = 3;
3443 exit.kind = "worker.exited".into();
3444 exit.node_id = Some(nid.clone());
3445 exit.data = serde_json::json!({ "exit_code": 0 });
3446 apply_event(&paths, &exit).unwrap();
3447
3448 let mut death = event(run_id);
3450 death.seq = 4;
3451 death.kind = "node.death_observed".into();
3452 death.node_id = Some(nid.clone());
3453 death.data = serde_json::json!({});
3454 apply_event(&paths, &death).expect("applies as no-op");
3455 assert_eq!(
3456 read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3457 None,
3458 "a told worker.exited makes the crash backstop moot; no anchor recorded"
3459 );
3460 }
3461
3462 #[test]
3467 fn node_retry_clears_worker_exit() {
3468 let tmp = TempDir::new().unwrap();
3469 let run_id = "01jxsnap000000000000000000";
3470 let paths = bootstrap_retry_node(&tmp, run_id);
3471 let nid = NodeId::parse_str("n-0001").unwrap();
3472
3473 let mut exit = event(run_id);
3475 exit.seq = 3;
3476 exit.kind = "worker.exited".into();
3477 exit.node_id = Some(nid.clone());
3478 exit.data = serde_json::json!({ "exit_code": 7 });
3479 apply_event(&paths, &exit).unwrap();
3480 assert!(read_node_opt(&paths, &nid)
3481 .unwrap()
3482 .unwrap()
3483 .worker_exit
3484 .is_some());
3485
3486 let mut retry = event(run_id);
3488 retry.seq = 4;
3489 retry.kind = "node.retry".into();
3490 retry.node_id = Some(nid.clone());
3491 retry.data = serde_json::json!({
3492 "attempt": 1,
3493 "reason": "agent-died",
3494 "branch": "wt/foo",
3495 "worktree_path": "/tmp/new-wt",
3496 "agent_pid": 222,
3497 });
3498 apply_event(&paths, &retry).unwrap();
3499
3500 let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3501 assert!(
3502 n.worker_exit.is_none(),
3503 "node.retry must clear the previous attempt's worker_exit"
3504 );
3505 assert_eq!(
3506 n.status,
3507 Status::Pending,
3508 "retry returns the node to Pending"
3509 );
3510 }
3511}