Skip to main content

taskfleet_core/
reducer.rs

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