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