Skip to main content

taskfleet_core/
reducer.rs

1//! Event → projection reducer (design.md §1.4).
2//!
3//! Each event mutates zero or more projection files. Unknown kinds are
4//! ignored for forward compatibility. The reducer expects to run under the
5//! per-run `flock`.
6//!
7//! Idempotency contract (per design.md §7.3 at-least-once delivery):
8//! `*.created` reducers short-circuit when their projection file already
9//! exists; status/resolution reducers are no-ops once the terminal state
10//! has been reached. Replaying the same event stream against existing
11//! projections is therefore a clean no-op-or-apply.
12//!
13//! This idempotence is load-bearing for the `applied_seq` watermark
14//! (append-then-apply atomicity; see [`crate::schema::Manifest::applied_seq`]
15//! and [`crate::events::append_and_apply_event`]). Of the two options the spec
16//! offered — make the reducer idempotent, OR have the writer skip events
17//! already reflected in the projection — we chose **idempotent reducer**: the
18//! existence/terminal guards already present here mean the catch-up replay can
19//! re-fold *any* tail event (one whose projection landed before a crash, or one
20//! whose projection did not) with the same no-op-or-apply outcome, so the
21//! writer needs no per-event "already applied?" probe. The watermark advances
22//! only after an event's projections are fsynced, so it can lag the projections
23//! but never lead them.
24//!
25//! ## Manifest counters are derived, not folded
26//!
27//! The reducers here deliberately do **not** touch the manifest's denormalized
28//! `node_count` counter. It is
29//! recomputed from the projection directories by
30//! [`derive_counters`](crate::projections), invoked from
31//! [`advance_applied_seq`](crate::events) at the
32//! watermark advance. An earlier design incremented/decremented them inside
33//! these reducers, but a crash between a projection write and the follow-on
34//! `manifest.json` write could permanently desync them: the replay re-folded
35//! the event, hit the `*.created`/terminal idempotency guard above, and skipped
36//! the counter mutation that never actually landed. Deriving the counts makes
37//! drift impossible — there is no delta to lose. See issue
38//! `manifest-counter-desync`. A count-affecting reducer still emits its manifest
39//! op to refresh `updated_at`; the counter fields it carries are overwritten by
40//! the derive step.
41//!
42//! Because a count-affecting event rewrites a projection file *and* the
43//! manifest counter under the same exclusive `flock`, a reader that scans both
44//! together must hold the shared `flock` (`LOCK_SH`) for the whole scan or it
45//! could see the projection change without the matching counter (or vice
46//! versa). See [`crate::projections`] and design.md §4.
47
48use std::path::{Path, PathBuf};
49
50use chrono::{DateTime, Utc};
51use serde_json::Value;
52
53use crate::error::{Error, Result};
54use crate::paths::RunPaths;
55use crate::projections::{read_manifest_opt, read_node_opt, write_manifest, write_node};
56use crate::report::ReportOrigin;
57use crate::schema::{
58    ChildRef, Event, EvidenceStatus, IdValidationError, Kind, Lifecycle, Manifest, MergeTxn, Node,
59    NodeId, RunId, Status, TmuxIdentity, WorkerEvidence, WorkerExit, STATE_SCHEMA_VERSION,
60};
61
62/// Map an id-validation failure on an event-sourced id to a [`CorruptEventLog`]
63/// error. An id that fails to parse here came off `events.jsonl` (or a forged
64/// event), so the log — not the caller — is the corrupt party.
65///
66/// [`CorruptEventLog`]: Error::CorruptEventLog
67fn corrupt_id(events_path: &Path, ev: &Event, e: &IdValidationError) -> Error {
68    Error::CorruptEventLog {
69        path: events_path.to_path_buf(),
70        reason: format!("event seq={} kind={}: {e}", ev.seq, ev.kind),
71    }
72}
73
74/// Parse an optional `RunId` from event-data field `field`: missing/null →
75/// `None`; a JSON string → validated `Some(RunId)`; a malformed id or a
76/// non-string value → [`Error::CorruptEventLog`].
77fn opt_run_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<RunId>> {
78    match d.get(field) {
79        None | Some(Value::Null) => Ok(None),
80        Some(Value::String(s)) => RunId::parse_str(s)
81            .map(Some)
82            .map_err(|e| corrupt_id(events_path, ev, &e)),
83        Some(_) => Err(Error::CorruptEventLog {
84            path: events_path.to_path_buf(),
85            reason: format!(
86                "event seq={} kind={} `{field}` must be a JSON string or null",
87                ev.seq, ev.kind
88            ),
89        }),
90    }
91}
92
93/// Parse an optional `NodeId` from event-data field `field`. See [`opt_run_id`].
94fn opt_node_id(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<NodeId>> {
95    match d.get(field) {
96        None | Some(Value::Null) => Ok(None),
97        Some(Value::String(s)) => NodeId::parse_str(s)
98            .map(Some)
99            .map_err(|e| corrupt_id(events_path, ev, &e)),
100        Some(_) => Err(Error::CorruptEventLog {
101            path: events_path.to_path_buf(),
102            reason: format!(
103                "event seq={} kind={} `{field}` must be a JSON string or null",
104                ev.seq, ev.kind
105            ),
106        }),
107    }
108}
109
110/// Parse a `kind` value from event `data` for a NEW append, failing closed on
111/// anything not a live, creatable kind.
112///
113/// `Kind`'s `#[serde(other)]` catch-all means every unrecognized string —
114/// a removed kind (`code`, `orchestrate`, …), a typo, or a future kind —
115/// deserializes to [`Kind::Unknown`] rather than erroring. That read-only
116/// catch-all exists so `run list` / `doctor` can decode a legacy on-disk run
117/// (ADR §D7); it must NOT let a garbage `kind` slip through the append gate as
118/// though it were valid. Mapping `Unknown` back to `None` keeps the reducer's
119/// `run.created` / `node.created` / `child.spawned` validation fail-closed, as
120/// it was before the 0.2 cut added the catch-all. (Legacy runs are never
121/// re-created through this path — their manifest/nodes already exist on disk and
122/// are read directly, not replayed from a fresh `*.created`.)
123fn data_kind(v: &Value) -> Option<Kind> {
124    match serde_json::from_value::<Kind>(v.clone()) {
125        Ok(Kind::Unknown) | Err(_) => None,
126        Ok(k) => Some(k),
127    }
128}
129
130fn data_status(v: &Value) -> Option<Status> {
131    serde_json::from_value(v.clone()).ok()
132}
133
134fn require_status(ev: &Event, path: PathBuf) -> Result<Status> {
135    data_status(ev.data.get("status").unwrap_or(&Value::Null)).ok_or_else(|| {
136        Error::CorruptEventLog {
137            path,
138            reason: format!("{} missing/invalid `status`", ev.kind),
139        }
140    })
141}
142
143fn want_str<'a>(events_path: &Path, ev: &Event, d: &'a Value, field: &str) -> Result<&'a str> {
144    d.get(field)
145        .and_then(Value::as_str)
146        .ok_or_else(|| Error::CorruptEventLog {
147            path: events_path.to_path_buf(),
148            reason: format!(
149                "event seq={} kind={} missing `{field}` string field",
150                ev.seq, ev.kind
151            ),
152        })
153}
154
155/// Read an optional boolean field with strict typing: missing/null → `None`,
156/// JSON bool → `Some(b)`, anything else → `CorruptEventLog`. Mirrors
157/// [`optional_str`] / [`optional_i32`]; prevents a non-boolean `success` /
158/// `cancelled` from being silently coerced to `false` and bypassing the
159/// success-XOR-cancelled invariant.
160fn optional_bool(events_path: &Path, ev: &Event, d: &Value, field: &str) -> Result<Option<bool>> {
161    match d.get(field) {
162        None | Some(Value::Null) => Ok(None),
163        Some(Value::Bool(b)) => Ok(Some(*b)),
164        Some(_) => Err(Error::CorruptEventLog {
165            path: events_path.to_path_buf(),
166            reason: format!(
167                "event seq={} kind={} `{field}` must be a JSON boolean or null",
168                ev.seq, ev.kind
169            ),
170        }),
171    }
172}
173
174fn optional_i32(d: &Value, field: &str, events_path: &Path, ev: &Event) -> Result<Option<i32>> {
175    match d.get(field) {
176        None | Some(Value::Null) => Ok(None),
177        Some(v) => {
178            let raw = v.as_i64().ok_or_else(|| Error::CorruptEventLog {
179                path: events_path.to_path_buf(),
180                reason: format!(
181                    "event seq={} kind={} `{field}` must be integer",
182                    ev.seq, ev.kind
183                ),
184            })?;
185            i32::try_from(raw)
186                .map(Some)
187                .map_err(|_| Error::CorruptEventLog {
188                    path: events_path.to_path_buf(),
189                    reason: format!(
190                        "event seq={} kind={} `{field}` out of i32 range: {raw}",
191                        ev.seq, ev.kind
192                    ),
193                })
194        }
195    }
196}
197
198fn optional_ts(
199    d: &Value,
200    field: &str,
201    events_path: &Path,
202    ev: &Event,
203) -> Result<Option<DateTime<Utc>>> {
204    match d.get(field) {
205        None | Some(Value::Null) => Ok(None),
206        Some(Value::String(s)) => DateTime::parse_from_rfc3339(s)
207            .map(|dt| Some(dt.with_timezone(&Utc)))
208            .map_err(|_| Error::CorruptEventLog {
209                path: events_path.to_path_buf(),
210                reason: format!(
211                    "event seq={} kind={} `{field}` not RFC3339",
212                    ev.seq, ev.kind
213                ),
214            }),
215        Some(_) => Err(Error::CorruptEventLog {
216            path: events_path.to_path_buf(),
217            reason: format!(
218                "event seq={} kind={} `{field}` must be RFC3339 string or null",
219                ev.seq, ev.kind
220            ),
221        }),
222    }
223}
224
225/// A projection write planned by [`reduce_event_to_ops`] and performed by
226/// [`commit_ops`].
227///
228/// Splitting the reducer into a pure *plan* phase (compute these ops from the
229/// current projection state, validating as it goes) and a *commit* phase
230/// (write them) means a single branch per kind implements both the pre-append
231/// validation gate and the post-append apply — there is no validate/apply
232/// mirror to drift out of lockstep, and the projection state is read once
233/// rather than twice.
234pub(crate) enum ProjectionOp {
235    /// Write the run manifest.
236    Manifest(Manifest),
237    /// Write a node projection.
238    Node(Node),
239}
240
241/// Commit a planned batch of projection writes, in order.
242///
243/// Caller must hold the run's [`crate::lock::RunLock`]. Pairs with
244/// [`reduce_event_to_ops`]: the ops were computed against the same locked
245/// state, and nothing mutates the projections between the plan and this commit
246/// (in the append path only `events.jsonl` is written in between), so the
247/// planned writes are still valid.
248pub(crate) fn commit_ops(paths: &RunPaths, ops: Vec<ProjectionOp>) -> Result<()> {
249    for op in ops {
250        match op {
251            ProjectionOp::Manifest(m) => write_manifest(paths, &m)?,
252            ProjectionOp::Node(n) => write_node(paths, &n)?,
253        }
254    }
255    Ok(())
256}
257
258/// Plan the projection writes one event implies, *without* performing them.
259///
260/// This is the single source of truth for both validation and application: it
261/// reads the current projection state, enforces every event-payload invariant
262/// (returning [`Error::CorruptEventLog`] for a malformed or cross-run event),
263/// and returns the exact [`ProjectionOp`]s to commit (empty for a no-op or an
264/// unknown `kind`). Because it never writes, it is also the transactional gate
265/// run *before* the durable append in
266/// [`crate::events::append_and_apply_unlocked`]: a reducer-rejected event is
267/// caught here and never reaches `events.jsonl`, so a later replay /
268/// `rebuild_projections` can't trip over a poison line. The state-dependent
269/// no-op guards live here too (a settled node/run/discussion swallows a late
270/// or even malformed event as a clean no-op rather than erroring).
271///
272/// Caller must hold the run's [`crate::lock::RunLock`].
273pub(crate) fn reduce_event_to_ops(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
274    // An event whose envelope `run_id` doesn't match the run we're folding it
275    // into means the log was copied/misrouted — folding it would silently
276    // cross-contaminate projections. Reject before planning anything.
277    if ev.run_id != paths.run_id {
278        return Err(Error::CorruptEventLog {
279            path: paths.events(),
280            reason: format!(
281                "event seq={} envelope run_id {:?} does not match run {:?}",
282                ev.seq,
283                ev.run_id.as_str(),
284                paths.run_id.as_str()
285            ),
286        });
287    }
288    // Each event kind is listed explicitly as documentation of the known set;
289    // `supervisor.exited` and the `_` fallthrough share a body intentionally.
290    #[allow(clippy::match_same_arms)]
291    match ev.kind.as_str() {
292        "run.created" => reduce_run_created(paths, ev),
293        "run.status" => reduce_run_status(paths, ev),
294        "node.created" => reduce_node_created(paths, ev),
295        "node.status" => reduce_node_status(paths, ev),
296        "node.report" => reduce_node_report(paths, ev),
297        "node.retry" => reduce_node_retry(paths, ev),
298        "worker.exited" => reduce_worker_exited(paths, ev),
299        "worker.evidence.archived" => reduce_worker_evidence_archived(paths, ev),
300        "worker.evidence.failed" => reduce_worker_evidence_failed(paths, ev),
301        "node.death_observed" => reduce_node_death_observed(paths, ev),
302        "node.awaiting_input" => reduce_node_awaiting_input(paths, ev),
303        "node.input_resolved" => reduce_node_input_resolved(paths, ev),
304        KIND_MERGE_STARTED => reduce_merge_started(paths, ev),
305        KIND_MERGE_ABORTED => reduce_merge_aborted(paths, ev),
306        "child.spawned" => reduce_child_spawned(paths, ev),
307        "supervisor.attached" => reduce_supervisor_attached(paths, ev),
308        "supervisor.cursor_advanced" => reduce_supervisor_cursor_advanced(paths, ev),
309        "supervisor.exited" => Ok(vec![]),
310        // Append-only audit records from `/orchestrate` (decision log +
311        // pakkopysäytys). They mutate no projection — the event log is their
312        // canonical home — so they fold to a clean no-op. Listed explicitly
313        // (rather than relying on the `_` fallthrough) so the append path's
314        // transactional gate runs the same no-op plan for them and the intent
315        // is documented at the match site. They are NOT `node.report`, so the
316        // supervisor never mistakes them for a terminal signal.
317        "orchestrator.decision" | "discuss.critical" => Ok(vec![]),
318        // At-most-once marker the supervisor appends the first time a run is
319        // observed terminal, gating the `run create --notify` completion hook so
320        // a restart never re-fires it (issue `no-completion-notification-to-parent`).
321        // Mutates no projection — the event log is its only home — so it folds to
322        // a clean no-op. Listed explicitly so the append path's transactional gate
323        // runs the same no-op plan and the intent is documented here.
324        "run.notified" | "run.awaiting_input_notified" => Ok(vec![]),
325        // Best-effort teardown audit records from the supervisor's cleanup
326        // path. Each mutates no projection — the event log is their only home —
327        // so they fold to a clean no-op. Listed explicitly so the append path's
328        // transactional gate runs the same no-op plan and the intent is
329        // documented here.
330        //   - `cleanup.window_missing`: the node's tmux window could not be
331        //     located to close it (typically a manually-resolved rebase renamed
332        //     the window — issue `worktree-merge-orphans-tmux-window`).
333        //   - `cleanup.worktree_missing`: the worktree dir was already gone at
334        //     teardown (e.g. removed manually), so nothing to `worktree remove`.
335        //   - `cleanup.branch_remove_failed`: `git branch -{d,D}` refused (e.g.
336        //     unmerged commits, or the branch is already gone); the run completes
337        //     anyway (issue `supervisor-worktree-remove-no-force`).
338        //   - `cleanup.branch_preserved`: a BLOCKED terminal report
339        //     (`success: false`, no explicit merge) intentionally left the branch
340        //     and worktree in place for the human to pick up, instead of tearing
341        //     them down (issue `blocked-report-deletes-branch`).
342        //   - `cleanup.session_killed`: the run's managed `--headless` tmux
343        //     session was torn down once its last managed window was gone, so an
344        //     empty session is not left behind (issue
345        //     `headless-tmux-session-not-torn-down`).
346        //   - `cleanup.session_retained`: the same teardown was skipped because a
347        //     human had attached to the session — never yanked out from under
348        //     them.
349        "cleanup.window_missing"
350        | "cleanup.worktree_missing"
351        | "cleanup.branch_remove_failed"
352        | "cleanup.branch_preserved"
353        | "cleanup.discard_authorized"
354        | "cleanup.session_killed"
355        | "cleanup.session_retained" => Ok(vec![]),
356        // Data-integrity audit record: the supervisor found a persisted child
357        // run id (in `supervisor.state.json`'s `spawned_children`) that fails
358        // `RunId` structural validation and quarantined it — a corrupt id that
359        // would otherwise resolve with `.ok()` and be silently skipped every
360        // tick, indistinguishable from a child that completed and was torn down
361        // (issue `wildly-glorious-food`). It mutates no projection — the event
362        // log is its only home — so it folds to a clean no-op. Listed
363        // explicitly so the append path's transactional gate runs the same
364        // no-op plan and the intent is documented here.
365        "supervisor.child_id_quarantined" => Ok(vec![]),
366        _ => Ok(vec![]),
367    }
368}
369
370/// The projection file [`commit_ops`] writes for `op`, keyed exactly as the
371/// `write_*` helpers key it internally. Shared by [`plan_projections`] (which
372/// reports the path) and conceptually by [`commit_ops`] (which writes it), so
373/// the enumerated path list can never name a different file than the one the
374/// reducer actually fsyncs.
375fn op_path(paths: &RunPaths, op: &ProjectionOp) -> PathBuf {
376    match op {
377        ProjectionOp::Manifest(_) => paths.manifest(),
378        ProjectionOp::Node(n) => paths.node(&n.node_id),
379    }
380}
381
382/// Enumerate the projection files the reducer would write for `event`, in the
383/// order `commit_ops` would write them, *without* performing any write.
384///
385/// This is the single source of truth that ends the CLI/reducer divergence the
386/// `projected-paths-into-reducer` issue describes: rather than a hand-maintained
387/// list in `taskfleet` that drifts whenever a new projection is added, both the
388/// reducer and a caller's preflight (`event create --dry-run`) read the *same*
389/// `reduce_event_to_ops` plan. This function maps that plan to file paths;
390/// `apply_event` commits it. A new projection added to a reducer arm is
391/// therefore reflected here automatically.
392///
393/// Because it runs the real reducer plan against current projection state, the
394/// result is exact, not a guess: a state-dependent no-op (a settled node, an
395/// already-created projection, a terminal-guarded transition) yields an empty
396/// list — precisely the files `apply_event` would touch, which is none. A
397/// malformed-payload event surfaces the same [`Error::CorruptEventLog`] the
398/// real apply would, so a dry-run preflight cannot report success for an event
399/// the write path would reject.
400///
401/// Caller should hold the run's [`crate::lock::RunLock`] for a snapshot
402/// consistent with a concurrent reducer; a lock-free read is best-effort.
403pub fn plan_projections(paths: &RunPaths, event: &Event) -> Result<Vec<PathBuf>> {
404    let ops = reduce_event_to_ops(paths, event)?;
405    Ok(ops.iter().map(|op| op_path(paths, op)).collect())
406}
407
408/// Apply one event to projections: plan via [`reduce_event_to_ops`], then
409/// [`commit_ops`]. No-op for unknown `kind`. Caller must hold the run's
410/// [`crate::lock::RunLock`].
411///
412/// Shares the one [`reduce_event_to_ops`] plan with [`plan_projections`]: the
413/// paths that function reports are exactly the files this one fsyncs, because
414/// both consume the same `ProjectionOp` vector (this commits it; that maps it to
415/// paths via [`op_path`]).
416///
417/// `pub(crate)`: applying an event in isolation (without the matching
418/// `events.jsonl` append) is an internal building block used by `cancel` (to
419/// re-fold a crash-stranded event) and a future `rebuild_projections_from_events`.
420/// External callers mutate state through
421/// [`crate::events::append_and_apply_event`] so the log and projections can
422/// never diverge.
423pub(crate) fn apply_event(paths: &RunPaths, ev: &Event) -> Result<()> {
424    let ops = reduce_event_to_ops(paths, ev)?;
425    commit_ops(paths, ops)
426}
427
428/// Validate one persisted event through the normal reducer plan without
429/// committing projection writes. It performs the same checks as projection
430/// application and then discards the planned writes. Migration preflight uses
431/// this under the run's exclusive lock so it never invents a second
432/// event-payload validator.
433pub fn validate_event(paths: &RunPaths, ev: &Event) -> Result<()> {
434    reduce_event_to_ops(paths, ev).map(|_| ())
435}
436
437/// The envelope `node_id` that a `node.*` event must carry, with the same
438/// `CorruptEventLog` message `apply_*` produces. Shared by validate/apply so
439/// the missing-id check can't drift between them.
440fn require_envelope_node_id(events_path: &Path, ev: &Event) -> Result<NodeId> {
441    ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
442        path: events_path.to_path_buf(),
443        reason: format!(
444            "event seq={} kind={} missing top-level `node_id`",
445            ev.seq, ev.kind
446        ),
447    })
448}
449
450fn reduce_run_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
451    // Idempotent: a replayed `run.created` against an existing manifest is a
452    // no-op (but validates that `run_id` matches; otherwise the event log
453    // is being applied to the wrong run).
454    if let Some(existing) = read_manifest_opt(paths)? {
455        if existing.run_id != ev.run_id {
456            return Err(Error::CorruptEventLog {
457                path: paths.manifest(),
458                reason: format!(
459                    "run.created run_id={} conflicts with existing manifest run_id={}",
460                    ev.run_id, existing.run_id
461                ),
462            });
463        }
464        return Ok(vec![]);
465    }
466    let events_path = paths.events();
467    let d = &ev.data;
468    let kind =
469        data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
470            path: events_path.clone(),
471            reason: "run.created missing/invalid `kind`".into(),
472        })?;
473    let lifecycle: Lifecycle = serde_json::from_value(
474        d.get("lifecycle").cloned().unwrap_or(Value::Null),
475    )
476    .map_err(|_| Error::CorruptEventLog {
477        path: events_path.clone(),
478        reason: "run.created missing/invalid `lifecycle`".into(),
479    })?;
480    let title = want_str(&events_path, ev, d, "title")?.to_string();
481    let agent_selection: Option<crate::schema::AgentSelection> = d
482        .get("agent_selection")
483        .cloned()
484        .map(serde_json::from_value)
485        .transpose()
486        .map_err(|e| Error::CorruptEventLog {
487            path: events_path.clone(),
488            reason: format!("run.created invalid `agent_selection`: {e}"),
489        })?;
490    if let Some(selection) = &agent_selection {
491        selection
492            .validate()
493            .map_err(|reason| Error::CorruptEventLog {
494                path: events_path.clone(),
495                reason: format!("run.created invalid `agent_selection`: {reason}"),
496            })?;
497    }
498    let m = Manifest {
499        schema_version: STATE_SCHEMA_VERSION,
500        // Created at the watermark floor; the append path advances it to this
501        // event's `seq` (after the manifest is fsynced) in `advance_applied_seq`.
502        applied_seq: 0,
503        // `run_id == paths.run_id` was verified at `reduce_event_to_ops` entry.
504        run_id: paths.run_id.clone(),
505        kind,
506        lifecycle,
507        title,
508        status: Status::Pending,
509        created_at: ev.ts,
510        updated_at: ev.ts,
511        source_repo: d
512            .get("source_repo")
513            .and_then(Value::as_str)
514            .map(str::to_string),
515        source_branch: d
516            .get("source_branch")
517            .and_then(Value::as_str)
518            .map(str::to_string),
519        worktree_root: d
520            .get("worktree_root")
521            .and_then(Value::as_str)
522            .map(str::to_string),
523        managed_tmux_session: d
524            .get("managed_tmux_session")
525            .and_then(Value::as_str)
526            .map(str::to_string),
527        notify_cmd: d
528            .get("notify_cmd")
529            .and_then(Value::as_str)
530            .map(str::to_string),
531        harness: d.get("harness").and_then(Value::as_str).map(str::to_string),
532        agent_selection,
533        node_count: 0,
534        parent_run_id: opt_run_id(&events_path, ev, d, "parent_run_id")?,
535        parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
536    };
537    Ok(vec![ProjectionOp::Manifest(m)])
538}
539
540fn reduce_run_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
541    let mut m = match read_manifest_opt(paths)? {
542        Some(m) => m,
543        None => return Ok(vec![]),
544    };
545    let new_status = require_status(ev, paths.events())?;
546    // Terminal-state guard with one narrow recovery: Failed -> Done is allowed
547    // only when this event names an authoritative merge report that was already
548    // adopted by the log-derived node fold before this event, every own node is
549    // now Done, and the supervisor attests that every linked child is Done.
550    // Cancelled (and every other terminal transition) remains immutable.
551    if m.status.is_terminal() {
552        let recovery_seq = ev
553            .data
554            .get("recovery_merge_report_seq")
555            .and_then(Value::as_u64);
556        let children_successful =
557            ev.data.get("children_successful").and_then(Value::as_bool) == Some(true);
558        let recovery_allowed = m.status == Status::Failed
559            && new_status == Status::Done
560            && children_successful
561            && recovery_seq.is_some_and(|wanted| {
562                crate::cancel::read_node_status_facts(paths, Some(ev.seq))
563                    .ok()
564                    .filter(|facts| {
565                        crate::aggregate_terminal_status(facts.iter().map(|fact| fact.status))
566                            == Some(Status::Done)
567                    })
568                    .is_some_and(|facts| {
569                        facts
570                            .iter()
571                            .any(|fact| fact.confirmed_merge_seq == Some(wanted))
572                    })
573            });
574        if !recovery_allowed {
575            trace_terminal_noop(ev, m.status, new_status);
576            return Ok(vec![]);
577        }
578    }
579    if m.status == new_status {
580        return Ok(vec![]);
581    }
582    m.status = new_status;
583    m.updated_at = ev.ts;
584    Ok(vec![ProjectionOp::Manifest(m)])
585}
586
587/// Reconstruct the fully-qualified tmux identity from `node.created` event
588/// data. Returns `Some` only when both `tmux_session` and `tmux_window_id` are
589/// present and non-empty — the minimum needed to match a window. `tmux_socket`
590/// is optional (a default-socket spawn may emit null); an empty socket is
591/// normalized to `None` so the watchdog never invokes `tmux -S ""`.
592/// `tmux_pane_id` is likewise optional (create.sh predating it emits nothing);
593/// agent-log capture falls back to the window's active pane when absent. Legacy
594/// events from a create.sh that predates the qualified fields (or that emit a
595/// partial/empty identity) yield `None`, so the node falls back to bare-name
596/// matching on `tmux_window`.
597fn tmux_identity_from_data(d: &Value) -> Option<TmuxIdentity> {
598    let nonempty = |key| {
599        d.get(key)
600            .and_then(Value::as_str)
601            .map(str::trim)
602            .filter(|s| !s.is_empty())
603            .map(str::to_string)
604    };
605    let session = nonempty("tmux_session")?;
606    let window_id = nonempty("tmux_window_id")?;
607    Some(TmuxIdentity {
608        socket: nonempty("tmux_socket"),
609        session,
610        window_id,
611        // Optional: create.sh predating the field (or a failed pane query)
612        // emits no `tmux_pane_id`; capture then falls back to `window_id`.
613        pane_id: nonempty("tmux_pane_id"),
614    })
615}
616
617fn worker_evidence_from_spawn_data(
618    events_path: &Path,
619    ev: &Event,
620    data: &Value,
621) -> Result<Option<WorkerEvidence>> {
622    let Some(session_id) = data.get("pi_session_id") else {
623        if data.get("pi_session_path").is_some() || data.get("pi_session_cwd").is_some() {
624            return Err(Error::CorruptEventLog {
625                path: events_path.to_path_buf(),
626                reason: format!(
627                    "event seq={} Pi evidence fields must be all-or-none",
628                    ev.seq
629                ),
630            });
631        }
632        return Ok(None);
633    };
634    if session_id.is_null() {
635        if data
636            .get("pi_session_path")
637            .is_some_and(|value| !value.is_null())
638            || data
639                .get("pi_session_cwd")
640                .is_some_and(|value| !value.is_null())
641        {
642            return Err(Error::CorruptEventLog {
643                path: events_path.to_path_buf(),
644                reason: format!(
645                    "event seq={} Pi evidence fields must be all-or-none",
646                    ev.seq
647                ),
648            });
649        }
650        return Ok(None);
651    }
652    let session_id = session_id
653        .as_str()
654        .filter(|v| !v.is_empty())
655        .ok_or_else(|| Error::CorruptEventLog {
656            path: events_path.to_path_buf(),
657            reason: format!(
658                "event seq={} pi_session_id must be a non-empty string",
659                ev.seq
660            ),
661        })?;
662    if session_id.len() != 36
663        || !session_id.bytes().enumerate().all(|(index, byte)| {
664            if matches!(index, 8 | 13 | 18 | 23) {
665                byte == b'-'
666            } else {
667                byte.is_ascii_hexdigit()
668            }
669        })
670    {
671        return Err(Error::CorruptEventLog {
672            path: events_path.to_path_buf(),
673            reason: format!("event seq={} pi_session_id must be a UUID", ev.seq),
674        });
675    }
676    let original_cwd = want_str(events_path, ev, data, "pi_session_cwd")?;
677    let live_session_path = want_str(events_path, ev, data, "pi_session_path")?;
678    let expected_live_path = format!(
679        ".creating/pi-sessions/{}/pi-session-{session_id}.jsonl",
680        ev.run_id.as_str()
681    );
682    if live_session_path != expected_live_path {
683        return Err(Error::CorruptEventLog {
684            path: events_path.to_path_buf(),
685            reason: format!(
686                "event seq={} pi_session_path is not the canonical state-relative path",
687                ev.seq
688            ),
689        });
690    }
691    let attempt = match data.get("attempt") {
692        None | Some(Value::Null) => 0,
693        Some(value) => value
694            .as_u64()
695            .and_then(|raw| u32::try_from(raw).ok())
696            .ok_or_else(|| Error::CorruptEventLog {
697                path: events_path.to_path_buf(),
698                reason: format!("event seq={} attempt must be a u32", ev.seq),
699            })?,
700    };
701    Ok(Some(WorkerEvidence {
702        attempt,
703        session_id: session_id.to_string(),
704        original_cwd: original_cwd.to_string(),
705        live_session_path: live_session_path.to_string(),
706        status: EvidenceStatus::Pending,
707        transcript_path: None,
708        resume_path: None,
709        pane_path: None,
710        report_path: None,
711        transcript_sha256: None,
712        error: None,
713    }))
714}
715
716fn reduce_node_created(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
717    let events_path = paths.events();
718    // The envelope `node_id` is already a validated `NodeId` (parsed on read),
719    // so take it directly — no re-parse needed.
720    let node_id = require_envelope_node_id(&events_path, ev)?;
721    let is_default_node = node_id.as_str() == "n-0001";
722    // Idempotent on replay: skip if the node already exists.
723    if read_node_opt(paths, &node_id)?.is_some() {
724        return Ok(vec![]);
725    }
726    let d = &ev.data;
727    let kind =
728        data_kind(d.get("kind").unwrap_or(&Value::Null)).ok_or_else(|| Error::CorruptEventLog {
729            path: events_path.clone(),
730            reason: format!(
731                "event seq={} kind=node.created missing/invalid `kind`",
732                ev.seq
733            ),
734        })?;
735    let n = Node {
736        schema_version: STATE_SCHEMA_VERSION,
737        node_id,
738        // `run_id == paths.run_id` was verified at `reduce_event_to_ops` entry.
739        run_id: paths.run_id.clone(),
740        parent_node_id: opt_node_id(&events_path, ev, d, "parent_node_id")?,
741        kind,
742        status: Status::Pending,
743        task: d.get("task").and_then(Value::as_str).map(str::to_string),
744        worktree_path: d
745            .get("worktree_path")
746            .and_then(Value::as_str)
747            .map(str::to_string),
748        branch: d.get("branch").and_then(Value::as_str).map(str::to_string),
749        base_sha: d
750            .get("base_sha")
751            .and_then(Value::as_str)
752            .filter(|s| !s.is_empty())
753            .map(str::to_string),
754        tmux_window: d
755            .get("tmux_window")
756            .and_then(Value::as_str)
757            .map(str::to_string),
758        tmux_identity: tmux_identity_from_data(d),
759        evidence: worker_evidence_from_spawn_data(&events_path, ev, d)?,
760        agent_pid: optional_i32(d, "agent_pid", &events_path, ev)?,
761        agent_pid_start_time: optional_ts(d, "agent_pid_start_time", &events_path, ev)?,
762        supervisor_pid: optional_i32(d, "supervisor_pid", &events_path, ev)?,
763        children: Vec::new(),
764        started_at: Some(ev.ts),
765        updated_at: ev.ts,
766        last_report: None,
767        last_processed_report_seq_by_child: serde_json::Map::default(),
768        retry_attempts: 0,
769        worker_exit: None,
770        pending_merge: None,
771        first_death_at: None,
772        awaiting_input: None,
773    };
774    let mut ops = vec![ProjectionOp::Node(n)];
775    if let Some(mut m) = read_manifest_opt(paths)? {
776        // Materialization is the point at which an implicit source branch is
777        // known. Preserve an explicit run.created value; otherwise fold the
778        // source discovered by the creator into the manifest in this same
779        // locked event application.
780        if is_default_node && m.source_branch.is_none() {
781            m.source_branch = d
782                .get("source_branch")
783                .and_then(Value::as_str)
784                .filter(|branch| !branch.is_empty())
785                .map(str::to_string);
786        }
787        // `node_count` is derived from the projection directories in
788        // `advance_applied_seq`, never incremented here — see the module note
789        // and issue `manifest-counter-desync`. This op also refreshes the run's
790        // last-activity timestamp.
791        m.updated_at = ev.ts;
792        ops.push(ProjectionOp::Manifest(m));
793    }
794    Ok(ops)
795}
796
797/// Rewire an existing node to a freshly re-spawned agent after an empty-handed
798/// `agent-died` bounded auto-retry (issue `autoretry-agent-died-worker`). The
799/// supervisor tore down the dead worker's stale worktree and `create.sh`'d a
800/// clean one at the run's source branch; this event carries the new spawn
801/// metadata (`branch`, `base_sha`, `worktree_path`, tmux identity, `agent_pid`)
802/// plus the audit fields (`attempt`, `reason`).
803///
804/// It updates the node in place: the new agent's coordinates replace the dead
805/// one's, `status` returns to `Pending`, `started_at` is re-stamped so the
806/// watchdog's spawn-grace window re-applies to the new agent, `last_report` is
807/// cleared, and `retry_attempts` is incremented — the DURABLE, restart-safe
808/// bound the watchdog checks before scheduling the next retry.
809///
810/// Guards, mirroring the other node reducers:
811/// - A missing node is a no-op (a retry event whose node was never created).
812/// - A TERMINAL node is never resurrected (a settled node is frozen): if a real
813///   `node.report` raced in and terminalized the node, the retry is a dead event.
814///   This keeps replay robust and preserves the terminal-state invariant.
815fn reduce_node_retry(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
816    let events_path = paths.events();
817    let node_id = require_envelope_node_id(&events_path, ev)?;
818    let mut n = match read_node_opt(paths, &node_id)? {
819        Some(n) => n,
820        None => return Ok(vec![]),
821    };
822    // Terminal-state guard: a settled node is frozen. A late `node.report` that
823    // beat this retry to the lock wins; the retry must not resurrect it.
824    if n.status.is_terminal() {
825        tracing::debug!(
826            target: "taskfleet_core::reducer",
827            seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
828            "no-op: node.retry against terminal node"
829        );
830        return Ok(vec![]);
831    }
832    let d = &ev.data;
833    // Rewire to the new agent. Each field mirrors `reduce_node_created`'s parsing
834    // so the projection shape is identical to a fresh spawn.
835    n.branch = d.get("branch").and_then(Value::as_str).map(str::to_string);
836    n.base_sha = d
837        .get("base_sha")
838        .and_then(Value::as_str)
839        .filter(|s| !s.is_empty())
840        .map(str::to_string);
841    n.worktree_path = d
842        .get("worktree_path")
843        .and_then(Value::as_str)
844        .map(str::to_string);
845    n.tmux_window = d
846        .get("tmux_window")
847        .and_then(Value::as_str)
848        .map(str::to_string);
849    n.tmux_identity = tmux_identity_from_data(d);
850    n.evidence = worker_evidence_from_spawn_data(&events_path, ev, d)?;
851    n.agent_pid = optional_i32(d, "agent_pid", &events_path, ev)?;
852    n.agent_pid_start_time = optional_ts(d, "agent_pid_start_time", &events_path, ev)?;
853    n.status = Status::Pending;
854    n.started_at = Some(ev.ts);
855    n.updated_at = ev.ts;
856    n.last_report = None;
857    // Drop any in-flight merge transaction from the PREVIOUS attempt: the retry
858    // rewires the node to a new branch/worktree/agent, so a `pending_merge` that
859    // referenced the dead attempt's branch must not carry forward — recovery would
860    // otherwise judge the new attempt from the old worker's merge state (issue
861    // `merge-transaction-recovery`, /llm-review finding).
862    n.pending_merge = None;
863    // Clear the previous attempt's told exit fact: the freshly re-spawned worker
864    // is a NEW process, so a stale `worker_exit` must not carry over — otherwise
865    // the supervisor's told-fact pass would instantly (mis)judge the new attempt
866    // from the dead one's exit (issue `thin-exit-status-launcher`).
867    n.worker_exit = None;
868    // Clear the previous attempt's first-death anchor: the residual crash backstop
869    // must measure the NEW attempt's own post-death grace from scratch, not inherit
870    // the dead attempt's timestamp (which would fire the backstop with no grace on
871    // the fresh worker's first confirmed death). Issue `typed-supervisor-outcomes`.
872    n.first_death_at = None;
873    // A retry is a new worker attempt. Never carry an unresolved question from
874    // the dead attempt onto the replacement worker.
875    n.awaiting_input = None;
876    // The event carries its ABSOLUTE attempt number (the supervisor set it to
877    // `retry_attempts + 1` at emit time). Assign it directly rather than a blind
878    // `+= 1`: this makes the projection a pure function of the event, so a
879    // full replay from seq 0, or a (guarded-against but defensive) double-apply,
880    // converges to the same `retry_attempts` the log declares — the audit count
881    // and the durable bound can never disagree. A legacy/malformed event with no
882    // parseable `attempt` falls back to the monotone increment.
883    n.retry_attempts = d
884        .get("attempt")
885        .and_then(Value::as_u64)
886        .map_or_else(|| n.retry_attempts.saturating_add(1), |a| a as u32);
887    Ok(vec![ProjectionOp::Node(n)])
888}
889
890fn reduce_node_status(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
891    let events_path = paths.events();
892    let node_id = require_envelope_node_id(&events_path, ev)?;
893    let mut n = match read_node_opt(paths, &node_id)? {
894        Some(n) => n,
895        None => return Ok(vec![]),
896    };
897    let new_status = require_status(ev, events_path)?;
898    // Terminal-state guard: a settled node never transitions again. See
899    // run-cli-read/handoff.md D5.
900    if n.status.is_terminal() {
901        trace_terminal_noop(ev, n.status, new_status);
902        return Ok(vec![]);
903    }
904    if n.status == new_status {
905        return Ok(vec![]);
906    }
907    n.status = new_status;
908    // A terminal `node.status` (e.g. a watchdog-synthesized failure) ends the
909    // node's lifecycle, so any in-flight merge transaction is moot — clear it so a
910    // `pending_merge` is not stranded on a terminal node (recovery skips terminal
911    // nodes, so an uncleared one would dangle forever). Issue
912    // `merge-transaction-recovery` (/llm-review finding).
913    if new_status.is_terminal() {
914        n.pending_merge = None;
915        n.awaiting_input = None;
916    }
917    n.updated_at = ev.ts;
918    Ok(vec![ProjectionOp::Node(n)])
919}
920
921fn reduce_node_report(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
922    let events_path = paths.events();
923    let node_id = require_envelope_node_id(&events_path, ev)?;
924    let mut n = match read_node_opt(paths, &node_id)? {
925        Some(n) => n,
926        None => return Ok(vec![]),
927    };
928    // Terminal-state guard *before* payload validation: a node that already
929    // reached a terminal state is settled, so a late-arriving report (e.g. an
930    // agent success racing a `run cancel`) is a dead event — it must not
931    // resurrect the node, and must not even decorate the projection, so
932    // `last_report` is left untouched. Guarding first also keeps replay
933    // robust: a malformed dead report against a settled node is a clean
934    // no-op rather than a `CorruptEventLog` that would brick rebuild of a
935    // log `append_and_apply_event` already committed. See run-cli-read/handoff.md
936    // D5. (3/4 of /llm-review preferred guard-before-validate over the
937    // reverse the issue spec sketched; the required CorruptEventLog cases
938    // all target live nodes, so validation still runs for them.)
939    if n.status.is_terminal() {
940        // ONE exception to the dead-event rule: a late, CONFIRMED explicit-merge
941        // report is adopted even against a terminal node (issue
942        // `reducer-adopt-explicit-merge`). A watchdog `agent-died` false positive
943        // on a long-lived interactive run can terminalize a node BEFORE the user's
944        // `run merge` report arrives; an explicit user merge carries strictly
945        // higher-fidelity ground truth (the branch demonstrably landed in source)
946        // than a watchdog timeout, so it wins. Overwriting `last_report` here is
947        // what lets `any_node_merged_explicitly` see the merge and the SUPERVISOR
948        // — invariant #5's canonical teardown actor — warrant teardown, instead of
949        // the CLI compensating inline (issues `merge-skips-teardown`,
950        // `agent-died-merge-no-teardown-interactive`).
951        //
952        // Scoped tightly, on BOTH sides:
953        //   - incoming: a CONFIRMED SUCCESSFUL explicit merge
954        //     (`via == "explicit-merge"`, `success == true`, not `cancelled`) —
955        //     matches exactly the force-`-D` teardown gate (`node_branch_merged`),
956        //     so a failed/cancelled or non-merge late report never resurrects a
957        //     settled node and unmerged-work preservation is untouched.
958        //   - prior: only a `Failed` or `Done` node (positive whitelist). A
959        //     `Cancelled` terminal is a DELIBERATE `run cancel` teardown, not a
960        //     watchdog false positive, so a later merge does not override it (it
961        //     stays cancelled — matching the existing "late success report keeps the
962        //     cancel" reducer contract). The whitelist (rather than `!= Cancelled`)
963        //     is future-safe: a new deliberate-teardown terminal added later is not
964        //     silently resurrected to Done.
965        // Idempotent: if this exact report is already the node's `last_report`,
966        // re-folding it on replay is a clean no-op (never churns `updated_at`).
967        //
968        // The node reducer does not directly project run status: the supervisor
969        // owns cross-node/child topology. Once that topology is wholly successful,
970        // it emits the narrowly-authorized Failed -> Done recovery `run.status`,
971        // whose reducer independently verifies this adopted merge evidence.
972        if ReportOrigin::permits_terminal_merge_recovery(n.status, &ev.data) {
973            if n.last_report.as_ref() == Some(&ev.data) && n.status == Status::Done {
974                return Ok(vec![]);
975            }
976            tracing::info!(
977                target: "taskfleet_core::reducer",
978                seq = ev.seq, kind = %ev.kind, node_id = %node_id, prior = ?n.status,
979                "adopting late explicit-merge report against terminal node (invariant #5 teardown)"
980            );
981            n.last_report = Some(ev.data.clone());
982            // A confirmed merge is a terminal SUCCESS: the work landed in source.
983            // (A false watchdog `Failed` is corrected to `Done`; a genuine `Done`
984            // stays `Done` with the merge marker adopted so teardown is warranted.)
985            n.status = Status::Done;
986            // The merge completed, so any in-flight merge transaction is resolved:
987            // clear it so recovery does not later re-examine a settled node
988            // (issue `merge-transaction-recovery`).
989            n.pending_merge = None;
990            n.awaiting_input = None;
991            n.updated_at = ev.ts;
992            return Ok(vec![ProjectionOp::Node(n)]);
993        }
994        tracing::debug!(
995            target: "taskfleet_core::reducer",
996            seq = ev.seq, kind = %ev.kind, node_id = %node_id, current = ?n.status,
997            "no-op: node.report against terminal node"
998        );
999        return Ok(vec![]);
1000    }
1001    // Live node: validate the report's terminal outcome. A `node.report`
1002    // must express exactly one terminal outcome — success/failure XOR
1003    // cancellation. Anything else (a bare `{}` with neither, or the
1004    // contradiction `success: true` + `cancelled: true`) is a corrupt event:
1005    // the reducer is the canonical gate, so reject it rather than silently
1006    // leaving the node in a dangling state. See design.md §7.7 and
1007    // node-cli-read/handoff.md D4.
1008    let new_status = report_terminal_status(&events_path, ev)?;
1009    n.last_report = Some(ev.data.clone());
1010    n.status = new_status;
1011    // A terminal report settles any open human-decision request. A blocked
1012    // report still preserves the discussion in `last_report`, while avoiding a
1013    // stale non-terminal awaiting-input flag on the settled node.
1014    n.awaiting_input = None;
1015    // Any terminal outcome resolves an in-flight merge transaction: a successful
1016    // `explicit-merge` report completes it here (the normal, no-crash path), and
1017    // any other terminal report ends the node's lifecycle so no merge recovery
1018    // should later fire (issue `merge-transaction-recovery`).
1019    n.pending_merge = None;
1020    n.updated_at = ev.ts;
1021    Ok(vec![ProjectionOp::Node(n)])
1022}
1023
1024/// Fold a `worker.exited` event onto the node's `worker_exit` field (design.md
1025/// §2.1 / A1). This records the launcher shim's **told** exit status as a
1026/// durable fact; it deliberately does NOT transition `status`. Terminalization
1027/// is the supervisor's decision via the typed outcome table (§2.6) — a non-zero
1028/// or signalled exit becomes `failed`, while a clean exit without a merge stays
1029/// non-terminal (attention-required). Keeping the status transition out of the
1030/// reducer is what lets the clean-but-unmerged worker remain a visible, resumable
1031/// state instead of an auto-failed one.
1032///
1033/// Payload contract: at least one of `exit_code` (JSON integer) or `signal`
1034/// (JSON integer) must be present; a payload carrying neither is a corrupt event
1035/// (the reducer is the canonical gate). The fold is idempotent and **first-write-
1036/// wins**: once `worker_exit` is set, a replay or a spurious duplicate is a clean
1037/// no-op, so a full replay from seq 0 converges to the same recorded fact.
1038fn reduce_worker_exited(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1039    let events_path = paths.events();
1040    let node_id = require_envelope_node_id(&events_path, ev)?;
1041    let attempt = match ev.data.get("attempt") {
1042        None | Some(Value::Null) => 0,
1043        Some(_) => evidence_attempt(&events_path, ev)?,
1044    };
1045    let code = optional_i32(&ev.data, "exit_code", &events_path, ev)?;
1046    let signal = optional_i32(&ev.data, "signal", &events_path, ev)?;
1047    // A worker exit is EXACTLY one of a normal return (code) or a signal death
1048    // (signal). Neither is meaningless; both is contradictory (a process cannot
1049    // both return a code and be killed) — reject either rather than record an
1050    // ambiguous fact the outcome classifier would then have to disambiguate.
1051    match (code, signal) {
1052        (Some(_), None) | (None, Some(_)) => {}
1053        _ => {
1054            return Err(Error::CorruptEventLog {
1055                path: events_path,
1056                reason: format!(
1057                    "event seq={} kind=worker.exited must carry EXACTLY one of `exit_code` or `signal`",
1058                    ev.seq
1059                ),
1060            });
1061        }
1062    }
1063    let mut n = match read_node_opt(paths, &node_id)? {
1064        Some(n) => n,
1065        // No projection to decorate. A `worker.exited` for a node that does not
1066        // exist folds to nothing — consistent with the other node reducers
1067        // (`node.report` / `node.status`). In practice the shim validates the node
1068        // exists before it can record an exit, and normal append ordering always
1069        // places `node.created` first, so this is only hit for a genuinely orphan
1070        // event.
1071        None => return Ok(vec![]),
1072    };
1073    // First-write-wins: the shim fires exactly once per worker, so an existing
1074    // record is a replay/duplicate. Leaving it untouched keeps the fold a pure
1075    // function of the first exit event and never churns `updated_at`.
1076    if attempt != n.retry_attempts || n.worker_exit.is_some() {
1077        return Ok(vec![]);
1078    }
1079    n.worker_exit = Some(WorkerExit {
1080        code,
1081        signal,
1082        at: ev.ts,
1083    });
1084    // A departed worker cannot proceed on its recommended default. Clear its
1085    // open request so clean-exit attention and crash/stall handling remain the
1086    // actionable read-surface verdicts.
1087    n.awaiting_input = None;
1088    n.updated_at = ev.ts;
1089    Ok(vec![ProjectionOp::Node(n)])
1090}
1091
1092fn evidence_attempt(events_path: &Path, ev: &Event) -> Result<u32> {
1093    ev.data
1094        .get("attempt")
1095        .and_then(Value::as_u64)
1096        .and_then(|raw| u32::try_from(raw).ok())
1097        .ok_or_else(|| Error::CorruptEventLog {
1098            path: events_path.to_path_buf(),
1099            reason: format!("event seq={} attempt must be a u32", ev.seq),
1100        })
1101}
1102
1103fn reduce_worker_evidence_archived(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1104    let events_path = paths.events();
1105    let node_id = require_envelope_node_id(&events_path, ev)?;
1106    let attempt = evidence_attempt(&events_path, ev)?;
1107    let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1108    let mut values = Vec::new();
1109    for field in [
1110        "transcript_path",
1111        "resume_path",
1112        "pane_path",
1113        "report_path",
1114        "transcript_sha256",
1115    ] {
1116        values.push(want_str(&events_path, ev, &ev.data, field)?.to_string());
1117    }
1118    let expected_prefix = format!("evidence/{}/", node_id.as_str());
1119    for (index, suffix) in [
1120        "pi-session.original.jsonl",
1121        "pi-session.resume.jsonl",
1122        "final-pane.log",
1123        "terminal-report.json",
1124    ]
1125    .iter()
1126    .enumerate()
1127    {
1128        if values[index] != format!("{expected_prefix}{suffix}") {
1129            return Err(Error::CorruptEventLog {
1130                path: events_path,
1131                reason: format!(
1132                    "event seq={} evidence artifact path is not canonical",
1133                    ev.seq
1134                ),
1135            });
1136        }
1137    }
1138    if values[4].len() != 64
1139        || !values[4]
1140            .bytes()
1141            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1142    {
1143        return Err(Error::CorruptEventLog {
1144            path: events_path,
1145            reason: format!(
1146                "event seq={} transcript_sha256 is not lowercase SHA-256",
1147                ev.seq
1148            ),
1149        });
1150    }
1151    let mut node = match read_node_opt(paths, &node_id)? {
1152        Some(node) => node,
1153        None => return Ok(vec![]),
1154    };
1155    let Some(evidence) = node.evidence.as_mut() else {
1156        return Err(Error::CorruptEventLog {
1157            path: events_path,
1158            reason: format!(
1159                "event seq={} evidence archive has no recorded Pi session",
1160                ev.seq
1161            ),
1162        });
1163    };
1164    if evidence.attempt != attempt || evidence.session_id != session_id {
1165        return Ok(vec![]);
1166    }
1167    if evidence.status == EvidenceStatus::Complete {
1168        return Ok(vec![]);
1169    }
1170    evidence.transcript_path = Some(values.remove(0));
1171    evidence.resume_path = Some(values.remove(0));
1172    evidence.pane_path = Some(values.remove(0));
1173    evidence.report_path = Some(values.remove(0));
1174    evidence.transcript_sha256 = Some(values.remove(0));
1175    evidence.status = EvidenceStatus::Complete;
1176    evidence.error = None;
1177    node.updated_at = ev.ts;
1178    Ok(vec![ProjectionOp::Node(node)])
1179}
1180
1181fn reduce_worker_evidence_failed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1182    let events_path = paths.events();
1183    let node_id = require_envelope_node_id(&events_path, ev)?;
1184    let attempt = evidence_attempt(&events_path, ev)?;
1185    let session_id = want_str(&events_path, ev, &ev.data, "session_id")?.to_string();
1186    let detail = want_str(&events_path, ev, &ev.data, "error")?.to_string();
1187    let mut node = match read_node_opt(paths, &node_id)? {
1188        Some(node) => node,
1189        None => return Ok(vec![]),
1190    };
1191    let Some(evidence) = node.evidence.as_mut() else {
1192        return Err(Error::CorruptEventLog {
1193            path: events_path,
1194            reason: format!(
1195                "event seq={} evidence failure has no recorded Pi session",
1196                ev.seq
1197            ),
1198        });
1199    };
1200    if evidence.attempt != attempt || evidence.session_id != session_id {
1201        return Ok(vec![]);
1202    }
1203    if evidence.status == EvidenceStatus::Complete {
1204        return Ok(vec![]);
1205    }
1206    evidence.status = EvidenceStatus::Failed;
1207    evidence.error = Some(detail);
1208    node.updated_at = ev.ts;
1209    Ok(vec![ProjectionOp::Node(node)])
1210}
1211
1212/// Fold a `node.death_observed` event onto [`Node::first_death_at`], recording
1213/// the FIRST tick on which the supervisor saw this node's worker confirmed-dead
1214/// with no told `worker.exited` and no merge — the durable anchor for the
1215/// residual crash backstop's fixed post-death grace (design.md §2.1a, issue
1216/// `typed-supervisor-outcomes`).
1217///
1218/// The anchor is the event's own timestamp (`ev.ts`) — no payload field needed.
1219/// The fold is **first-write-wins**: the anchor is monotonic, so a later
1220/// re-observation (a supervisor restart still seeing the dead pid) never resets
1221/// the clock, and a full replay from seq 0 converges to the first observation. A
1222/// `node.death_observed` for a missing or already terminal node folds to nothing
1223/// (the backstop is moot once the node settles).
1224fn reduce_node_death_observed(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1225    let events_path = paths.events();
1226    let node_id = require_envelope_node_id(&events_path, ev)?;
1227    let mut n = match read_node_opt(paths, &node_id)? {
1228        Some(n) => n,
1229        None => return Ok(vec![]),
1230    };
1231    // First-write-wins (monotonic anchor). No-op once the backstop is moot or a
1232    // higher-fidelity fact exists — a terminal node, a told `worker.exited`, a
1233    // landed report, or an in-flight merge transaction. The supervisor's emitter
1234    // already gates on all of these under the exclusive lock; mirroring them here
1235    // keeps a from-scratch replay convergent regardless of caller.
1236    if n.first_death_at.is_some()
1237        || n.status.is_terminal()
1238        || n.worker_exit.is_some()
1239        || n.last_report.is_some()
1240        || n.pending_merge.is_some()
1241    {
1242        return Ok(vec![]);
1243    }
1244    n.first_death_at = Some(ev.ts);
1245    n.updated_at = ev.ts;
1246    Ok(vec![ProjectionOp::Node(n)])
1247}
1248
1249/// Fold an agent's explicit request for a human decision onto the node without
1250/// changing its status. This is deliberately non-terminal: the worker may still
1251/// resolve the fork itself or proceed with its stated default.
1252///
1253/// The first open signal wins until a matching `node.input_resolved` clears it,
1254/// so retries or duplicate writes cannot move the grace clock forward. The
1255/// event timestamp, not a caller-supplied payload timestamp, is the durable
1256/// restart-safe anchor.
1257fn reduce_node_awaiting_input(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1258    const MAX_ITEMS: usize = 8;
1259    const MAX_TOPIC_CHARS: usize = 512;
1260    const MAX_OPTIONS: usize = 16;
1261    const MAX_OPTION_CHARS: usize = 256;
1262
1263    let events_path = paths.events();
1264    let node_id = require_envelope_node_id(&events_path, ev)?;
1265    // Validate before every state-dependent no-op. An ignored duplicate or late
1266    // event must still be structurally valid before it enters the durable log.
1267    let items = ev
1268        .data
1269        .get("discussion_items")
1270        .and_then(Value::as_array)
1271        .filter(|items| !items.is_empty() && items.len() <= MAX_ITEMS)
1272        .ok_or_else(|| Error::CorruptEventLog {
1273            path: events_path.clone(),
1274            reason: format!(
1275                "event seq={} kind=node.awaiting_input requires 1..={MAX_ITEMS} `discussion_items`",
1276                ev.seq
1277            ),
1278        })?;
1279    for (index, item) in items.iter().enumerate() {
1280        let obj = item.as_object().ok_or_else(|| Error::CorruptEventLog {
1281            path: events_path.clone(),
1282            reason: format!(
1283                "event seq={} discussion_items[{index}] must be an object",
1284                ev.seq
1285            ),
1286        })?;
1287        let topic = obj.get("topic").and_then(Value::as_str).unwrap_or("");
1288        let default = obj
1289            .get("recommended_default")
1290            .and_then(Value::as_str)
1291            .unwrap_or("");
1292        let options = obj.get("options").and_then(Value::as_array);
1293        let options_valid = options.is_some_and(|values| {
1294            !values.is_empty()
1295                && values.len() <= MAX_OPTIONS
1296                && values.iter().all(|v| {
1297                    v.as_str().is_some_and(|s| {
1298                        !s.trim().is_empty() && s.chars().count() <= MAX_OPTION_CHARS
1299                    })
1300                })
1301                && values.iter().any(|v| v.as_str() == Some(default))
1302        });
1303        if topic.trim().is_empty()
1304            || topic.chars().count() > MAX_TOPIC_CHARS
1305            || default.trim().is_empty()
1306            || !options_valid
1307        {
1308            return Err(Error::CorruptEventLog {
1309                path: events_path.clone(),
1310                reason: format!(
1311                    "event seq={} discussion_items[{index}] requires bounded non-empty `topic`, 1..={MAX_OPTIONS} bounded string `options`, and a `recommended_default` present in options",
1312                    ev.seq
1313                ),
1314            });
1315        }
1316    }
1317
1318    let mut n = match read_node_opt(paths, &node_id)? {
1319        Some(n) => n,
1320        None => return Ok(vec![]),
1321    };
1322    if n.status.is_terminal() || n.worker_exit.is_some() {
1323        return Ok(vec![]);
1324    }
1325    if let Some(open) = n.awaiting_input.as_mut() {
1326        // A later fork joins the current open generation without moving its
1327        // restart-safe clock or notification key. Bound the aggregate too.
1328        if open.discussion_items.len() + items.len() > MAX_ITEMS {
1329            return Err(Error::CorruptEventLog {
1330                path: events_path,
1331                reason: format!(
1332                    "event seq={} would exceed {MAX_ITEMS} open discussion items",
1333                    ev.seq
1334                ),
1335            });
1336        }
1337        open.discussion_items.extend(items.iter().cloned());
1338    } else {
1339        n.awaiting_input = Some(Box::new(crate::schema::AwaitingInput {
1340            opened_at: ev.ts,
1341            event_seq: ev.seq,
1342            discussion_items: items.clone(),
1343        }));
1344    }
1345    n.updated_at = ev.ts;
1346    Ok(vec![ProjectionOp::Node(n)])
1347}
1348
1349/// Clear the current open decision request. `event_seq` is mandatory and
1350/// fences the resolve to the generation the worker observed, so a delayed
1351/// timeout cannot clear a newer question opened after the old one resolved.
1352fn reduce_node_input_resolved(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1353    let events_path = paths.events();
1354    let node_id = require_envelope_node_id(&events_path, ev)?;
1355    // Validate before the no-open no-op so malformed events never enter the log
1356    // merely because their projection happens to be absent today.
1357    let seq = ev
1358        .data
1359        .get("event_seq")
1360        .and_then(Value::as_u64)
1361        .ok_or_else(|| Error::CorruptEventLog {
1362            path: events_path.clone(),
1363            reason: format!(
1364                "event seq={} kind=node.input_resolved requires unsigned `event_seq`",
1365                ev.seq
1366            ),
1367        })?;
1368    let mut n = match read_node_opt(paths, &node_id)? {
1369        Some(n) => n,
1370        None => return Ok(vec![]),
1371    };
1372    let Some(open) = n.awaiting_input.as_ref() else {
1373        return Ok(vec![]);
1374    };
1375    if seq != open.event_seq {
1376        return Ok(vec![]);
1377    }
1378    n.awaiting_input = None;
1379    n.updated_at = ev.ts;
1380    Ok(vec![ProjectionOp::Node(n)])
1381}
1382
1383/// The event kind `run merge` appends BEFORE mutating git to record the
1384/// in-flight merge transaction (design.md §2.1b / A2). Its `data` payload is a
1385/// serialized [`MergeTxn`]; the reducer folds it onto [`Node::pending_merge`].
1386pub const KIND_MERGE_STARTED: &str = "merge.started";
1387
1388/// The event kind recovery appends when it resolves a pending merge transaction
1389/// by REJECTING it — the recorded source ref never moved (the git mutation never
1390/// landed), so the worker's branch + work are preserved and the transaction is
1391/// cleared. Its `data` carries `op_id` (which transaction) and `reason`.
1392pub const KIND_MERGE_ABORTED: &str = "merge.aborted";
1393
1394/// Fold a `merge.started` event onto [`Node::pending_merge`], recording the
1395/// in-flight `run merge` transaction BEFORE the git mutation so a crash between
1396/// the git merge and the terminal `explicit-merge` report can be resolved
1397/// deterministically by OID (design.md §2.1b / A2, issue
1398/// `merge-transaction-recovery`). The reducer deliberately does NOT transition
1399/// `status`: recording a transaction is not a terminal outcome.
1400///
1401/// Payload contract: the `data` is a serialized [`MergeTxn`]; a payload missing
1402/// required fields is a corrupt event (the reducer is the canonical gate, so the
1403/// append is rejected before any byte is written). The fold is idempotent —
1404/// re-folding the same `op_id` on replay is a clean no-op — and last-write-wins
1405/// across a fresh attempt's larger `op_id` (each `run merge` re-reads
1406/// `expected_source_oid`, so the newest record is authoritative).
1407fn reduce_merge_started(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1408    let events_path = paths.events();
1409    let node_id = require_envelope_node_id(&events_path, ev)?;
1410    let txn: MergeTxn =
1411        serde_json::from_value(ev.data.clone()).map_err(|e| Error::CorruptEventLog {
1412            path: events_path.clone(),
1413            reason: format!(
1414                "event seq={} kind=merge.started has an invalid MergeTxn payload: {e}",
1415                ev.seq
1416            ),
1417        })?;
1418    let mut n = match read_node_opt(paths, &node_id)? {
1419        Some(n) => n,
1420        None => return Ok(vec![]),
1421    };
1422    // A terminal node has no in-flight merge to track — `run merge` is refused on
1423    // a terminal run at the CLI, so this is a dead/duplicate event. Ignore it
1424    // (never resurrect the projection).
1425    if n.status.is_terminal() {
1426        return Ok(vec![]);
1427    }
1428    // Idempotent: re-folding the SAME transaction on replay must not churn
1429    // `updated_at`.
1430    if n.pending_merge.as_ref().map(|t| t.op_id.as_str()) == Some(txn.op_id.as_str()) {
1431        return Ok(vec![]);
1432    }
1433    n.pending_merge = Some(Box::new(txn));
1434    n.updated_at = ev.ts;
1435    Ok(vec![ProjectionOp::Node(n)])
1436}
1437
1438/// Fold a `merge.aborted` event, clearing [`Node::pending_merge`] iff it names
1439/// the transaction being aborted (`op_id` match). Recovery appends this when it
1440/// determines a pending merge's git mutation never landed (the source ref is
1441/// still at `expected_source_oid`) or moved unexpectedly — the transaction is
1442/// rejected, the worker's branch + work are preserved, and the node stays
1443/// whatever non-terminal status it was (a retry may re-attempt the merge).
1444///
1445/// The `op_id` guard is what keeps this from clobbering a *newer* transaction: a
1446/// stale `merge.aborted` for a superseded attempt (a different `op_id`) is a
1447/// clean no-op. Deliberately
1448/// does NOT transition `status`.
1449fn reduce_merge_aborted(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1450    let events_path = paths.events();
1451    let node_id = require_envelope_node_id(&events_path, ev)?;
1452    let op_id = ev
1453        .data
1454        .get("op_id")
1455        .and_then(Value::as_str)
1456        .ok_or_else(|| Error::CorruptEventLog {
1457            path: events_path.clone(),
1458            reason: format!(
1459                "event seq={} kind=merge.aborted is missing string `op_id`",
1460                ev.seq
1461            ),
1462        })?;
1463    let mut n = match read_node_opt(paths, &node_id)? {
1464        Some(n) => n,
1465        None => return Ok(vec![]),
1466    };
1467    // Clear only the transaction this event names. A mismatch (already resolved,
1468    // or a newer attempt is pending) is a clean no-op.
1469    match n.pending_merge.as_ref() {
1470        Some(t) if t.op_id == op_id => {}
1471        _ => return Ok(vec![]),
1472    }
1473    n.pending_merge = None;
1474    n.updated_at = ev.ts;
1475    Ok(vec![ProjectionOp::Node(n)])
1476}
1477
1478/// Emit an observability trace for a status event dropped by the terminal
1479/// guard. Re-applying the *same* terminal status is routine idempotent replay
1480/// (`debug`); an event carrying a *different* status is a real conflict that
1481/// should not occur on a well-formed log (`warn`) — e.g. a `done` node being
1482/// told to go `cancelled`. The guard no-ops either way; the level is the only
1483/// difference, so a genuine corruption signal is visible without flooding
1484/// logs on every replay.
1485fn trace_terminal_noop(ev: &Event, current: Status, incoming: Status) {
1486    if current == incoming {
1487        tracing::debug!(
1488            target: "taskfleet_core::reducer",
1489            seq = ev.seq, kind = %ev.kind, status = ?current,
1490            "no-op: status re-applied to terminal target"
1491        );
1492    } else {
1493        tracing::warn!(
1494            target: "taskfleet_core::reducer",
1495            seq = ev.seq, kind = %ev.kind, current = ?current, incoming = ?incoming,
1496            "no-op: ignored conflicting transition from terminal target"
1497        );
1498    }
1499}
1500
1501/// Derive the terminal status a `node.report` event asserts, enforcing the
1502/// success-XOR-cancelled invariant with strict boolean typing.
1503///
1504/// `cancelled: true` (with `success: false` or absent) → [`Status::Cancelled`].
1505/// Otherwise `success` must be present: `true` → [`Status::Done`], `false` →
1506/// [`Status::Failed`]. Neither field (bare `{}`), the contradiction
1507/// `success: true` + `cancelled: true`, or a non-boolean `success` /
1508/// `cancelled` is a [`Error::CorruptEventLog`].
1509fn report_terminal_status(events_path: &Path, ev: &Event) -> Result<Status> {
1510    let corrupt = |reason: String| Error::CorruptEventLog {
1511        path: events_path.to_path_buf(),
1512        reason,
1513    };
1514    let cancelled = optional_bool(events_path, ev, &ev.data, "cancelled")?.unwrap_or(false);
1515    let success = optional_bool(events_path, ev, &ev.data, "success")?;
1516    if cancelled {
1517        if success == Some(true) {
1518            return Err(corrupt(format!(
1519                "event seq={} kind=node.report has contradictory `success: true` with `cancelled: true`",
1520                ev.seq
1521            )));
1522        }
1523        Ok(Status::Cancelled)
1524    } else {
1525        match success {
1526            Some(true) => Ok(Status::Done),
1527            Some(false) => Ok(Status::Failed),
1528            None => Err(corrupt(format!(
1529                "event seq={} kind=node.report must set boolean `success` or `cancelled: true`",
1530                ev.seq
1531            ))),
1532        }
1533    }
1534}
1535
1536fn reduce_child_spawned(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1537    // `child.spawned` is written to the PARENT run's events; the parent
1538    // spawning node is `ev.node_id`, the child run/node lives in `data`.
1539    let events_path = paths.events();
1540    let parent_node_id = ev.node_id.clone().ok_or_else(|| Error::CorruptEventLog {
1541        path: events_path.clone(),
1542        reason: format!(
1543            "event seq={} kind=child.spawned missing parent `node_id`",
1544            ev.seq
1545        ),
1546    })?;
1547    let child_run_id = RunId::parse_str(want_str(&events_path, ev, &ev.data, "child_run_id")?)
1548        .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1549    let child_node_id = NodeId::parse_str(
1550        ev.data
1551            .get("child_node_id")
1552            .and_then(Value::as_str)
1553            .unwrap_or("n-0001"),
1554    )
1555    .map_err(|e| corrupt_id(&events_path, ev, &e))?;
1556    let mut n = match read_node_opt(paths, &parent_node_id)? {
1557        Some(n) => n,
1558        None => return Ok(vec![]),
1559    };
1560    let new_ref = ChildRef {
1561        run_id: child_run_id,
1562        node_id: child_node_id,
1563    };
1564    if n.children.iter().any(|c| c == &new_ref) {
1565        // Already recorded — pure no-op so replayed events don't churn
1566        // `updated_at` or the projection file.
1567        return Ok(vec![]);
1568    }
1569    n.children.push(new_ref);
1570    n.updated_at = ev.ts;
1571    Ok(vec![ProjectionOp::Node(n)])
1572}
1573
1574/// `supervisor.attached` records the supervisor PID watching the envelope
1575/// node onto `Node.supervisor_pid`. Event-sourced replacement for the
1576/// supervisor's former direct `write_node` (issue
1577/// `supervisor-state-not-event-sourced`), so a from-scratch projection
1578/// rebuild reproduces the field.
1579///
1580/// Latest-wins: a later attach (a supervisor restart binds a fresh PID)
1581/// overrides the recorded value. Re-applying an event that carries the
1582/// already-recorded PID is a pure no-op, so replay never churns the
1583/// projection file's `updated_at`.
1584fn reduce_supervisor_attached(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1585    let events_path = paths.events();
1586    let node_id = require_envelope_node_id(&events_path, ev)?;
1587    let raw = ev
1588        .data
1589        .get("pid")
1590        .and_then(Value::as_i64)
1591        .ok_or_else(|| Error::CorruptEventLog {
1592            path: events_path.clone(),
1593            reason: format!(
1594                "event seq={} kind=supervisor.attached missing/invalid `pid`",
1595                ev.seq
1596            ),
1597        })?;
1598    let pid = i32::try_from(raw).map_err(|_| Error::CorruptEventLog {
1599        path: events_path.clone(),
1600        reason: format!(
1601            "event seq={} kind=supervisor.attached `pid` out of i32 range: {raw}",
1602            ev.seq
1603        ),
1604    })?;
1605    let mut n = match read_node_opt(paths, &node_id)? {
1606        Some(n) => n,
1607        None => return Ok(vec![]),
1608    };
1609    if n.supervisor_pid == Some(pid) {
1610        return Ok(vec![]);
1611    }
1612    n.supervisor_pid = Some(pid);
1613    n.updated_at = ev.ts;
1614    Ok(vec![ProjectionOp::Node(n)])
1615}
1616
1617/// `supervisor.cursor_advanced` mirrors the supervisor's per-child report
1618/// cursor onto the envelope (parent) node's `last_processed_report_seq_by_child`
1619/// map. Event-sourced replacement for the supervisor's former direct
1620/// `write_node` of that map (issue `supervisor-state-not-event-sourced`).
1621///
1622/// The cursor is monotonic: a `report_seq` at or below the recorded
1623/// high-water mark for this child is a no-op, so replaying the same event —
1624/// or an older out-of-order one — never moves the cursor backward or churns
1625/// the projection. This is the §7.3 idempotency guarantee at the reducer
1626/// boundary.
1627fn reduce_supervisor_cursor_advanced(paths: &RunPaths, ev: &Event) -> Result<Vec<ProjectionOp>> {
1628    let events_path = paths.events();
1629    let node_id = require_envelope_node_id(&events_path, ev)?;
1630    let child_run_id = want_str(&events_path, ev, &ev.data, "child_run_id")?;
1631    // Validate the child id even though it only becomes a map key — a forged
1632    // event must not smuggle a path-shaped or malformed run id into the
1633    // projection.
1634    RunId::parse_str(child_run_id).map_err(|e| corrupt_id(&events_path, ev, &e))?;
1635    let report_seq = ev
1636        .data
1637        .get("report_seq")
1638        .and_then(Value::as_u64)
1639        .ok_or_else(|| Error::CorruptEventLog {
1640            path: events_path.clone(),
1641            reason: format!(
1642                "event seq={} kind=supervisor.cursor_advanced missing/invalid `report_seq`",
1643                ev.seq
1644            ),
1645        })?;
1646    let mut n = match read_node_opt(paths, &node_id)? {
1647        Some(n) => n,
1648        None => return Ok(vec![]),
1649    };
1650    if let Some(prev) = n
1651        .last_processed_report_seq_by_child
1652        .get(child_run_id)
1653        .and_then(Value::as_u64)
1654    {
1655        if report_seq <= prev {
1656            return Ok(vec![]);
1657        }
1658    }
1659    n.last_processed_report_seq_by_child
1660        .insert(child_run_id.to_string(), Value::from(report_seq));
1661    n.updated_at = ev.ts;
1662    Ok(vec![ProjectionOp::Node(n)])
1663}
1664
1665#[cfg(test)]
1666mod tests {
1667    use super::*;
1668    use crate::schema::Event;
1669    use chrono::Utc;
1670    use tempfile::TempDir;
1671
1672    fn event(run_id: &str) -> Event {
1673        Event {
1674            ts: Utc::now(),
1675            seq: 1,
1676            kind: "run.status".into(),
1677            run_id: RunId::parse_str(run_id).unwrap(),
1678            node_id: None,
1679            idempotency_key: None,
1680            data: serde_json::json!({ "status": "running" }),
1681        }
1682    }
1683
1684    #[test]
1685    fn orchestrator_decision_and_discuss_critical_reduce_to_noop() {
1686        // The /orchestrate audit kinds are append-only: the reducer must plan
1687        // ZERO projection ops for them regardless of payload, so the event log
1688        // is their sole home and no projection is created or mutated.
1689        let tmp = TempDir::new().unwrap();
1690        let run_id = "01jxsnap000000000000000000";
1691        let rid = RunId::parse_str(run_id).unwrap();
1692        let dir = crate::run_dir(tmp.path(), &rid);
1693        std::fs::create_dir_all(&dir).unwrap();
1694        let paths = RunPaths::new(dir, run_id).unwrap();
1695
1696        // Bootstrap a manifest so we can prove the audit events leave it
1697        // byte-for-byte untouched (no counter churn, no status drift).
1698        let mut created = event(run_id);
1699        created.kind = "run.created".into();
1700        created.data = serde_json::json!({
1701            "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
1702        });
1703        apply_event(&paths, &created).expect("run.created applies");
1704        let manifest_before = std::fs::read(paths.manifest()).unwrap();
1705
1706        for (seq, kind) in [(10u64, "orchestrator.decision"), (11, "discuss.critical")] {
1707            let mut ev = event(run_id);
1708            ev.seq = seq;
1709            ev.kind = kind.into();
1710            // A non-trivial payload to prove the reducer ignores it wholesale.
1711            ev.data = serde_json::json!({ "summary": "x", "arbitrary": [1, 2, 3] });
1712            let ops = reduce_event_to_ops(&paths, &ev).expect("audit kind reduces cleanly");
1713            assert!(ops.is_empty(), "{kind} must plan no projection ops");
1714            // apply_event is the plan+commit path; it must also be a clean no-op.
1715            apply_event(&paths, &ev).expect("audit kind applies as no-op");
1716        }
1717
1718        // The manifest is unchanged and no stray projection dirs appeared.
1719        assert_eq!(
1720            std::fs::read(paths.manifest()).unwrap(),
1721            manifest_before,
1722            "audit events must not mutate the manifest"
1723        );
1724        assert!(!paths.nodes_dir().exists(), "no node projection created");
1725    }
1726
1727    #[test]
1728    fn run_created_folds_harness_when_present_and_defaults_none() {
1729        let tmp = TempDir::new().unwrap();
1730
1731        // A `run.created` carrying `harness` folds it onto the manifest.
1732        let run_id = "01jxhrnsaa0000000000000001";
1733        let rid = RunId::parse_str(run_id).unwrap();
1734        let dir = crate::run_dir(tmp.path(), &rid);
1735        std::fs::create_dir_all(&dir).unwrap();
1736        let paths = RunPaths::new(dir, run_id).unwrap();
1737        let mut created = event(run_id);
1738        created.kind = "run.created".into();
1739        created.data = serde_json::json!({
1740            "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
1741            "harness": "pi", "harness_source": "flag",
1742        });
1743        apply_event(&paths, &created).expect("run.created applies");
1744        let m = read_manifest_opt(&paths).unwrap().unwrap();
1745        assert_eq!(m.harness.as_deref(), Some("pi"));
1746
1747        // A `run.created` WITHOUT `harness` (legacy / claude) leaves it `None`.
1748        let run_id2 = "01jxhrnsaa0000000000000002";
1749        let rid2 = RunId::parse_str(run_id2).unwrap();
1750        let dir2 = crate::run_dir(tmp.path(), &rid2);
1751        std::fs::create_dir_all(&dir2).unwrap();
1752        let paths2 = RunPaths::new(dir2, run_id2).unwrap();
1753        let mut created2 = event(run_id2);
1754        created2.kind = "run.created".into();
1755        created2.data =
1756            serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1757        apply_event(&paths2, &created2).expect("run.created applies");
1758        let m2 = read_manifest_opt(&paths2).unwrap().unwrap();
1759        assert_eq!(m2.harness, None);
1760    }
1761
1762    /// Bootstrap a run manifest + one live `n-0001` spinoff node, returning its
1763    /// paths. Used by the `node.retry` reducer tests.
1764    fn bootstrap_retry_node(tmp: &TempDir, run_id: &str) -> RunPaths {
1765        let rid = RunId::parse_str(run_id).unwrap();
1766        let dir = crate::run_dir(tmp.path(), &rid);
1767        std::fs::create_dir_all(&dir).unwrap();
1768        let paths = RunPaths::new(dir, run_id).unwrap();
1769        let mut created = event(run_id);
1770        created.kind = "run.created".into();
1771        created.data =
1772            serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
1773        apply_event(&paths, &created).expect("run.created applies");
1774        let mut node = event(run_id);
1775        node.seq = 2;
1776        node.kind = "node.created".into();
1777        node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1778        node.data = serde_json::json!({
1779            "kind": "spinoff",
1780            "branch": "wt/foo",
1781            "worktree_path": "/tmp/old-wt",
1782            "agent_pid": 111,
1783        });
1784        apply_event(&paths, &node).expect("node.created applies");
1785        paths
1786    }
1787
1788    /// `node.retry` rewires the node to the freshly re-spawned agent, returns it to
1789    /// `Pending`, re-stamps `started_at`, and increments the durable
1790    /// `retry_attempts` bound (issue `autoretry-agent-died-worker`).
1791    #[test]
1792    fn node_retry_rewires_node_and_increments_attempts() {
1793        let tmp = TempDir::new().unwrap();
1794        let run_id = "01jxsnap000000000000000000";
1795        let paths = bootstrap_retry_node(&tmp, run_id);
1796
1797        let mut retry = event(run_id);
1798        retry.seq = 3;
1799        retry.kind = "node.retry".into();
1800        retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1801        retry.data = serde_json::json!({
1802            "attempt": 1,
1803            "reason": "agent-died",
1804            "branch": "wt/foo-r1",
1805            "base_sha": "a".repeat(40),
1806            "worktree_path": "/tmp/new-wt",
1807            "agent_pid": 222,
1808            "tmux_session": "s",
1809            "tmux_window_id": "@9",
1810        });
1811        apply_event(&paths, &retry).expect("node.retry applies");
1812
1813        let n = read_n0001(&paths);
1814        assert_eq!(n.retry_attempts, 1, "attempt bound incremented");
1815        assert_eq!(
1816            n.branch.as_deref(),
1817            Some("wt/foo-r1"),
1818            "rewired to new branch"
1819        );
1820        assert_eq!(n.worktree_path.as_deref(), Some("/tmp/new-wt"));
1821        assert_eq!(n.agent_pid, Some(222), "rewired to new agent pid");
1822        assert_eq!(n.status, Status::Pending, "node returns to pending");
1823        assert!(n.last_report.is_none());
1824        assert_eq!(
1825            n.tmux_identity.as_ref().map(|t| t.window_id.as_str()),
1826            Some("@9"),
1827            "rewired tmux identity"
1828        );
1829
1830        // A second retry increments again — the bound is monotone.
1831        let mut retry2 = event(run_id);
1832        retry2.seq = 4;
1833        retry2.kind = "node.retry".into();
1834        retry2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1835        retry2.data = serde_json::json!({
1836            "attempt": 2, "reason": "agent-died", "branch": "wt/foo-r2",
1837            "worktree_path": "/tmp/new-wt-2", "agent_pid": 333,
1838        });
1839        apply_event(&paths, &retry2).expect("node.retry applies");
1840        assert_eq!(read_n0001(&paths).retry_attempts, 2);
1841    }
1842
1843    #[test]
1844    fn stale_worker_evidence_and_exit_do_not_cross_retry_generation() {
1845        let tmp = TempDir::new().unwrap();
1846        let run_id = "01jxsnap000000000000000000";
1847        let paths = bootstrap_retry_node(&tmp, run_id);
1848        let nid = NodeId::parse_str("n-0001").unwrap();
1849        let current_session = "018f5f64-b137-7d44-b2b4-4f02c3f646e8";
1850
1851        let mut retry = event(run_id);
1852        retry.seq = 3;
1853        retry.kind = "node.retry".into();
1854        retry.node_id = Some(nid.clone());
1855        retry.data = serde_json::json!({
1856            "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1857            "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1858            "pi_session_id": current_session,
1859            "pi_session_path": format!(".creating/pi-sessions/{run_id}/pi-session-{current_session}.jsonl"),
1860            "pi_session_cwd": "/tmp/new-wt"
1861        });
1862        apply_event(&paths, &retry).unwrap();
1863
1864        let mut stale_failure = event(run_id);
1865        stale_failure.seq = 4;
1866        stale_failure.kind = "worker.evidence.failed".into();
1867        stale_failure.node_id = Some(nid.clone());
1868        stale_failure.data = serde_json::json!({
1869            "attempt": 0,
1870            "session_id": "118f5f64-b137-7d44-b2b4-4f02c3f646e8",
1871            "error": "old attempt"
1872        });
1873        apply_event(&paths, &stale_failure).unwrap();
1874        assert_eq!(
1875            read_n0001(&paths).evidence.unwrap().status,
1876            EvidenceStatus::Pending
1877        );
1878
1879        let mut stale_exit = event(run_id);
1880        stale_exit.seq = 5;
1881        stale_exit.kind = "worker.exited".into();
1882        stale_exit.node_id = Some(nid.clone());
1883        stale_exit.data = serde_json::json!({"attempt":0,"exit_code":9});
1884        apply_event(&paths, &stale_exit).unwrap();
1885        assert!(read_n0001(&paths).worker_exit.is_none());
1886
1887        let mut archived = event(run_id);
1888        archived.seq = 6;
1889        archived.kind = "worker.evidence.archived".into();
1890        archived.node_id = Some(nid);
1891        archived.data = serde_json::json!({
1892            "attempt":1, "session_id":current_session,
1893            "transcript_path":"evidence/n-0001/pi-session.original.jsonl",
1894            "resume_path":"evidence/n-0001/pi-session.resume.jsonl",
1895            "pane_path":"evidence/n-0001/final-pane.log",
1896            "report_path":"evidence/n-0001/terminal-report.json",
1897            "transcript_sha256":"0".repeat(64)
1898        });
1899        apply_event(&paths, &archived).unwrap();
1900        assert_eq!(
1901            read_n0001(&paths).evidence.unwrap().status,
1902            EvidenceStatus::Complete
1903        );
1904    }
1905
1906    /// A `node.retry` against an already-terminal node is a dead event: the
1907    /// terminal-state invariant holds, so a late retry never resurrects a settled
1908    /// node (a real report that raced in wins).
1909    #[test]
1910    fn node_retry_against_terminal_node_is_noop() {
1911        let tmp = TempDir::new().unwrap();
1912        let run_id = "01jxsnap000000000000000000";
1913        let paths = bootstrap_retry_node(&tmp, run_id);
1914
1915        // Terminalize the node via a success report.
1916        let mut report = event(run_id);
1917        report.seq = 3;
1918        report.kind = "node.report".into();
1919        report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1920        report.data = serde_json::json!({ "success": true });
1921        apply_event(&paths, &report).expect("node.report applies");
1922        assert_eq!(read_n0001(&paths).status, Status::Done);
1923
1924        let mut retry = event(run_id);
1925        retry.seq = 4;
1926        retry.kind = "node.retry".into();
1927        retry.node_id = Some(NodeId::parse_str("n-0001").unwrap());
1928        retry.data = serde_json::json!({
1929            "attempt": 1, "reason": "agent-died", "branch": "wt/foo-r1",
1930            "worktree_path": "/tmp/new-wt", "agent_pid": 222,
1931        });
1932        apply_event(&paths, &retry).expect("node.retry applies as no-op");
1933
1934        let n = read_n0001(&paths);
1935        assert_eq!(n.status, Status::Done, "terminal node not resurrected");
1936        assert_eq!(n.retry_attempts, 0, "no increment against terminal node");
1937        assert_eq!(n.agent_pid, Some(111), "not rewired");
1938    }
1939
1940    #[test]
1941    fn apply_event_rejects_event_from_a_different_run() {
1942        let tmp = TempDir::new().unwrap();
1943        let run_id = "01jxsnap000000000000000000";
1944        let rid = RunId::parse_str(run_id).unwrap();
1945        let dir = crate::run_dir(tmp.path(), &rid);
1946        std::fs::create_dir_all(&dir).unwrap();
1947        let paths = RunPaths::new(dir, run_id).unwrap();
1948
1949        // An event whose envelope names a different run must not be folded.
1950        let foreign = event("02jxsnap000000000000000000");
1951        let err = apply_event(&paths, &foreign).expect_err("cross-run event must be rejected");
1952        assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
1953
1954        // The matching run_id is accepted (no projection exists yet, so
1955        // `run.status` is a clean no-op rather than an error).
1956        let mine = event(run_id);
1957        apply_event(&paths, &mine).expect("matching run_id must be accepted");
1958    }
1959
1960    #[test]
1961    fn tmux_identity_from_data_reads_qualified_fields() {
1962        let d = serde_json::json!({
1963            "tmux_socket": "/private/tmp/tmux-501/default",
1964            "tmux_session": "taskfleet",
1965            "tmux_window_id": "@42",
1966        });
1967        let id = tmux_identity_from_data(&d).expect("qualified identity");
1968        assert_eq!(id.socket.as_deref(), Some("/private/tmp/tmux-501/default"));
1969        assert_eq!(id.session, "taskfleet");
1970        assert_eq!(id.window_id, "@42");
1971        // No pane_id in this event → None (back-compat / older create.sh).
1972        assert_eq!(id.pane_id, None);
1973
1974        // Null socket is tolerated — session + window_id are the minimum.
1975        let d2 = serde_json::json!({
1976            "tmux_socket": null,
1977            "tmux_session": "taskfleet",
1978            "tmux_window_id": "@7",
1979        });
1980        let id2 = tmux_identity_from_data(&d2).expect("identity without socket");
1981        assert_eq!(id2.socket, None);
1982        assert_eq!(id2.window_id, "@7");
1983
1984        // A create.sh that emits `tmux_pane_id` is folded into the identity.
1985        let d3 = serde_json::json!({
1986            "tmux_session": "taskfleet",
1987            "tmux_window_id": "@42",
1988            "tmux_pane_id": "%7",
1989        });
1990        let id3 = tmux_identity_from_data(&d3).expect("identity with pane");
1991        assert_eq!(id3.pane_id.as_deref(), Some("%7"));
1992        assert_eq!(id3.capture_target(), "%7");
1993
1994        // Explicit `tmux_pane_id: null` (create.sh emits null when its pane
1995        // query failed) must fold to None — never `Some("null")`.
1996        let d4 = serde_json::json!({
1997            "tmux_session": "taskfleet",
1998            "tmux_window_id": "@42",
1999            "tmux_pane_id": null,
2000        });
2001        let id4 = tmux_identity_from_data(&d4).expect("identity with null pane");
2002        assert_eq!(id4.pane_id, None);
2003        assert_eq!(id4.capture_target(), "@42");
2004    }
2005
2006    #[test]
2007    fn tmux_identity_from_data_back_compat_is_none() {
2008        // Legacy create.sh: no qualified fields at all.
2009        let legacy = serde_json::json!({ "tmux_window": "🚀 wt/x" });
2010        assert!(tmux_identity_from_data(&legacy).is_none());
2011        // Partial (window_id without session) is also insufficient → None.
2012        let partial = serde_json::json!({ "tmux_window_id": "@42" });
2013        assert!(tmux_identity_from_data(&partial).is_none());
2014    }
2015
2016    /// End-to-end: a `node.created` event carrying the qualified fields folds
2017    /// them into `Node.tmux_identity`; one without them leaves it `None`.
2018    #[test]
2019    fn node_created_populates_tmux_identity() {
2020        let tmp = TempDir::new().unwrap();
2021        let run_id = "01jxsnap000000000000000000";
2022        let rid = RunId::parse_str(run_id).unwrap();
2023        let dir = crate::run_dir(tmp.path(), &rid);
2024        std::fs::create_dir_all(&dir).unwrap();
2025        let paths = RunPaths::new(dir, run_id).unwrap();
2026
2027        let mut ev = event(run_id);
2028        ev.seq = 2;
2029        ev.kind = "node.created".into();
2030        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2031        ev.data = serde_json::json!({
2032            "kind": "spinoff",
2033            "tmux_window": "🚀 wt/x",
2034            "tmux_socket": "/private/tmp/tmux-501/default",
2035            "tmux_session": "taskfleet",
2036            "tmux_window_id": "@42",
2037        });
2038        apply_event(&paths, &ev).expect("node.created applies");
2039        let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2040            .unwrap()
2041            .unwrap();
2042        let id = n.tmux_identity.expect("qualified identity recorded");
2043        assert_eq!(id.session, "taskfleet");
2044        assert_eq!(id.window_id, "@42");
2045        assert_eq!(n.tmux_window.as_deref(), Some("🚀 wt/x"));
2046
2047        // A second run with a legacy event leaves tmux_identity None.
2048        let run2 = "02jxsnap000000000000000000";
2049        let rid2 = RunId::parse_str(run2).unwrap();
2050        let dir2 = crate::run_dir(tmp.path(), &rid2);
2051        std::fs::create_dir_all(&dir2).unwrap();
2052        let paths2 = RunPaths::new(dir2, run2).unwrap();
2053        let mut ev2 = event(run2);
2054        ev2.seq = 2;
2055        ev2.kind = "node.created".into();
2056        ev2.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2057        ev2.data = serde_json::json!({ "kind": "spinoff", "tmux_window": "🚀 wt/y" });
2058        apply_event(&paths2, &ev2).expect("legacy node.created applies");
2059        let n2 = read_node_opt(&paths2, &NodeId::parse_str("n-0001").unwrap())
2060            .unwrap()
2061            .unwrap();
2062        assert!(n2.tmux_identity.is_none());
2063        assert_eq!(n2.tmux_window.as_deref(), Some("🚀 wt/y"));
2064    }
2065
2066    #[test]
2067    fn node_materialization_populates_missing_manifest_source_branch() {
2068        let tmp = TempDir::new().unwrap();
2069        let run_id = "01jxsnap000000000000000001";
2070        let rid = RunId::parse_str(run_id).unwrap();
2071        let dir = crate::run_dir(tmp.path(), &rid);
2072        std::fs::create_dir_all(&dir).unwrap();
2073        let paths = RunPaths::new(dir, run_id).unwrap();
2074
2075        let mut created = event(run_id);
2076        created.kind = "run.created".into();
2077        created.data = serde_json::json!({
2078            "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2079        });
2080        apply_event(&paths, &created).unwrap();
2081        assert!(read_manifest_opt(&paths)
2082            .unwrap()
2083            .unwrap()
2084            .source_branch
2085            .is_none());
2086
2087        let mut node = event(run_id);
2088        node.seq = 2;
2089        node.kind = "node.created".into();
2090        node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2091        node.data = serde_json::json!({
2092            "kind": "spinoff",
2093            "source_branch": "main",
2094            "worktree_path": "/tmp/wt/pending"
2095        });
2096        apply_event(&paths, &node).unwrap();
2097
2098        let manifest = read_manifest_opt(&paths).unwrap().unwrap();
2099        assert_eq!(manifest.status, Status::Pending);
2100        assert_eq!(manifest.source_branch.as_deref(), Some("main"));
2101        let node = read_n0001(&paths);
2102        assert_eq!(node.worktree_path.as_deref(), Some("/tmp/wt/pending"));
2103    }
2104
2105    #[test]
2106    fn node_materialization_preserves_explicit_manifest_source_branch() {
2107        let tmp = TempDir::new().unwrap();
2108        let run_id = "01jxsnap000000000000000002";
2109        let rid = RunId::parse_str(run_id).unwrap();
2110        let dir = crate::run_dir(tmp.path(), &rid);
2111        std::fs::create_dir_all(&dir).unwrap();
2112        let paths = RunPaths::new(dir, run_id).unwrap();
2113
2114        let mut created = event(run_id);
2115        created.kind = "run.created".into();
2116        created.data = serde_json::json!({
2117            "kind": "spinoff", "lifecycle": "autonomous", "title": "t",
2118            "source_branch": "release"
2119        });
2120        apply_event(&paths, &created).unwrap();
2121
2122        let mut node = event(run_id);
2123        node.seq = 2;
2124        node.kind = "node.created".into();
2125        node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2126        node.data = serde_json::json!({
2127            "kind": "spinoff", "source_branch": "main"
2128        });
2129        apply_event(&paths, &node).unwrap();
2130
2131        assert_eq!(
2132            read_manifest_opt(&paths)
2133                .unwrap()
2134                .unwrap()
2135                .source_branch
2136                .as_deref(),
2137            Some("release")
2138        );
2139    }
2140
2141    /// Bootstrap a run with a single `n-0001` node via the event-sourced path,
2142    /// returning its paths. Shared by the supervisor-state replay tests below.
2143    fn seed_run_with_node(tmp: &TempDir, run_id: &str) -> RunPaths {
2144        let rid = RunId::parse_str(run_id).unwrap();
2145        let dir = crate::run_dir(tmp.path(), &rid);
2146        std::fs::create_dir_all(&dir).unwrap();
2147        let paths = RunPaths::new(dir, run_id).unwrap();
2148
2149        let mut created = event(run_id);
2150        created.kind = "run.created".into();
2151        created.data = serde_json::json!({
2152            "kind": "spinoff", "lifecycle": "autonomous", "title": "t"
2153        });
2154        apply_event(&paths, &created).expect("run.created applies");
2155
2156        let mut node = event(run_id);
2157        node.seq = 2;
2158        node.kind = "node.created".into();
2159        node.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2160        node.data = serde_json::json!({ "kind": "spinoff" });
2161        apply_event(&paths, &node).expect("node.created applies");
2162        paths
2163    }
2164
2165    fn read_n0001(paths: &RunPaths) -> Node {
2166        read_node_opt(paths, &NodeId::parse_str("n-0001").unwrap())
2167            .unwrap()
2168            .unwrap()
2169    }
2170
2171    #[test]
2172    fn awaiting_input_clock_is_durable_first_write_wins_and_resolve_is_fenced() {
2173        let tmp = TempDir::new().unwrap();
2174        let run_id = "01jxwd0000000000000000000w";
2175        let paths = seed_run_with_node(&tmp, run_id);
2176        let nid = Some(NodeId::parse_str("n-0001").unwrap());
2177        let opened_at: chrono::DateTime<Utc> = "2026-08-16T12:00:00Z".parse().unwrap();
2178        let mut open = event(run_id);
2179        open.seq = 3;
2180        open.ts = opened_at;
2181        open.kind = "node.awaiting_input".into();
2182        open.node_id = nid.clone();
2183        open.data = serde_json::json!({ "discussion_items": [{
2184            "topic": "Which scope?",
2185            "options": ["small", "large"],
2186            "recommended_default": "small"
2187        }] });
2188        apply_event(&paths, &open).unwrap();
2189        let first = read_n0001(&paths).awaiting_input.unwrap();
2190        assert_eq!(first.opened_at, opened_at);
2191        assert_eq!(first.event_seq, 3);
2192
2193        // A duplicate/restarted supervisor observation cannot restart the clock.
2194        let mut duplicate = open.clone();
2195        duplicate.seq = 4;
2196        duplicate.ts = opened_at + chrono::Duration::hours(1);
2197        apply_event(&paths, &duplicate).unwrap();
2198        let still_first = read_n0001(&paths).awaiting_input.unwrap();
2199        assert_eq!(still_first.opened_at, opened_at);
2200        assert_eq!(still_first.event_seq, 3);
2201
2202        // A stale timeout for another generation cannot clear this request.
2203        let mut stale = event(run_id);
2204        stale.seq = 5;
2205        stale.kind = "node.input_resolved".into();
2206        stale.node_id = nid.clone();
2207        stale.data = serde_json::json!({ "event_seq": 2 });
2208        apply_event(&paths, &stale).unwrap();
2209        assert!(read_n0001(&paths).awaiting_input.is_some());
2210
2211        let mut resolved = stale;
2212        resolved.seq = 6;
2213        resolved.data = serde_json::json!({ "event_seq": 3 });
2214        apply_event(&paths, &resolved).unwrap();
2215        assert!(read_n0001(&paths).awaiting_input.is_none());
2216    }
2217
2218    #[test]
2219    fn awaiting_input_rejects_missing_default_without_mutating_projection() {
2220        let tmp = TempDir::new().unwrap();
2221        let run_id = "01jxwd0000000000000000000x";
2222        let paths = seed_run_with_node(&tmp, run_id);
2223        let mut open = event(run_id);
2224        open.seq = 3;
2225        open.kind = "node.awaiting_input".into();
2226        open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2227        open.data = serde_json::json!({ "discussion_items": [{
2228            "topic": "Which scope?", "options": ["small", "large"]
2229        }] });
2230        assert!(reduce_event_to_ops(&paths, &open).is_err());
2231        assert!(read_n0001(&paths).awaiting_input.is_none());
2232    }
2233
2234    #[test]
2235    fn awaiting_input_validation_is_state_independent_and_worker_exit_clears_it() {
2236        let tmp = TempDir::new().unwrap();
2237        let run_id = "01jxwd0000000000000000000y";
2238        let paths = seed_run_with_node(&tmp, run_id);
2239        let nid = Some(NodeId::parse_str("n-0001").unwrap());
2240
2241        let mut open = event(run_id);
2242        open.seq = 3;
2243        open.kind = "node.awaiting_input".into();
2244        open.node_id = nid.clone();
2245        open.data = serde_json::json!({ "discussion_items": [{
2246            "topic": "Which scope?", "options": ["small", "large"],
2247            "recommended_default": "small"
2248        }] });
2249        apply_event(&paths, &open).unwrap();
2250
2251        let mut malformed_duplicate = open.clone();
2252        malformed_duplicate.seq = 4;
2253        malformed_duplicate.data = serde_json::json!({ "discussion_items": [] });
2254        assert!(reduce_event_to_ops(&paths, &malformed_duplicate).is_err());
2255
2256        let mut exited = event(run_id);
2257        exited.seq = 5;
2258        exited.kind = "worker.exited".into();
2259        exited.node_id = nid;
2260        exited.data = serde_json::json!({ "exit_code": 0 });
2261        apply_event(&paths, &exited).unwrap();
2262        let node = read_n0001(&paths);
2263        assert!(node.awaiting_input.is_none());
2264        assert!(node.worker_exit.is_some());
2265
2266        let mut delayed_open = open;
2267        delayed_open.seq = 6;
2268        assert!(reduce_event_to_ops(&paths, &delayed_open)
2269            .unwrap()
2270            .is_empty());
2271    }
2272
2273    #[test]
2274    fn input_resolved_requires_generation_even_when_nothing_is_open() {
2275        let tmp = TempDir::new().unwrap();
2276        let run_id = "01jxwd0000000000000000000z";
2277        let paths = seed_run_with_node(&tmp, run_id);
2278        let mut resolved = event(run_id);
2279        resolved.seq = 3;
2280        resolved.kind = "node.input_resolved".into();
2281        resolved.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2282        resolved.data = serde_json::json!({});
2283        assert!(reduce_event_to_ops(&paths, &resolved).is_err());
2284    }
2285
2286    #[test]
2287    fn awaiting_input_default_must_be_one_of_options() {
2288        let tmp = TempDir::new().unwrap();
2289        let run_id = "01jxwd00000000000000000010";
2290        let paths = seed_run_with_node(&tmp, run_id);
2291        let mut open = event(run_id);
2292        open.seq = 3;
2293        open.kind = "node.awaiting_input".into();
2294        open.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2295        open.data = serde_json::json!({ "discussion_items": [{
2296            "topic": "Which scope?", "options": ["small", "large"],
2297            "recommended_default": "other"
2298        }] });
2299        assert!(reduce_event_to_ops(&paths, &open).is_err());
2300    }
2301
2302    fn merge_started_event(run_id: &str, seq: u64, op_id: &str, expected: &str) -> Event {
2303        let mut ev = event(run_id);
2304        ev.seq = seq;
2305        ev.kind = KIND_MERGE_STARTED.into();
2306        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2307        ev.data = serde_json::json!({
2308            "op_id": op_id,
2309            "source_branch": "main",
2310            "worker_branch": "wt/worker",
2311            "expected_source_oid": expected,
2312            "worker_oid": "cafebabecafebabecafebabecafebabecafebabe",
2313            "base_sha": null,
2314            "driver_pid": 4242,
2315            "driver_pid_start_secs": null,
2316            "started_at": "2026-08-15T00:00:00Z",
2317        });
2318        ev
2319    }
2320
2321    /// `merge.started` records the in-flight transaction on `pending_merge`
2322    /// without transitioning the node's status.
2323    #[test]
2324    fn merge_started_records_pending_transaction() {
2325        let tmp = TempDir::new().unwrap();
2326        let run_id = "01jxsnap000000000000000000";
2327        let paths = seed_run_with_node(&tmp, run_id);
2328
2329        apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2330        let n = read_n0001(&paths);
2331        assert_eq!(
2332            n.status,
2333            Status::Pending,
2334            "recording a merge is not terminal"
2335        );
2336        let txn = n.pending_merge.expect("transaction recorded");
2337        assert_eq!(txn.op_id, "op-1");
2338        assert_eq!(txn.expected_source_oid, "aaa");
2339    }
2340
2341    /// `merge.aborted` clears the pending transaction it names, leaving the node
2342    /// live; a stale abort for a different `op_id` is a clean no-op.
2343    #[test]
2344    fn merge_aborted_clears_matching_transaction_only() {
2345        let tmp = TempDir::new().unwrap();
2346        let run_id = "01jxsnap000000000000000000";
2347        let paths = seed_run_with_node(&tmp, run_id);
2348        apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2349
2350        // A stale abort for a different op_id does nothing.
2351        let mut stale = event(run_id);
2352        stale.seq = 4;
2353        stale.kind = KIND_MERGE_ABORTED.into();
2354        stale.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2355        stale.data = serde_json::json!({ "op_id": "op-OTHER", "reason": "x" });
2356        apply_event(&paths, &stale).unwrap();
2357        assert!(
2358            read_n0001(&paths).pending_merge.is_some(),
2359            "stale abort is a no-op"
2360        );
2361
2362        // The matching abort clears it; the node stays live.
2363        let mut abort = event(run_id);
2364        abort.seq = 5;
2365        abort.kind = KIND_MERGE_ABORTED.into();
2366        abort.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2367        abort.data = serde_json::json!({ "op_id": "op-1", "reason": "no mutation" });
2368        apply_event(&paths, &abort).unwrap();
2369        let n = read_n0001(&paths);
2370        assert!(
2371            n.pending_merge.is_none(),
2372            "matching abort clears the transaction"
2373        );
2374        assert_eq!(n.status, Status::Pending, "abort does not terminalize");
2375    }
2376
2377    /// A terminal `node.report` (the normal, no-crash completion) clears any
2378    /// pending merge transaction.
2379    #[test]
2380    fn terminal_report_clears_pending_merge() {
2381        let tmp = TempDir::new().unwrap();
2382        let run_id = "01jxsnap000000000000000000";
2383        let paths = seed_run_with_node(&tmp, run_id);
2384        apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2385
2386        let mut report = event(run_id);
2387        report.seq = 4;
2388        report.kind = "node.report".into();
2389        report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2390        report.data = serde_json::json!({ "success": true, "via": "explicit-merge" });
2391        apply_event(&paths, &report).unwrap();
2392        let n = read_n0001(&paths);
2393        assert_eq!(n.status, Status::Done);
2394        assert!(
2395            n.pending_merge.is_none(),
2396            "completed merge clears the transaction"
2397        );
2398    }
2399
2400    /// Regression (issue `retire-via-string`): the terminal-node adoption
2401    /// exception now keys on the typed `RunMerge` origin, NOT a forgeable `via`
2402    /// string. A late report against a `Failed` node that carries an `Agent`
2403    /// origin (as every `node report` self-submission does) plus a forged
2404    /// `via: "explicit-merge"` must NOT be adopted — the node stays `Failed`. A
2405    /// present-but-malformed origin is likewise not adopted. Only a genuine
2406    /// `RunMerge`-origin report (or a legacy report with NO origin field) is
2407    /// adopted and corrects the node to `Done`.
2408    #[test]
2409    fn late_merge_adoption_requires_run_merge_origin_not_forged_via() {
2410        let tmp = TempDir::new().unwrap();
2411
2412        // Helper: seed a fresh run (distinct id), drive n-0001 to Failed, apply a
2413        // late report, and return the resulting node status.
2414        let drive = |run_id: &str, report_data: Value| -> Status {
2415            let paths = seed_run_with_node(&tmp, run_id);
2416            // Terminalize the node as Failed (a watchdog-synthesized failure).
2417            let mut fail = event(run_id);
2418            fail.seq = 3;
2419            fail.kind = "node.status".into();
2420            fail.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2421            fail.data = serde_json::json!({ "status": "failed" });
2422            apply_event(&paths, &fail).unwrap();
2423            assert_eq!(read_n0001(&paths).status, Status::Failed);
2424            // The late report under test.
2425            let mut report = event(run_id);
2426            report.seq = 4;
2427            report.kind = "node.report".into();
2428            report.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2429            report.data = report_data;
2430            apply_event(&paths, &report).unwrap();
2431            read_n0001(&paths).status
2432        };
2433
2434        // Forged: Agent origin + a hand-set `via` — NOT adopted, stays Failed.
2435        let mut agent_forged = serde_json::json!({ "success": true, "via": "explicit-merge" });
2436        crate::ReportOrigin::Agent.stamp(&mut agent_forged);
2437        assert_eq!(
2438            drive("01jxsnap000000000000000001", agent_forged),
2439            Status::Failed,
2440            "an Agent-origin report with a forged via must not be adopted"
2441        );
2442
2443        // Present-but-malformed origin + forged via — NOT adopted, stays Failed.
2444        let malformed = serde_json::json!({
2445            "success": true, "via": "explicit-merge", "origin": "garbage-not-an-object"
2446        });
2447        assert_eq!(
2448            drive("01jxsnap000000000000000002", malformed),
2449            Status::Failed,
2450            "a malformed origin must not re-unlock the legacy via adoption path"
2451        );
2452
2453        // Genuine RunMerge origin (no `via` at all) — adopted, corrected to Done.
2454        let mut run_merge = serde_json::json!({ "success": true });
2455        crate::ReportOrigin::RunMerge {
2456            op_id: Some("op-1".into()),
2457            worker_oid: Some("cafebabe".into()),
2458        }
2459        .stamp(&mut run_merge);
2460        assert_eq!(
2461            drive("01jxsnap000000000000000003", run_merge),
2462            Status::Done,
2463            "a genuine RunMerge-origin report is adopted and corrects Failed→Done"
2464        );
2465
2466        // Legacy report (no origin field) with `via` — still adopted (backward
2467        // compat with pre-typed-origin on-disk runs).
2468        let legacy = serde_json::json!({ "success": true, "via": "explicit-merge" });
2469        assert_eq!(
2470            drive("01jxsnap000000000000000004", legacy),
2471            Status::Done,
2472            "a legacy via-only report (no origin field) is still adopted"
2473        );
2474    }
2475
2476    /// A terminal `node.status` (e.g. a watchdog-synthesized failure) clears any
2477    /// in-flight merge transaction, so `pending_merge` is never stranded on a
2478    /// terminal node where recovery would refuse to look (/llm-review finding).
2479    #[test]
2480    fn terminal_node_status_clears_pending_merge() {
2481        let tmp = TempDir::new().unwrap();
2482        let run_id = "01jxsnap000000000000000000";
2483        let paths = seed_run_with_node(&tmp, run_id);
2484        apply_event(&paths, &merge_started_event(run_id, 3, "op-1", "aaa")).unwrap();
2485        assert!(read_n0001(&paths).pending_merge.is_some());
2486
2487        let mut status = event(run_id);
2488        status.seq = 4;
2489        status.kind = "node.status".into();
2490        status.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2491        status.data = serde_json::json!({ "status": "failed" });
2492        apply_event(&paths, &status).unwrap();
2493        let n = read_n0001(&paths);
2494        assert_eq!(n.status, Status::Failed);
2495        assert!(
2496            n.pending_merge.is_none(),
2497            "terminal status clears the transaction"
2498        );
2499    }
2500
2501    /// Replaying `supervisor.attached` from scratch reproduces
2502    /// `Node.supervisor_pid` — the field is now event-sourced, not a
2503    /// projection-only write (issue `supervisor-state-not-event-sourced`).
2504    #[test]
2505    fn supervisor_attached_sets_supervisor_pid() {
2506        let tmp = TempDir::new().unwrap();
2507        let run_id = "01jxsnap000000000000000000";
2508        let paths = seed_run_with_node(&tmp, run_id);
2509        assert_eq!(read_n0001(&paths).supervisor_pid, None);
2510
2511        let mut ev = event(run_id);
2512        ev.seq = 3;
2513        ev.kind = "supervisor.attached".into();
2514        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2515        ev.data = serde_json::json!({ "pid": 47820 });
2516        apply_event(&paths, &ev).expect("supervisor.attached applies");
2517        assert_eq!(read_n0001(&paths).supervisor_pid, Some(47820));
2518    }
2519
2520    /// A second attach with a different pid overrides (latest-wins); a replay
2521    /// of the *same* pid is a pure no-op that does not churn `updated_at`.
2522    #[test]
2523    fn supervisor_attached_latest_wins_and_idempotent_on_replay() {
2524        let tmp = TempDir::new().unwrap();
2525        let run_id = "01jxsnap000000000000000000";
2526        let paths = seed_run_with_node(&tmp, run_id);
2527
2528        let mut ev = event(run_id);
2529        ev.kind = "supervisor.attached".into();
2530        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2531
2532        ev.seq = 3;
2533        ev.data = serde_json::json!({ "pid": 100 });
2534        apply_event(&paths, &ev).expect("first attach applies");
2535        assert_eq!(read_n0001(&paths).supervisor_pid, Some(100));
2536
2537        // A restart binds a fresh pid: latest-wins.
2538        ev.seq = 4;
2539        ev.data = serde_json::json!({ "pid": 200 });
2540        apply_event(&paths, &ev).expect("second attach applies");
2541        let after_second = read_n0001(&paths);
2542        assert_eq!(after_second.supervisor_pid, Some(200));
2543
2544        // Replaying the latest event again is a no-op: the planned ops are
2545        // empty and the projection bytes (including `updated_at`) are unchanged.
2546        let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2547        assert!(ops.is_empty(), "re-applying same pid must plan no ops");
2548        apply_event(&paths, &ev).expect("replay applies as no-op");
2549        assert_eq!(read_n0001(&paths).updated_at, after_second.updated_at);
2550    }
2551
2552    /// Replaying `supervisor.cursor_advanced` from scratch reproduces
2553    /// `Node.last_processed_report_seq_by_child`.
2554    #[test]
2555    fn supervisor_cursor_advanced_sets_report_cursor() {
2556        let tmp = TempDir::new().unwrap();
2557        let run_id = "01jxsnap000000000000000000";
2558        let paths = seed_run_with_node(&tmp, run_id);
2559        let child = "02jxsnap000000000000000000";
2560
2561        let mut ev = event(run_id);
2562        ev.seq = 3;
2563        ev.kind = "supervisor.cursor_advanced".into();
2564        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2565        ev.data = serde_json::json!({ "child_run_id": child, "report_seq": 7 });
2566        apply_event(&paths, &ev).expect("cursor_advanced applies");
2567
2568        let n = read_n0001(&paths);
2569        assert_eq!(
2570            n.last_processed_report_seq_by_child.get(child),
2571            Some(&Value::from(7u64))
2572        );
2573    }
2574
2575    /// The cursor is monotonic and idempotent: re-applying the same
2576    /// `(child_run_id, report_seq)` is a no-op, an older seq never moves the
2577    /// cursor backward, and a higher seq advances it. A second distinct child
2578    /// gets its own independent entry.
2579    #[test]
2580    fn supervisor_cursor_advanced_is_monotonic_and_idempotent() {
2581        let tmp = TempDir::new().unwrap();
2582        let run_id = "01jxsnap000000000000000000";
2583        let paths = seed_run_with_node(&tmp, run_id);
2584        let child_a = "02jxsnap000000000000000000";
2585        let child_b = "03jxsnap000000000000000000";
2586
2587        let mut ev = event(run_id);
2588        ev.kind = "supervisor.cursor_advanced".into();
2589        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2590
2591        ev.seq = 3;
2592        ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 5 });
2593        apply_event(&paths, &ev).expect("seq 5 applies");
2594
2595        // Replay the exact same event — no-op, plans zero ops.
2596        let ops = reduce_event_to_ops(&paths, &ev).expect("replay reduces cleanly");
2597        assert!(ops.is_empty(), "re-applying same cursor must plan no ops");
2598
2599        // An older seq must not move the cursor backward.
2600        ev.seq = 4;
2601        ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 3 });
2602        let ops = reduce_event_to_ops(&paths, &ev).expect("older seq reduces cleanly");
2603        assert!(ops.is_empty(), "older seq must plan no ops");
2604        apply_event(&paths, &ev).expect("older seq applies as no-op");
2605        assert_eq!(
2606            read_n0001(&paths)
2607                .last_processed_report_seq_by_child
2608                .get(child_a),
2609            Some(&Value::from(5u64))
2610        );
2611
2612        // A higher seq advances; an independent child gets its own entry.
2613        ev.seq = 5;
2614        ev.data = serde_json::json!({ "child_run_id": child_a, "report_seq": 9 });
2615        apply_event(&paths, &ev).expect("higher seq applies");
2616        ev.seq = 6;
2617        ev.data = serde_json::json!({ "child_run_id": child_b, "report_seq": 1 });
2618        apply_event(&paths, &ev).expect("second child applies");
2619
2620        let n = read_n0001(&paths);
2621        assert_eq!(
2622            n.last_processed_report_seq_by_child.get(child_a),
2623            Some(&Value::from(9u64))
2624        );
2625        assert_eq!(
2626            n.last_processed_report_seq_by_child.get(child_b),
2627            Some(&Value::from(1u64))
2628        );
2629    }
2630
2631    /// Both new kinds reject a malformed payload at the reducer boundary so a
2632    /// forged event can never write a corrupt projection.
2633    #[test]
2634    fn supervisor_state_events_reject_malformed_payloads() {
2635        let tmp = TempDir::new().unwrap();
2636        let run_id = "01jxsnap000000000000000000";
2637        let paths = seed_run_with_node(&tmp, run_id);
2638        let nid = Some(NodeId::parse_str("n-0001").unwrap());
2639
2640        // Missing pid.
2641        let mut ev = event(run_id);
2642        ev.seq = 3;
2643        ev.kind = "supervisor.attached".into();
2644        ev.node_id = nid.clone();
2645        ev.data = serde_json::json!({});
2646        assert!(matches!(
2647            reduce_event_to_ops(&paths, &ev),
2648            Err(Error::CorruptEventLog { .. })
2649        ));
2650
2651        // Missing envelope node_id.
2652        ev.node_id = None;
2653        ev.data = serde_json::json!({ "pid": 1 });
2654        assert!(matches!(
2655            reduce_event_to_ops(&paths, &ev),
2656            Err(Error::CorruptEventLog { .. })
2657        ));
2658
2659        // cursor_advanced: malformed child_run_id.
2660        let mut ev2 = event(run_id);
2661        ev2.seq = 4;
2662        ev2.kind = "supervisor.cursor_advanced".into();
2663        ev2.node_id = nid.clone();
2664        ev2.data = serde_json::json!({ "child_run_id": "../etc", "report_seq": 1 });
2665        assert!(matches!(
2666            reduce_event_to_ops(&paths, &ev2),
2667            Err(Error::CorruptEventLog { .. })
2668        ));
2669
2670        // cursor_advanced: missing report_seq.
2671        ev2.data = serde_json::json!({ "child_run_id": "02jxsnap000000000000000000" });
2672        assert!(matches!(
2673            reduce_event_to_ops(&paths, &ev2),
2674            Err(Error::CorruptEventLog { .. })
2675        ));
2676    }
2677
2678    /// The append gate stays fail-closed after the 0.2 cut added
2679    /// `Kind`'s `#[serde(other)]` catch-all: a `run.created` / `node.created`
2680    /// whose `kind` is a removed kind (`code`, …) or plain garbage must still be
2681    /// rejected as `CorruptEventLog`, NOT silently accepted as `Kind::Unknown`.
2682    /// (Legacy runs are never re-created through the reducer — their manifest is
2683    /// read directly from disk via the permissive `Kind::Unknown` decode.)
2684    #[test]
2685    fn removed_or_garbage_kind_in_created_events_is_rejected() {
2686        let tmp = TempDir::new().unwrap();
2687        let run_id = "01jxsnap000000000000000000";
2688        let rid = RunId::parse_str(run_id).unwrap();
2689        let dir = crate::run_dir(tmp.path(), &rid);
2690        std::fs::create_dir_all(&dir).unwrap();
2691        let paths = RunPaths::new(dir, run_id).unwrap();
2692
2693        for bad in ["code", "orchestrate", "bugfix", "make-skill", "garbage"] {
2694            let mut ev = event(run_id);
2695            ev.kind = "run.created".into();
2696            ev.node_id = None;
2697            ev.data = serde_json::json!({ "kind": bad, "lifecycle": "autonomous", "title": "t" });
2698            assert!(
2699                matches!(
2700                    reduce_event_to_ops(&paths, &ev),
2701                    Err(Error::CorruptEventLog { .. })
2702                ),
2703                "run.created with kind {bad:?} must be rejected, not folded to Unknown"
2704            );
2705        }
2706
2707        // A surviving creatable kind still folds cleanly (guards against a
2708        // false positive that rejects everything).
2709        let mut ok = event(run_id);
2710        ok.kind = "run.created".into();
2711        ok.node_id = None;
2712        ok.data = serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" });
2713        assert!(reduce_event_to_ops(&paths, &ok).is_ok());
2714    }
2715
2716    /// Snapshot every projection file under `paths` to a `path → inode` map.
2717    ///
2718    /// An atomic projection write is temp-file + rename, so a rewritten file
2719    /// always lands a *fresh inode* — even when its bytes are byte-for-byte
2720    /// identical (e.g. a manifest op that refreshes `updated_at` to the same
2721    /// timestamp). Comparing inodes therefore detects every write the reducer
2722    /// makes, with no false negatives a content diff would suffer. `events.jsonl`
2723    /// and `.lock` are excluded: `apply_event` never touches them.
2724    #[cfg(unix)]
2725    fn projection_inodes(paths: &RunPaths) -> std::collections::BTreeMap<PathBuf, u64> {
2726        use std::os::unix::fs::MetadataExt;
2727        let mut consider = vec![paths.manifest()];
2728        for dir in [paths.nodes_dir()] {
2729            if let Ok(rd) = std::fs::read_dir(&dir) {
2730                for ent in rd.flatten() {
2731                    let p = ent.path();
2732                    if p.extension().and_then(|s| s.to_str()) == Some("json") {
2733                        consider.push(p);
2734                    }
2735                }
2736            }
2737        }
2738        let mut map = std::collections::BTreeMap::new();
2739        for p in consider {
2740            if let Ok(md) = std::fs::symlink_metadata(&p) {
2741                if md.file_type().is_file() {
2742                    map.insert(p, md.ino());
2743                }
2744            }
2745        }
2746        map
2747    }
2748
2749    /// The exhaustive parity guarantee `projected-paths-into-reducer` requires:
2750    /// for an event applied against a given state, the paths
2751    /// [`plan_projections`] reports MUST equal the files [`apply_event`]
2752    /// actually writes. Plan first (against pre-apply state), apply, then diff
2753    /// the projection inodes — a file is "written" iff it is newly present or
2754    /// its inode changed. `expect_writes` guards the test itself: when set, the
2755    /// touched set must be non-empty, so a kind that silently stopped writing
2756    /// can't pass by matching an empty plan against an empty diff.
2757    #[cfg(unix)]
2758    fn assert_plan_matches_apply(paths: &RunPaths, ev: &Event, expect_writes: bool) {
2759        use std::collections::BTreeSet;
2760        let before = projection_inodes(paths);
2761        let planned: BTreeSet<PathBuf> = plan_projections(paths, ev)
2762            .unwrap_or_else(|e| panic!("plan_projections({}) errored: {e:?}", ev.kind))
2763            .into_iter()
2764            .collect();
2765        apply_event(paths, ev)
2766            .unwrap_or_else(|e| panic!("apply_event({}) errored: {e:?}", ev.kind));
2767        let after = projection_inodes(paths);
2768        let touched: BTreeSet<PathBuf> = after
2769            .iter()
2770            .filter(|(p, ino)| before.get(*p) != Some(*ino))
2771            .map(|(p, _)| p.clone())
2772            .collect();
2773        assert_eq!(
2774            planned, touched,
2775            "kind={}: plan_projections must name exactly the files apply_event writes",
2776            ev.kind
2777        );
2778        if expect_writes {
2779            assert!(
2780                !touched.is_empty(),
2781                "kind={}: expected this event to write at least one projection",
2782                ev.kind
2783            );
2784        }
2785    }
2786
2787    /// Drive every event kind through a dependency-ordered lifecycle on real
2788    /// runs, asserting plan/apply parity at each step. Covers the writing kinds
2789    /// (run/node/supervisor/child) in states where they
2790    /// project, plus the no-op kinds (audit records, `supervisor.exited`,
2791    /// terminal-guarded transitions) where both the plan and the apply touch
2792    /// nothing.
2793    #[cfg(unix)]
2794    #[test]
2795    fn plan_projections_matches_apply_for_every_kind() {
2796        let tmp = TempDir::new().unwrap();
2797        let run_id = "01jxsnap000000000000000000";
2798        let rid = RunId::parse_str(run_id).unwrap();
2799        let dir = crate::run_dir(tmp.path(), &rid);
2800        std::fs::create_dir_all(&dir).unwrap();
2801        let paths = RunPaths::new(dir, run_id).unwrap();
2802        let nid = || Some(NodeId::parse_str("n-0001").unwrap());
2803        let child = "02jxsnap000000000000000000";
2804
2805        // Helper to build a fresh envelope at a monotonic seq.
2806        let mut next_seq = 0u64;
2807        let mut at = |kind: &str, node_id, data| {
2808            next_seq += 1;
2809            Event {
2810                ts: Utc::now(),
2811                seq: next_seq,
2812                kind: kind.into(),
2813                run_id: rid.clone(),
2814                node_id,
2815                idempotency_key: None,
2816                data,
2817            }
2818        };
2819
2820        // run.created → manifest.json
2821        assert_plan_matches_apply(
2822            &paths,
2823            &at(
2824                "run.created",
2825                None,
2826                serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
2827            ),
2828            true,
2829        );
2830        // run.status (pending → running) → manifest.json
2831        assert_plan_matches_apply(
2832            &paths,
2833            &at(
2834                "run.status",
2835                None,
2836                serde_json::json!({ "status": "running" }),
2837            ),
2838            true,
2839        );
2840        // node.created → nodes/n-0001.json + manifest.json
2841        assert_plan_matches_apply(
2842            &paths,
2843            &at(
2844                "node.created",
2845                nid(),
2846                serde_json::json!({ "kind": "spinoff" }),
2847            ),
2848            true,
2849        );
2850        // node.status (pending → running) → nodes/n-0001.json
2851        assert_plan_matches_apply(
2852            &paths,
2853            &at(
2854                "node.status",
2855                nid(),
2856                serde_json::json!({ "status": "running" }),
2857            ),
2858            true,
2859        );
2860        // supervisor.attached → nodes/n-0001.json (still non-terminal)
2861        assert_plan_matches_apply(
2862            &paths,
2863            &at(
2864                "supervisor.attached",
2865                nid(),
2866                serde_json::json!({ "pid": 4242 }),
2867            ),
2868            true,
2869        );
2870        // supervisor.cursor_advanced → nodes/n-0001.json
2871        assert_plan_matches_apply(
2872            &paths,
2873            &at(
2874                "supervisor.cursor_advanced",
2875                nid(),
2876                serde_json::json!({ "child_run_id": child, "report_seq": 3 }),
2877            ),
2878            true,
2879        );
2880        // child.spawned → nodes/n-0001.json (parent node)
2881        assert_plan_matches_apply(
2882            &paths,
2883            &at(
2884                "child.spawned",
2885                nid(),
2886                serde_json::json!({ "child_run_id": child, "child_node_id": "n-0001" }),
2887            ),
2888            true,
2889        );
2890        // node.report success → nodes/n-0001.json (now terminal)
2891        assert_plan_matches_apply(
2892            &paths,
2893            &at("node.report", nid(), serde_json::json!({ "success": true })),
2894            true,
2895        );
2896        // Terminal-guarded no-ops: a settled node swallows further transitions,
2897        // so both the plan and the apply touch nothing.
2898        assert_plan_matches_apply(
2899            &paths,
2900            &at(
2901                "node.status",
2902                nid(),
2903                serde_json::json!({ "status": "failed" }),
2904            ),
2905            false,
2906        );
2907        // No-op audit / lifecycle kinds: zero projections by design.
2908        for kind in [
2909            "supervisor.exited",
2910            "orchestrator.decision",
2911            "discuss.critical",
2912            "cleanup.window_missing",
2913        ] {
2914            assert_plan_matches_apply(&paths, &at(kind, None, serde_json::json!({})), false);
2915        }
2916    }
2917
2918    /// A `worker.exited` carrying a clean `exit_code: 0` folds onto the node's
2919    /// `worker_exit` field as a clean exit — and does NOT transition `status`
2920    /// (terminalization is the supervisor's decision via the typed table).
2921    #[test]
2922    fn worker_exited_records_clean_exit_without_transitioning_status() {
2923        let tmp = TempDir::new().unwrap();
2924        let run_id = "01jxsnap000000000000000000";
2925        let paths = bootstrap_retry_node(&tmp, run_id);
2926
2927        let mut ev = event(run_id);
2928        ev.seq = 3;
2929        ev.kind = "worker.exited".into();
2930        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2931        ev.data = serde_json::json!({ "exit_code": 0 });
2932        apply_event(&paths, &ev).expect("worker.exited applies");
2933
2934        let n = read_node_opt(&paths, &NodeId::parse_str("n-0001").unwrap())
2935            .unwrap()
2936            .unwrap();
2937        let exit = n.worker_exit.expect("worker_exit recorded");
2938        assert_eq!(exit.code, Some(0));
2939        assert_eq!(exit.signal, None);
2940        assert!(exit.is_clean());
2941        assert_eq!(
2942            n.status,
2943            Status::Pending,
2944            "the exit fact never transitions status"
2945        );
2946    }
2947
2948    /// A `worker.exited` carrying a `signal` records it as a failure; and the fold
2949    /// is first-write-wins — a replayed/duplicate exit event never overwrites the
2950    /// first recorded fact (replay-safety for the `applied_seq` watermark).
2951    #[test]
2952    fn worker_exited_records_signal_and_is_first_write_wins() {
2953        let tmp = TempDir::new().unwrap();
2954        let run_id = "01jxsnap000000000000000000";
2955        let paths = bootstrap_retry_node(&tmp, run_id);
2956        let nid = NodeId::parse_str("n-0001").unwrap();
2957
2958        let mut ev = event(run_id);
2959        ev.seq = 3;
2960        ev.kind = "worker.exited".into();
2961        ev.node_id = Some(nid.clone());
2962        ev.data = serde_json::json!({ "signal": 9 });
2963        apply_event(&paths, &ev).expect("worker.exited applies");
2964
2965        let n = read_node_opt(&paths, &nid).unwrap().unwrap();
2966        let exit = n.worker_exit.expect("worker_exit recorded");
2967        assert_eq!(exit.signal, Some(9));
2968        assert!(exit.is_failure());
2969
2970        // A later, conflicting exit event (e.g. a replay of a different value) is a
2971        // clean no-op: the first fact stands.
2972        let mut dup = event(run_id);
2973        dup.seq = 4;
2974        dup.kind = "worker.exited".into();
2975        dup.node_id = Some(nid.clone());
2976        dup.data = serde_json::json!({ "exit_code": 0 });
2977        apply_event(&paths, &dup).expect("duplicate worker.exited applies as no-op");
2978        let n2 = read_node_opt(&paths, &nid).unwrap().unwrap();
2979        assert_eq!(
2980            n2.worker_exit.unwrap().signal,
2981            Some(9),
2982            "first-write-wins: the replayed exit must not overwrite the recorded fact"
2983        );
2984    }
2985
2986    /// A `worker.exited` carrying neither `exit_code` nor `signal` is malformed —
2987    /// the reducer is the canonical gate and rejects it as `CorruptEventLog` rather
2988    /// than record an empty fact.
2989    #[test]
2990    fn worker_exited_without_code_or_signal_is_corrupt() {
2991        let tmp = TempDir::new().unwrap();
2992        let run_id = "01jxsnap000000000000000000";
2993        let paths = bootstrap_retry_node(&tmp, run_id);
2994
2995        let mut ev = event(run_id);
2996        ev.seq = 3;
2997        ev.kind = "worker.exited".into();
2998        ev.node_id = Some(NodeId::parse_str("n-0001").unwrap());
2999        ev.data = serde_json::json!({});
3000        match reduce_event_to_ops(&paths, &ev) {
3001            Err(Error::CorruptEventLog { .. }) => {}
3002            Ok(_) => panic!("an empty worker.exited payload must be rejected, not applied"),
3003            Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3004        }
3005
3006        // Carrying BOTH is contradictory (a process cannot both return a code and
3007        // be killed) — also rejected.
3008        ev.data = serde_json::json!({ "exit_code": 0, "signal": 9 });
3009        match reduce_event_to_ops(&paths, &ev) {
3010            Err(Error::CorruptEventLog { .. }) => {}
3011            Ok(_) => panic!("a worker.exited with both fields must be rejected"),
3012            Err(other) => panic!("expected CorruptEventLog, got {other:?}"),
3013        }
3014    }
3015
3016    /// `node.death_observed` records the residual crash backstop's first-death
3017    /// anchor (`first_death_at`) as `ev.ts`, is **first-write-wins** (a later
3018    /// re-observation never resets the monotonic anchor), and is a no-op against a
3019    /// terminal node (the backstop is moot once settled). Issue
3020    /// `typed-supervisor-outcomes`.
3021    #[test]
3022    fn node_death_observed_records_first_death_first_write_wins() {
3023        let tmp = TempDir::new().unwrap();
3024        let run_id = "01jxsnap000000000000000000";
3025        let paths = bootstrap_retry_node(&tmp, run_id);
3026        let nid = NodeId::parse_str("n-0001").unwrap();
3027
3028        let mut ev = event(run_id);
3029        ev.seq = 3;
3030        ev.kind = "node.death_observed".into();
3031        ev.node_id = Some(nid.clone());
3032        ev.data = serde_json::json!({});
3033        apply_event(&paths, &ev).expect("node.death_observed applies");
3034        let first = read_node_opt(&paths, &nid)
3035            .unwrap()
3036            .unwrap()
3037            .first_death_at
3038            .expect("first_death_at recorded");
3039        assert_eq!(first, ev.ts, "the anchor is the event's own timestamp");
3040
3041        // A later re-observation is first-write-wins: the monotonic anchor holds.
3042        let mut later = event(run_id);
3043        later.seq = 4;
3044        later.kind = "node.death_observed".into();
3045        later.node_id = Some(nid.clone());
3046        later.ts = ev.ts + chrono::Duration::seconds(30);
3047        later.data = serde_json::json!({});
3048        apply_event(&paths, &later).expect("re-observation applies as no-op");
3049        assert_eq!(
3050            read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3051            Some(first),
3052            "first-write-wins: a re-observation must not reset the anchor"
3053        );
3054    }
3055
3056    /// `node.death_observed` is a no-op once a higher-fidelity fact exists — here a
3057    /// told `worker.exited` — so a from-scratch replay converges to the same state
3058    /// the supervisor's lock-guarded emitter would produce (the backstop is moot
3059    /// once the shim recorded a real exit). Issue `typed-supervisor-outcomes`.
3060    #[test]
3061    fn node_death_observed_noop_when_worker_exit_present() {
3062        let tmp = TempDir::new().unwrap();
3063        let run_id = "01jxsnap000000000000000000";
3064        let paths = bootstrap_retry_node(&tmp, run_id);
3065        let nid = NodeId::parse_str("n-0001").unwrap();
3066
3067        // A told exit lands first.
3068        let mut exit = event(run_id);
3069        exit.seq = 3;
3070        exit.kind = "worker.exited".into();
3071        exit.node_id = Some(nid.clone());
3072        exit.data = serde_json::json!({ "exit_code": 0 });
3073        apply_event(&paths, &exit).unwrap();
3074
3075        // A death observation for the same node folds to nothing.
3076        let mut death = event(run_id);
3077        death.seq = 4;
3078        death.kind = "node.death_observed".into();
3079        death.node_id = Some(nid.clone());
3080        death.data = serde_json::json!({});
3081        apply_event(&paths, &death).expect("applies as no-op");
3082        assert_eq!(
3083            read_node_opt(&paths, &nid).unwrap().unwrap().first_death_at,
3084            None,
3085            "a told worker.exited makes the crash backstop moot; no anchor recorded"
3086        );
3087    }
3088
3089    /// `node.retry` clears the previous attempt's told exit fact: the re-spawned
3090    /// worker is a NEW process, so a stale `worker_exit` must not carry over (it
3091    /// would make the supervisor mis-judge the fresh attempt from the dead one's
3092    /// exit). Issue `thin-exit-status-launcher`.
3093    #[test]
3094    fn node_retry_clears_worker_exit() {
3095        let tmp = TempDir::new().unwrap();
3096        let run_id = "01jxsnap000000000000000000";
3097        let paths = bootstrap_retry_node(&tmp, run_id);
3098        let nid = NodeId::parse_str("n-0001").unwrap();
3099
3100        // Record a failing exit on the first attempt.
3101        let mut exit = event(run_id);
3102        exit.seq = 3;
3103        exit.kind = "worker.exited".into();
3104        exit.node_id = Some(nid.clone());
3105        exit.data = serde_json::json!({ "exit_code": 7 });
3106        apply_event(&paths, &exit).unwrap();
3107        assert!(read_node_opt(&paths, &nid)
3108            .unwrap()
3109            .unwrap()
3110            .worker_exit
3111            .is_some());
3112
3113        // Retry re-spawns the node — the stale exit fact must be gone.
3114        let mut retry = event(run_id);
3115        retry.seq = 4;
3116        retry.kind = "node.retry".into();
3117        retry.node_id = Some(nid.clone());
3118        retry.data = serde_json::json!({
3119            "attempt": 1,
3120            "reason": "agent-died",
3121            "branch": "wt/foo",
3122            "worktree_path": "/tmp/new-wt",
3123            "agent_pid": 222,
3124        });
3125        apply_event(&paths, &retry).unwrap();
3126
3127        let n = read_node_opt(&paths, &nid).unwrap().unwrap();
3128        assert!(
3129            n.worker_exit.is_none(),
3130            "node.retry must clear the previous attempt's worker_exit"
3131        );
3132        assert_eq!(
3133            n.status,
3134            Status::Pending,
3135            "retry returns the node to Pending"
3136        );
3137    }
3138}