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