Skip to main content

taskfleet_core/
cancel.rs

1//! Single-lock run cancellation.
2//!
3//! `run cancel` does three things under **one** held [`RunLock`]: refuse a run
4//! that is already in a non-cancelled terminal state, synthesize a terminal
5//! `node.report` for every still-live node, and append `run.status: cancelled`
6//! once. Holding one lock for the whole operation serializes it against other
7//! *cooperating* writers (those that honor the lock) so the node reads and the
8//! node-report appends can't interleave — which is what made the pre-refactor
9//! CLI loop both racy and prone to over-reporting `cancelled_nodes` (it pushed
10//! a node id even when the per-node append landed after another process had
11//! already settled the node, so the reducer dropped it). Under one lock the
12//! node we read is the node we cancel, so the reported count is honest.
13//!
14//! This is **not crash-atomic**: each `append_and_apply_unlocked` is its own
15//! durable append, so a crash or I/O error partway through can leave some nodes
16//! cancelled and `run.status` not yet appended. Recovery is convergent — a
17//! re-`cancel` of an already-`Cancelled` run scans the still-live stragglers
18//! and finishes the job — not transactional rollback.
19//!
20//! Two consistency properties beyond the single lock:
21//!
22//! - **Enumeration *and per-node liveness* are from the event log, not the
23//!   projection directory.** The node set and each node's current status are
24//!   both replayed from `events.jsonl` (the source of truth) in one streaming
25//!   pass rather than scanned from `nodes/*.json`. A `node.created` can be
26//!   appended+fsynced while its projection write is crash-interrupted
27//!   (`events.rs` documents the log leading the projections); a `nodes/` scan
28//!   would silently drop that node, mark the run `cancelled`, and let a future
29//!   `rebuild_projections` resurrect it as live under a `Cancelled` run. Walking
30//!   the log closes that window. Crucially, replaying `node.status` / `node.report`
31//!   to derive each node's status (rather than trusting `read_node_opt`) closes a
32//!   second window: a *non-cancel* terminal event (e.g. a `node.report`
33//!   `success: true`) fsynced but not yet folded leaves a stale-live projection,
34//!   and a projection-derived liveness check would over-write that node with a
35//!   fresh cancel that diverges on rebuild. The log-derived status settles it as
36//!   already-terminal instead — the log wins. (The manifest's `node_count` is
37//!   *also* a projection written in the same interrupted fold, so it is no more
38//!   authoritative than `nodes/` — and it carries no node ids — which is why we
39//!   replay the log rather than trust the counter.)
40//!
41//! - **Each synthesized event carries a deterministic idempotency key**
42//!   (`run-cancel:<run_id>:node:<node_id>` and `run-cancel:<run_id>:run-status`).
43//!   If a crash lands an append+fsync but interrupts the projection fold, the
44//!   node/run still reads non-terminal, so a re-`cancel` would append a *second*
45//!   logical-cancel event (duplicating it for auditors, metrics, and rebuild).
46//!   The prior cancel events (scoped by `(kind, key)` for this run) are captured
47//!   in the same replay pass, so instead of re-appending, the loop **re-folds
48//!   the already-logged event** via [`apply_event`](crate::reducer) — converging a projection
49//!   the crash left non-terminal *without* a duplicate log line (a re-fold is a
50//!   clean no-op when the projection already agrees). The whole transaction is
51//!   then both non-duplicating and projection-convergent.
52//!
53//! The cancel ledger is built by a *streaming* pass: [`for_each_event_probe`](crate::events)
54//! walks `events.jsonl` line by line, parsing only the small envelope + status
55//! fields each line needs and materializing a full [`Event`] payload solely for
56//! the handful of lines in this run's `run-cancel:<run_id>:` key namespace (the
57//! events the re-fold path replays). The whole log is never held in memory, so
58//! lock-hold time and peak memory stay bounded even for a run with hundreds of
59//! nodes and multi-KB `node.report` payloads.
60//!
61//! What is still *not* derived from the log here: **run-level** liveness (the
62//! terminal-refusal check and `run_was_already_cancelled`) is read from the
63//! manifest projection, with the prior-cancel re-fold converging a crash-stranded
64//! `run.status`. Deriving the run status from the log too would conflate a
65//! crash-stranded `run.status: cancelled` (manifest stale, must re-fold and
66//! report a *fresh* cancel) with an already-folded one, since the log is
67//! identical in both cases — so the manifest read stays authoritative for the
68//! run-level decision, exactly as the per-node convergence path consults
69//! `read_node_opt` only to tell those two cases apart.
70
71use std::collections::HashMap;
72
73use serde::Deserialize;
74use serde_json::{json, Value};
75
76use crate::error::{Error, Result};
77use crate::events::{append_and_apply_unlocked, excerpt, for_each_event_probe};
78use crate::lock::{LockedRun, RunLock};
79use crate::paths::RunPaths;
80use crate::projections::{read_manifest, read_node_opt};
81use crate::reducer::apply_event;
82use crate::report::ReportOrigin;
83use crate::schema::{Event, NodeId, RunId, Status};
84
85/// Outcome of a [`cancel_run`] transaction. Lets a thin CLI wrapper report
86/// honestly what actually changed: which live nodes it converged, which were
87/// already settled (skipped, not double-reported), and whether the run itself
88/// was already cancelled (a convergence-only no-op rather than a fresh cancel).
89#[derive(Debug, Clone, PartialEq, Eq)]
90#[must_use]
91pub struct CancelOutcome {
92    /// True when the run's manifest was already `Cancelled` on entry, so no
93    /// `run.status: cancelled` event was appended. The call still scans and
94    /// converges any straggler nodes (an interrupted earlier cancel), so this
95    /// is a SUCCESS, not an error: "no-op: run was already cancelled,
96    /// converged N additional nodes".
97    pub run_was_already_cancelled: bool,
98    /// Nodes this cancel transaction ensured are terminally cancelled: live
99    /// nodes for which it synthesized and durably appended a terminal cancel
100    /// `node.report` (and folded it), plus any node whose cancel `node.report` a
101    /// prior interrupted cancel had already durably appended (matched by
102    /// `(kind, idempotency_key)`) and which this call converged by *re-folding*
103    /// that event rather than re-appending. Either way the node carries a
104    /// terminal cancel in the source-of-truth log and its projection is folded
105    /// (or, for a still-missing projection, will fold on rebuild — see the
106    /// module docs); none is double-reported against a node that was already
107    /// terminal on entry.
108    pub nodes_cancelled: Vec<NodeId>,
109    /// Nodes whose *status* was already terminal on entry and so were skipped —
110    /// never double-reported as freshly cancelled.
111    pub nodes_already_terminal: Vec<NodeId>,
112}
113
114/// Cancel a run in a single locked transaction. Acquires the run's
115/// [`RunLock`] once for the whole operation, then delegates to
116/// [`cancel_run_unlocked`].
117///
118/// # Errors
119///
120/// - [`Error::RunAlreadyTerminal`] if the run is `Done`/`Failed` — refused
121///   without mutating state.
122/// - I/O / corrupt-log errors from reading the manifest, listing nodes, or
123///   appending events.
124pub fn cancel_run(paths: &RunPaths, note: Option<&str>) -> Result<CancelOutcome> {
125    RunLock::with_lock(paths, |lock| cancel_run_unlocked(lock, paths, note))
126}
127
128/// The locked body of [`cancel_run`]. The `lock: &LockedRun` witness proves the
129/// caller already holds the run's exclusive [`RunLock`]; this is the sanctioned
130/// lock-held composition path so the manifest read, the per-node
131/// read-then-append loop, and the final `run.status` append all share one
132/// critical section (it calls [`append_and_apply_unlocked`], never
133/// [`crate::append_and_apply_event`], which would deadlock by re-locking).
134pub fn cancel_run_unlocked(
135    lock: &LockedRun<'_>,
136    paths: &RunPaths,
137    note: Option<&str>,
138) -> Result<CancelOutcome> {
139    let started = std::time::Instant::now();
140    let manifest = read_manifest(paths)?;
141
142    // Refuse a non-cancelled terminal run BEFORE touching any node: cancelling
143    // a Done/Failed run would synthesize node reports and append a
144    // `run.status: cancelled` the reducer's terminal-state guard then drops,
145    // so the CLI would claim a transition that never happened. An already-
146    // `Cancelled` run is not refused — it falls through to converge stragglers.
147    if manifest.status.is_terminal() && manifest.status != Status::Cancelled {
148        return Err(Error::RunAlreadyTerminal {
149            status: manifest.status,
150        });
151    }
152    let run_was_already_cancelled = manifest.status == Status::Cancelled;
153
154    // Normalize the cancel reason ONCE up front (see [`normalize_cancel_reason`]):
155    // a blank `--note` would flow in as `reason: ""`, which the reducer rejects
156    // and would brick the run's cancellability. It falls back to the default.
157    let reason = normalize_cancel_reason(note);
158
159    // One streaming replay pass over the source-of-truth log: the authoritative
160    // node set *and each node's current status* (both immune to the projection
161    // crash window), plus the prior cancel events already recorded (so a prior
162    // interrupted cancel isn't duplicated — it is re-folded instead).
163    let CancelLedger {
164        node_status,
165        prior_cancel,
166    } = read_cancel_ledger(paths)?;
167
168    let mut nodes_cancelled = Vec::new();
169    let mut nodes_already_terminal = Vec::new();
170
171    for (nid, log_status) in node_status {
172        let key = node_cancel_key(&paths.run_id, &nid);
173        // Convergence path first: this run's cancel already logged a
174        // `node.report` for this node (a prior, possibly crash-interrupted,
175        // cancel). The log is identical whether that report's projection fold
176        // landed or not, so `read_node_opt` is what tells the two apart — an
177        // already-folded terminal projection is a clean no-op reported as
178        // already-terminal, while a crash-stranded still-live projection is
179        // converged by re-folding the already-logged event (no duplicate append)
180        // and reported as cancelled. This is the only remaining projection read,
181        // and it serves convergence, not the liveness decision below.
182        if let Some(prior) = prior_cancel.get(&("node.report".to_owned(), key.clone())) {
183            if let Some(n) = read_node_opt(paths, &nid)? {
184                if n.status.is_terminal() {
185                    nodes_already_terminal.push(nid);
186                    continue;
187                }
188            }
189            apply_event(paths, prior)?;
190            nodes_cancelled.push(nid);
191            continue;
192        }
193        // No prior cancel for this node: the event log is authoritative for
194        // liveness. A node the log replays as terminal — a non-cancel terminal
195        // (`node.report success` / a `node.status` to a terminal value), or a
196        // cancel logged outside this run's key namespace — is already settled
197        // and skipped, even if a stale projection still reads live (the window
198        // cancel-liveness-from-log closes: the log wins). Only a node the log
199        // shows non-terminal (including a `node.created` whose projection write
200        // was interrupted — the crash window a `nodes/*.json` scan would drop)
201        // gets a synthesized terminal cancel report so the log records it as
202        // cancelled and a future rebuild can't resurrect it as live.
203        if log_status.is_terminal() {
204            nodes_already_terminal.push(nid);
205            continue;
206        }
207        let data = json!({
208            "success": false,
209            "cancelled": true,
210            "reason": reason,
211            "summary": "Run cancelled before agent reported.",
212            "discussion_items": [],
213            "spinoff_proposals": [],
214            "wrap_up_recommendations": []
215        });
216        append_and_apply_unlocked(lock, paths, "node.report", Some(&nid), Some(&key), data)?;
217        nodes_cancelled.push(nid);
218    }
219
220    if !run_was_already_cancelled {
221        let key = run_status_cancel_key(&paths.run_id);
222        if let Some(prior) = prior_cancel.get(&("run.status".to_owned(), key.clone())) {
223            // A prior interrupted cancel already logged the terminal `run.status`
224            // (fsynced before its manifest fold). Re-fold it to converge the
225            // manifest instead of appending a duplicate `run.status: cancelled`.
226            apply_event(paths, prior)?;
227        } else {
228            let mut status_data = serde_json::Map::new();
229            status_data.insert("status".into(), "cancelled".into());
230            // Record the operator note only when one was actually supplied (the
231            // trimmed, non-blank value); a blank `--note` leaves the field unset
232            // rather than writing an empty string.
233            if let Some(n) = note.map(str::trim).filter(|s| !s.is_empty()) {
234                status_data.insert("note".into(), n.into());
235            }
236            append_and_apply_unlocked(
237                lock,
238                paths,
239                "run.status",
240                None,
241                Some(&key),
242                serde_json::Value::Object(status_data),
243            )?;
244        }
245    }
246
247    tracing::debug!(
248        target: "taskfleet_core::cancel",
249        run_id = %paths.run_id,
250        held_ms = started.elapsed().as_millis() as u64,
251        nodes_cancelled = nodes_cancelled.len(),
252        nodes_already_terminal = nodes_already_terminal.len(),
253        "cancel transaction complete",
254    );
255
256    Ok(CancelOutcome {
257        run_was_already_cancelled,
258        nodes_cancelled,
259        nodes_already_terminal,
260    })
261}
262
263/// Outcome of a [`cancel_node`] transaction — a single-node, branch-preserving
264/// cancel for one live fan-out child.
265///
266/// Per-node cancel deliberately leaves the run non-terminal *while any sibling is
267/// still live* (design §2.5 — a stuck child is unblocked without killing the
268/// batch). But when this cancel settles the **last** live node, the run is rolled
269/// up **in the same locked transaction** ([`rolled_up`](Self::rolled_up)) rather
270/// than deferred to the supervisor — so a run whose supervisor has died is never
271/// stranded non-terminal (llm-review C1).
272#[derive(Debug, Clone, PartialEq, Eq)]
273#[must_use]
274pub struct NodeCancelOutcome {
275    /// The node this call targeted (fully resolved).
276    pub node_id: NodeId,
277    /// True when this call ensured the node carries a terminal cancel in the
278    /// source-of-truth log — either by synthesizing and durably appending a fresh
279    /// cancel `node.report`, or by re-folding a prior interrupted cancel's
280    /// already-logged event (crash convergence, no duplicate append). False when
281    /// the node was already terminal on entry (see `already_terminal`).
282    pub cancelled: bool,
283    /// True when the node was *already* terminal on entry (merged, failed, or a
284    /// prior cancel already folded) — a clean idempotent no-op, never a fresh
285    /// cancel. Mutually exclusive with `cancelled`.
286    pub already_terminal: bool,
287    /// `Some(status)` when this cancel settled the last live node and therefore
288    /// rolled the whole run up to a terminal status **under the same lock**
289    /// (`Cancelled` when no sibling failed, `Failed` when one did). `None` when
290    /// siblings remain live (the run stays live) or the run was already terminal
291    /// on entry. Lets the CLI tell the operator whether the run itself is now
292    /// settled rather than implying a rollup that might never come.
293    pub rolled_up: Option<Status>,
294}
295
296/// Cancel exactly ONE live node of a run, preserving its branch + worktree.
297/// Acquires the run's [`RunLock`] once and delegates to
298/// [`cancel_node_unlocked`].
299///
300/// This is the fan-out selectivity primitive (design §2.5, issue
301/// `per-node-run`): where [`cancel_run`] settles every live node and rolls the
302/// run up to `Cancelled` in one shot, this settles a single named node. While
303/// any sibling is still live it appends **only** that node's terminal cancel
304/// `node.report` (no `run.status`) — the run stays live so the batch keeps
305/// running. When it settles the **last** live node it also rolls the run up to a
306/// terminal status **in the same locked transaction** (so a dead supervisor can
307/// never strand the run non-terminal — llm-review C1). The synthesized terminal
308/// cancel `node.report` classifies as [`Cancelled`](crate::Status) →
309/// `Teardown::SourceRelative`, so invariant 5 preserves the node's committed work
310/// rather than force-deleting it.
311///
312/// # Errors
313///
314/// - [`Error::RunAlreadyTerminal`] if the run is `Done`/`Failed` — refused
315///   without mutating state (mirrors [`cancel_run`]; an already-`Cancelled` run
316///   is *not* refused — its live nodes can still be settled).
317/// - [`Error::NodeNotFound`] if `node_id` names no node in the run's log.
318/// - I/O / corrupt-log errors from reading the manifest, replaying the log, or
319///   appending the report.
320pub fn cancel_node(
321    paths: &RunPaths,
322    node_id: &NodeId,
323    note: Option<&str>,
324) -> Result<NodeCancelOutcome> {
325    RunLock::with_lock(paths, |lock| {
326        cancel_node_unlocked(lock, paths, node_id, note)
327    })
328}
329
330/// The locked body of [`cancel_node`]. The `lock: &LockedRun` witness proves the
331/// caller already holds the run's exclusive [`RunLock`], so the manifest read,
332/// the log replay, the convergence read, the single report append, AND the
333/// optional last-node roll-up all share one critical section (it calls
334/// [`append_and_apply_unlocked`], never [`crate::append_and_apply_event`], which
335/// would deadlock by re-locking).
336pub fn cancel_node_unlocked(
337    lock: &LockedRun<'_>,
338    paths: &RunPaths,
339    node_id: &NodeId,
340    note: Option<&str>,
341) -> Result<NodeCancelOutcome> {
342    // Refuse a non-cancelled terminal run up front, mirroring `cancel_run`: a
343    // `Done`/`Failed` run's nodes are all terminal, so appending a fresh cancel
344    // `node.report` would either bloat the log with a dead post-terminal event or
345    // (on rebuild) flip a settled node — divergence. An already-`Cancelled` run
346    // is NOT refused: it may still carry a straggler live node to settle
347    // (llm-review C4).
348    let manifest = read_manifest(paths)?;
349    if manifest.status.is_terminal() && manifest.status != Status::Cancelled {
350        return Err(Error::RunAlreadyTerminal {
351            status: manifest.status,
352        });
353    }
354
355    // One streaming replay pass over the source-of-truth log gives the
356    // authoritative node set *and* each node's log-derived status (both immune to
357    // the projection crash window), plus this run's already-logged cancel events
358    // so a prior interrupted cancel is re-folded, never duplicated.
359    let CancelLedger {
360        node_status,
361        prior_cancel,
362    } = read_cancel_ledger(paths)?;
363
364    // The log is authoritative for the node set: a node whose `node.created` was
365    // fsynced but whose projection write was crash-interrupted is still
366    // resolvable here (a `nodes/*.json` scan would miss it). A genuinely absent
367    // id is a caller error.
368    let log_status = node_status
369        .iter()
370        .find(|(nid, _)| nid == node_id)
371        .map(|(_, s)| *s)
372        .ok_or_else(|| Error::NodeNotFound {
373            node_id: node_id.as_str().to_owned(),
374        })?;
375
376    let key = node_cancel_key(&paths.run_id, node_id);
377
378    // Settle the target node. `cancelled` = this call ensured a terminal cancel
379    // in the log (fresh append or crash-convergence re-fold); `already_terminal`
380    // = the node was already terminal on entry (idempotent no-op).
381    let (cancelled, already_terminal) =
382        if let Some(prior) = prior_cancel.get(&("node.report".to_owned(), key.clone())) {
383            // Convergence path: this run's cancel already logged a `node.report`
384            // for this node (a prior, possibly crash-interrupted, per-node or
385            // whole-run cancel). The log is identical whether that report's
386            // projection fold landed or not, so `read_node_opt` tells the two
387            // apart — an already-folded terminal projection is a clean no-op,
388            // while a crash-stranded still-live projection is converged by
389            // re-folding the already-logged event (no duplicate append).
390            let folded_terminal =
391                read_node_opt(paths, node_id)?.is_some_and(|n| n.status.is_terminal());
392            if folded_terminal {
393                (false, true)
394            } else {
395                apply_event(paths, prior)?;
396                (true, false)
397            }
398        } else if log_status.is_terminal() {
399            // No prior cancel: the log is authoritative for liveness. A node the
400            // log replays as terminal — a natural success/failure, or a cancel
401            // logged outside this run's key namespace — is already settled, even
402            // if a stale projection still reads live (the log wins).
403            (false, true)
404        } else {
405            let reason = normalize_cancel_reason(note);
406            let data = json!({
407                "success": false,
408                "cancelled": true,
409                "reason": reason,
410                "summary": "Node cancelled before agent reported.",
411                "discussion_items": [],
412                "spinoff_proposals": [],
413                "wrap_up_recommendations": []
414            });
415            append_and_apply_unlocked(lock, paths, "node.report", Some(node_id), Some(&key), data)?;
416            (true, false)
417        };
418
419    // Last-node roll-up (llm-review C1): if the run is still non-terminal but
420    // every node is now terminal, terminalize the run HERE, under the same lock,
421    // rather than deferring to the supervisor — which may be dead, leaving the
422    // run stranded `pending` with no live node. The aggregate is derived from the
423    // log-authoritative `node_status` (with the target's post-cancel status
424    // overridden), so it never mis-terminalizes over a stale projection. The
425    // decision runs only when NO sibling is live, so the "don't terminalize while
426    // a sibling runs" invariant holds.
427    let rolled_up = maybe_roll_up_run(
428        lock,
429        paths,
430        manifest.status,
431        &node_status,
432        node_id,
433        cancelled,
434        log_status,
435        &prior_cancel,
436    )?;
437
438    Ok(NodeCancelOutcome {
439        node_id: node_id.clone(),
440        cancelled,
441        already_terminal,
442        rolled_up,
443    })
444}
445
446/// Roll the run up to a terminal status when this per-node cancel settled the
447/// last live node. Returns the status appended, or `None` when the run stays
448/// live (a sibling is still live) or was already terminal on entry.
449///
450/// Shares the whole-run cancel's `run-cancel:<run>:run-status` idempotency key
451/// (via [`run_status_cancel_key`]) so the three run-status producers in the
452/// cancel family — whole-run cancel, this last-node roll-up, and a crash-retry of
453/// either — converge on ONE logical `run.status` rather than duplicating it: a
454/// prior interrupted append captured in `prior_cancel` is re-folded, and a later
455/// `cancel_run` finds the same key and skips (llm-review C2/C4). The supervisor's
456/// own `supervise::cleanup::rollup_status` uses a different key, but it fires only
457/// while the manifest is non-terminal, so once this append lands the supervisor's
458/// tick reads the run terminal and no-ops.
459#[allow(clippy::too_many_arguments)]
460fn maybe_roll_up_run(
461    lock: &LockedRun<'_>,
462    paths: &RunPaths,
463    run_status_on_entry: Status,
464    node_status: &[(NodeId, Status)],
465    target: &NodeId,
466    target_cancelled: bool,
467    target_log_status: Status,
468    prior_cancel: &HashMap<(String, String), Event>,
469) -> Result<Option<Status>> {
470    if run_status_on_entry.is_terminal() {
471        // The run was already terminal on entry (an already-`Cancelled` run whose
472        // straggler we just settled) — nothing to roll up.
473        return Ok(None);
474    }
475    // The target's effective post-cancel status: `Cancelled` if this call settled
476    // it, else its log-derived status (an already-terminal node).
477    let effective_target = if target_cancelled {
478        Status::Cancelled
479    } else {
480        target_log_status
481    };
482    let statuses = node_status
483        .iter()
484        .map(|(nid, s)| if nid == target { effective_target } else { *s });
485    let Some(agg) = crate::aggregate_terminal_status(statuses) else {
486        // A sibling is still live — the run stays live, as designed.
487        return Ok(None);
488    };
489
490    let key = run_status_cancel_key(&paths.run_id);
491    if let Some(prior) = prior_cancel.get(&("run.status".to_owned(), key.clone())) {
492        // A prior interrupted cancel already logged the terminal `run.status`
493        // (fsynced before its manifest fold). Re-fold it to converge the manifest
494        // instead of appending a duplicate.
495        apply_event(paths, prior)?;
496    } else {
497        // `Status` serializes kebab-case (`cancelled`/`failed`/`done`) — the
498        // exact shape the reducer's `run.status` handler parses.
499        append_and_apply_unlocked(
500            lock,
501            paths,
502            "run.status",
503            None,
504            Some(&key),
505            json!({ "status": agg }),
506        )?;
507    }
508    Ok(Some(agg))
509}
510
511/// Normalize a `--note` into the terminal cancel report's `reason`. An empty or
512/// whitespace-only note would otherwise flow in as `reason: ""`, which the
513/// reducer rejects (`CancelledRequiresReason`) — aborting the transaction and,
514/// since a retry reuses the same bad note, leaving the node/run permanently
515/// un-cancellable. A blank note falls back to the default.
516fn normalize_cancel_reason(note: Option<&str>) -> &str {
517    note.map(str::trim)
518        .filter(|s| !s.is_empty())
519        .unwrap_or("cancelled by user")
520}
521
522/// Cancel-relevant facts replayed from `events.jsonl` in one streaming pass
523/// under the held lock.
524struct CancelLedger {
525    /// Every node a `node.created` event introduced, paired with the status the
526    /// log replays for it, deduped and sorted by numeric suffix. The
527    /// authoritative live-node set *and* per-node liveness: replayed from the
528    /// source of truth, so it includes a node whose projection write was
529    /// crash-interrupted (the node a `nodes/*.json` scan would miss) and reports
530    /// a node terminal whenever the log says so even if the projection still
531    /// reads live (the window [`crate::events`] documents the log leading the
532    /// projections through).
533    node_status: Vec<(NodeId, Status)>,
534    /// Cancel events this run already logged, keyed by `(kind, idempotency_key)`
535    /// and limited to this run's `run-cancel:<run_id>:` key namespace (first
536    /// occurrence wins, mirroring [`crate::events::find_prior_with_key`]). The
537    /// cancel loop looks an entry up by its deterministic key to (a) avoid
538    /// re-appending a duplicate and (b) re-fold the event so a crash-stranded
539    /// projection converges. Keying by `(kind, key)` — not the bare string —
540    /// keeps a coincidental or forged key on an unrelated `kind` from masking a
541    /// real cancel append. Only these few lines have their full [`Event`] payload
542    /// materialized; every other line is skimmed envelope-only.
543    prior_cancel: HashMap<(String, String), Event>,
544}
545
546/// Envelope + the few small `data` fields the cancel ledger needs from each
547/// line, skimmed by [`for_each_event_probe`] without materializing the
548/// (potentially multi-KB) full `data` payload. serde ignores every other field,
549/// so a rich `node.report` is scanned but never allocated.
550#[derive(Deserialize)]
551struct EventSeqProbe {
552    seq: u64,
553}
554
555#[derive(Deserialize)]
556struct CancelProbe {
557    seq: u64,
558    kind: String,
559    #[serde(default)]
560    node_id: Option<NodeId>,
561    #[serde(default)]
562    idempotency_key: Option<String>,
563    #[serde(default)]
564    data: CancelProbeData,
565}
566
567/// The status-bearing `data` fields of `node.status` / `node.report`. All
568/// optional: any other event kind simply leaves them `None`.
569#[derive(Deserialize, Default)]
570struct CancelProbeData {
571    // Keep these as raw JSON values. Besides matching the reducer's strict
572    // interpretation exactly (for example, an advisory non-string `via` does
573    // not invalidate a typed RunMerge origin), this lets a bounded replay skip
574    // future malformed report fields without deserializing them into a typed
575    // shape that could reject an earlier recovery event.
576    #[serde(default)]
577    status: Option<Value>,
578    #[serde(default)]
579    success: Option<Value>,
580    #[serde(default)]
581    cancelled: Option<Value>,
582    #[serde(default)]
583    via: Option<Value>,
584    #[serde(default)]
585    origin: OriginProbe,
586}
587
588#[derive(Default)]
589struct OriginProbe {
590    present: bool,
591    value: Value,
592}
593
594impl<'de> Deserialize<'de> for OriginProbe {
595    fn deserialize<D>(deserializer: D) -> std::result::Result<Self, D::Error>
596    where
597        D: serde::Deserializer<'de>,
598    {
599        Ok(Self {
600            present: true,
601            value: Value::deserialize(deserializer)?,
602        })
603    }
604}
605
606impl CancelProbeData {
607    fn report_value(&self) -> Value {
608        let mut report = serde_json::Map::new();
609        if let Some(success) = &self.success {
610            report.insert("success".into(), success.clone());
611        }
612        if let Some(cancelled) = &self.cancelled {
613            report.insert("cancelled".into(), cancelled.clone());
614        }
615        if let Some(via) = &self.via {
616            report.insert("via".into(), via.clone());
617        }
618        if self.origin.present {
619            report.insert("origin".into(), self.origin.value.clone());
620        }
621        Value::Object(report)
622    }
623}
624
625/// Replay `events.jsonl` once, streaming, to build the [`CancelLedger`].
626///
627/// Reads through [`RunPaths::checked_events`] so a symlinked event log is
628/// refused, matching the mutation path — the cancel decision must not be made
629/// from content redirected outside the run tree. Uses [`for_each_event_probe`],
630/// which shares the crate's torn-tail policy: a crash-truncated final line is
631/// dropped as an uncommitted partial write, while any *interior* unparseable
632/// line is surfaced as [`Error::CorruptEventLog`] — so a corrupt log fails the
633/// cancel loudly rather than silently dropping a node. A missing log yields an
634/// empty ledger (run never appended an event — nothing to cancel).
635///
636/// Per-node status is accumulated by the shared [`NodeStatusAcc`] state machine
637/// (which mirrors the reducer's terminal-state guard) — the same accumulator the
638/// supervisor's log-authoritative [`read_node_statuses`] uses, so the cancel path
639/// and the supervisor roll-up can never derive a different node set or status
640/// from the same log. Node ids come out sorted by numeric suffix.
641fn read_cancel_ledger(paths: &RunPaths) -> Result<CancelLedger> {
642    let events_path = paths.checked_events()?;
643    let prefix = format!("run-cancel:{}:", paths.run_id.as_str());
644    let mut acc = NodeStatusAcc::default();
645    let mut prior_cancel: HashMap<(String, String), Event> = HashMap::new();
646
647    for_each_event_probe::<CancelProbe, _>(&events_path, |probe, raw| {
648        acc.observe(&probe);
649        // Capture only this run's cancel events, keyed by (kind, key), and only
650        // for those materialize the full payload the re-fold path needs.
651        if let Some(key) = probe
652            .idempotency_key
653            .as_deref()
654            .filter(|k| k.starts_with(&prefix))
655        {
656            let entry = (probe.kind.clone(), key.to_owned());
657            if let std::collections::hash_map::Entry::Vacant(slot) = prior_cancel.entry(entry) {
658                let ev: Event =
659                    serde_json::from_slice(raw).map_err(|e| Error::CorruptEventLog {
660                        path: events_path.clone(),
661                        reason: format!(
662                            "cancel ledger: line matched a run-cancel key but is not a \
663                         replayable event: {} [{e}]",
664                            excerpt(raw)
665                        ),
666                    })?;
667                slot.insert(ev);
668            }
669        }
670        Ok(())
671    })?;
672
673    Ok(CancelLedger {
674        node_status: acc.finish(),
675        prior_cancel,
676    })
677}
678
679/// The log-authoritative per-node status set for a run, replayed once from
680/// `events.jsonl` — the source-of-truth alternative to a `nodes/*.json`
681/// projection scan.
682///
683/// Returns every node a `node.created` event introduced, paired with the status
684/// the log replays for it, deduped and sorted by numeric suffix ([`NodeId`]
685/// order). Both the node set and each node's status come from the log, so the
686/// result includes a node whose `node.created` was fsynced while its projection
687/// write was crash-interrupted (the node a `nodes/*.json` scan would silently
688/// drop) and reports a node terminal whenever the log says so even if its
689/// projection still reads live (the window [`crate::events`] documents the log
690/// leading the projections through).
691///
692/// This is what lets a supervisor's run-status roll-up stay log-authoritative:
693/// terminalizing a run from the projection subset can miss a log-visible live
694/// node and roll the run up while it is still running, which a later
695/// `rebuild_projections` would then resurrect as live under a terminal run
696/// (violating "a run must not terminalize while a log-visible node is live" —
697/// issue `rollup-status-log-authoritative`). Feeding this into
698/// [`aggregate_terminal_status`](crate::aggregate_terminal_status) closes that
699/// window. It is the read half [`cancel_node`]'s in-lock self-roll-up already
700/// uses via the cancel ledger; both now share the `NodeStatusAcc` state machine
701/// so the supervisor tick and the cancel path can never diverge.
702///
703/// Reads through `RunPaths::checked_events` (a symlinked log is refused) and
704/// shares the crate's streaming torn-tail policy: a crash-truncated final line
705/// is dropped, an interior unparseable line surfaces as
706/// [`Error::CorruptEventLog`], and a missing log yields an empty set.
707///
708/// **Cost.** One streaming pass over the whole log — `O(total events)`, not
709/// `O(nodes)`. Memory stays bounded (each line is skimmed envelope + a few small
710/// status fields, never the full `node.report` payload — a run with hundreds of
711/// nodes and multi-KB reports is scanned without ever holding a report in
712/// memory), but the *work* is linear in the event count, not the node count. A
713/// caller that polls this every tick (the supervisor roll-up) re-reads from byte
714/// 0 each time; at the tool's scale (tens of nodes, hundreds of events,
715/// multi-second ticks) that is negligible, but an incremental fold that resumes
716/// from the last consumed offset would be the optimization if a run's log ever
717/// grows large enough to matter.
718///
719/// # Errors
720///
721/// I/O errors reading the log, a rejected symlinked path, or an interior corrupt
722/// event line.
723pub fn read_node_statuses(paths: &RunPaths) -> Result<Vec<(NodeId, Status)>> {
724    Ok(read_node_status_facts(paths, None)?
725        .into_iter()
726        .map(|fact| (fact.node_id, fact.status))
727        .collect())
728}
729
730/// One node's log-derived terminal facts. `confirmed_merge_seq` is present only
731/// when the shared terminal-recovery predicate actually adopted that report;
732/// seeing an authoritative-looking report elsewhere in the log is insufficient.
733#[derive(Debug, Clone, PartialEq, Eq)]
734pub struct NodeStatusFact {
735    /// Node whose facts were folded.
736    pub node_id: NodeId,
737    /// Current status after applying the reducer-equivalent transition rules.
738    pub status: Status,
739    /// Sequence of the authoritative merge report adopted for this node.
740    pub confirmed_merge_seq: Option<u64>,
741}
742
743/// Replay node statuses and adopted merge authority, optionally considering
744/// only events strictly before `before_seq`. The bound is load-bearing for
745/// reducer replay: a future merge report must never authorize an earlier
746/// `run.status` event.
747pub fn read_node_status_facts(
748    paths: &RunPaths,
749    before_seq: Option<u64>,
750) -> Result<Vec<NodeStatusFact>> {
751    let events_path = paths.checked_events()?;
752    let mut acc = NodeStatusAcc::default();
753    for_each_event_probe::<EventSeqProbe, _>(&events_path, |seq_probe, raw| {
754        if before_seq.is_none_or(|bound| seq_probe.seq < bound) {
755            let probe: CancelProbe =
756                serde_json::from_slice(raw).map_err(|e| Error::CorruptEventLog {
757                    path: events_path.clone(),
758                    reason: format!(
759                        "node-status fold: event before bound is malformed: {} [{e}]",
760                        excerpt(raw)
761                    ),
762                })?;
763            acc.observe(&probe);
764        }
765        Ok(())
766    })?;
767    Ok(acc.finish_facts())
768}
769
770/// Streaming accumulator for log-derived per-node status: the shared state
771/// machine behind both the cancel ledger ([`read_cancel_ledger`]) and the
772/// supervisor's log-authoritative roll-up ([`read_node_statuses`]), so the two
773/// can never derive a different node set or status from the same log.
774///
775/// Mirrors the reducer's terminal-state guard exactly: `node.created` seeds
776/// [`Status::Pending`] (idempotent on replay — a second `node.created` for the
777/// same id is a no-op, matching the reducer's existence guard), and
778/// `node.status` / `node.report` transition a node only while it is still
779/// non-terminal. A `node.status` / `node.report` for an id never introduced by a
780/// `node.created` is ignored, exactly as the reducer no-ops a status/report
781/// against a non-existent node. Malformed status fields degrade gracefully (the
782/// node is left at its current status — for the cancel path that keeps the node
783/// non-terminal, hence cancellable) rather than aborting — the append path
784/// already validates every committed event, so a `None` transition only arises
785/// from a hand-corrupted log.
786///
787/// **`node.retry` is deliberately not folded, and that is exact.** In the reducer
788/// (`reduce_node_retry`) a retry against an already-terminal node is a no-op (a
789/// settled node is frozen) and a retry against a *live* node only rewires it back
790/// to `Pending`. So `node.retry` never crosses the terminal/live boundary in
791/// either direction: a node that is terminal here is terminal in the reducer, and
792/// one that is live here (whatever its exact non-terminal value) is live there.
793/// Since the only consumers ([`read_node_statuses`] → `aggregate_terminal_status`,
794/// and the cancel loop) classify solely on terminal-vs-live, ignoring
795/// `node.retry` derives the same answer the reducer would — it is not a missed
796/// transition.
797#[derive(Default)]
798struct NodeStatusAcc {
799    /// Node ids in creation order (re-sorted by numeric suffix at [`finish`]).
800    order: Vec<NodeId>,
801    /// Each node's current log-derived status.
802    status: HashMap<NodeId, Status>,
803    /// Sequence of the confirmed merge report that the shared recovery
804    /// predicate adopted for this node.
805    confirmed_merge_seq: HashMap<NodeId, u64>,
806}
807
808impl NodeStatusAcc {
809    /// Fold one skimmed event line ([`CancelProbe`]) into the per-node map.
810    fn observe(&mut self, probe: &CancelProbe) {
811        match probe.kind.as_str() {
812            "node.created" => {
813                if let Some(nid) = &probe.node_id {
814                    if !self.status.contains_key(nid) {
815                        self.order.push(nid.clone());
816                        self.status.insert(nid.clone(), Status::Pending);
817                    }
818                }
819            }
820            "node.status" => {
821                if let Some(nid) = &probe.node_id {
822                    if let Some(cur) = self.status.get_mut(nid) {
823                        if !cur.is_terminal() {
824                            if let Some(ns) = probe
825                                .data
826                                .status
827                                .as_ref()
828                                .and_then(Value::as_str)
829                                .and_then(parse_status)
830                            {
831                                *cur = ns;
832                            }
833                        }
834                    }
835                }
836            }
837            "node.report" => {
838                if let Some(nid) = &probe.node_id {
839                    if let Some(cur) = self.status.get_mut(nid) {
840                        let report = probe.data.report_value();
841                        if cur.is_terminal() {
842                            // Keep this exactly aligned with `reduce_node_report`:
843                            // only an authoritative successful merge may repair
844                            // Failed/Done, and cancellation is immutable.
845                            if ReportOrigin::permits_terminal_merge_recovery(*cur, &report) {
846                                *cur = Status::Done;
847                                self.confirmed_merge_seq.insert(nid.clone(), probe.seq);
848                            }
849                        } else if let Some(ns) = report_terminal_status(
850                            probe.data.success.as_ref().and_then(Value::as_bool),
851                            probe.data.cancelled.as_ref().and_then(Value::as_bool),
852                        ) {
853                            *cur = ns;
854                            if ns == Status::Done
855                                && ReportOrigin::report_is_confirmed_merge(&report)
856                            {
857                                self.confirmed_merge_seq.insert(nid.clone(), probe.seq);
858                            }
859                        }
860                    }
861                }
862            }
863            _ => {}
864        }
865    }
866
867    /// Consume into the sorted `(NodeId, Status)` list. Node ids are sorted by
868    /// numeric suffix (not lexically), so a run past the digit-width boundary
869    /// where `n-10000` would otherwise sort before `n-9999` stays intuitive (see
870    /// [`NodeId`]). A validated `NodeId` is `n-` + ASCII digits (≤10, so it fits
871    /// in u64); the `unwrap_or` keeps the sort total for a hypothetical
872    /// unparseable body.
873    fn finish(self) -> Vec<(NodeId, Status)> {
874        self.finish_facts()
875            .into_iter()
876            .map(|fact| (fact.node_id, fact.status))
877            .collect()
878    }
879
880    fn finish_facts(mut self) -> Vec<NodeStatusFact> {
881        self.order.sort_by_key(|id| {
882            id.as_str()
883                .strip_prefix("n-")
884                .and_then(|d| d.parse::<u64>().ok())
885                .unwrap_or(0)
886        });
887        self.order
888            .into_iter()
889            .map(|id| NodeStatusFact {
890                status: self.status[&id],
891                confirmed_merge_seq: self.confirmed_merge_seq.get(&id).copied(),
892                node_id: id,
893            })
894            .collect()
895    }
896}
897
898/// Parse a `node.status` / `run.status` status string into a [`Status`],
899/// returning `None` for an unrecognized value (treated as "no transition" so a
900/// corrupt status never aborts the cancel). Goes through serde so the kebab-case
901/// mapping can never drift from the [`Status`] enum.
902fn parse_status(s: &str) -> Option<Status> {
903    serde_json::from_value(Value::String(s.to_owned())).ok()
904}
905
906/// Derive the terminal status a `node.report` asserts from its `success` /
907/// `cancelled` flags, mirroring the reducer's success-XOR-cancelled rule but
908/// *lenient*: a bare/contradictory report yields `None` (no transition — the
909/// node stays live and is cancelled) rather than the reducer's
910/// [`Error::CorruptEventLog`]. The append path rejects such a report before it
911/// is ever committed against a live node, so a `None` here only arises from a
912/// hand-corrupted log, where leaving the node cancellable is the safe default.
913fn report_terminal_status(success: Option<bool>, cancelled: Option<bool>) -> Option<Status> {
914    if cancelled.unwrap_or(false) {
915        // `cancelled: true` with `success: true` is contradictory → no transition.
916        if success == Some(true) {
917            return None;
918        }
919        Some(Status::Cancelled)
920    } else {
921        match success {
922            Some(true) => Some(Status::Done),
923            Some(false) => Some(Status::Failed),
924            None => None,
925        }
926    }
927}
928
929/// Deterministic idempotency key for the synthesized cancel `node.report` of
930/// one node. Stable in `(run_id, node_id)` so a re-`cancel` after a crash that
931/// fsynced the report but never folded its projection finds the prior event and
932/// does not append a duplicate logical-cancel.
933fn node_cancel_key(run_id: &RunId, node_id: &NodeId) -> String {
934    format!("run-cancel:{}:node:{}", run_id.as_str(), node_id.as_str())
935}
936
937/// Deterministic idempotency key for the run's terminal `run.status: cancelled`
938/// event. Stable in `run_id` for the same crash-retry reason as
939/// [`node_cancel_key`].
940fn run_status_cancel_key(run_id: &RunId) -> String {
941    format!("run-cancel:{}:run-status", run_id.as_str())
942}
943
944#[cfg(test)]
945mod tests {
946    use super::*;
947    use crate::events::{append_and_apply_event, append_event_with_seq, read_all_events};
948    use crate::lock::ACQUIRE_COUNT;
949    use tempfile::TempDir;
950
951    /// Count `node.report` events recorded in the log for one node id.
952    fn report_count(paths: &RunPaths, nid: &str) -> usize {
953        read_all_events(&paths.events())
954            .unwrap()
955            .iter()
956            .filter(|e| {
957                e.kind == "node.report" && e.node_id.as_ref().map(NodeId::as_str) == Some(nid)
958            })
959            .count()
960    }
961
962    fn fresh_run(tmp: &TempDir) -> RunPaths {
963        let run_id = "01jxsnap000000000000000000";
964        let dir = tmp.path().join(run_id);
965        std::fs::create_dir_all(&dir).unwrap();
966        RunPaths::new(dir, run_id).unwrap()
967    }
968
969    /// Parse a `NodeId` for a test append call (the typed envelope id).
970    fn nid(s: &str) -> NodeId {
971        NodeId::parse_str(s).unwrap()
972    }
973
974    /// Drive a run to `count` live nodes (n-0001..) under a created manifest.
975    fn bootstrap(paths: &RunPaths, count: usize) {
976        append_and_apply_event(
977            paths,
978            "run.created",
979            None,
980            None,
981            json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
982        )
983        .unwrap();
984        for i in 1..=count {
985            let node_id = nid(&format!("n-{i:04}"));
986            append_and_apply_event(
987                paths,
988                "node.created",
989                Some(&node_id),
990                None,
991                json!({ "kind": "spinoff" }),
992            )
993            .unwrap();
994        }
995    }
996
997    fn node_status(paths: &RunPaths, nid: &str) -> Status {
998        let id = NodeId::parse_str(nid).unwrap();
999        crate::read_node(paths, &id).unwrap().status
1000    }
1001
1002    #[test]
1003    fn cancel_running_run_converges_live_nodes_and_settles_run() {
1004        let tmp = TempDir::new().unwrap();
1005        let paths = fresh_run(&tmp);
1006        bootstrap(&paths, 2);
1007
1008        let out = cancel_run(&paths, Some("stop")).unwrap();
1009        assert!(!out.run_was_already_cancelled);
1010        assert_eq!(
1011            out.nodes_cancelled
1012                .iter()
1013                .map(NodeId::as_str)
1014                .collect::<Vec<_>>(),
1015            vec!["n-0001", "n-0002"]
1016        );
1017        assert_eq!(out.nodes_already_terminal.len(), 0);
1018        assert_eq!(node_status(&paths, "n-0001"), Status::Cancelled);
1019        assert_eq!(
1020            crate::read_manifest(&paths).unwrap().status,
1021            Status::Cancelled
1022        );
1023    }
1024
1025    #[test]
1026    fn cancel_done_run_is_refused_without_mutation() {
1027        let tmp = TempDir::new().unwrap();
1028        let paths = fresh_run(&tmp);
1029        bootstrap(&paths, 1);
1030        // Settle the single node, then the run, to Done.
1031        append_and_apply_event(
1032            &paths,
1033            "node.report",
1034            Some(&nid("n-0001")),
1035            None,
1036            json!({ "success": true }),
1037        )
1038        .unwrap();
1039        append_and_apply_event(
1040            &paths,
1041            "run.status",
1042            None,
1043            None,
1044            json!({ "status": "done" }),
1045        )
1046        .unwrap();
1047        let before = read_all_events(&paths.events()).unwrap().len();
1048
1049        let err = cancel_run(&paths, None).unwrap_err();
1050        assert!(
1051            matches!(
1052                err,
1053                Error::RunAlreadyTerminal {
1054                    status: Status::Done
1055                }
1056            ),
1057            "got {err:?}"
1058        );
1059        assert_eq!(
1060            read_all_events(&paths.events()).unwrap().len(),
1061            before,
1062            "a refused cancel must not append any event"
1063        );
1064        assert_eq!(crate::read_manifest(&paths).unwrap().status, Status::Done);
1065    }
1066
1067    #[test]
1068    fn recancel_cancelled_run_converges_straggler_node() {
1069        let tmp = TempDir::new().unwrap();
1070        let paths = fresh_run(&tmp);
1071        bootstrap(&paths, 2);
1072        // Simulate an interrupted cancel: run is Cancelled, but n-0002 is still
1073        // live (its node.report never landed).
1074        append_and_apply_event(
1075            &paths,
1076            "node.report",
1077            Some(&nid("n-0001")),
1078            None,
1079            json!({ "success": false, "cancelled": true, "reason": "x" }),
1080        )
1081        .unwrap();
1082        append_and_apply_event(
1083            &paths,
1084            "run.status",
1085            None,
1086            None,
1087            json!({ "status": "cancelled" }),
1088        )
1089        .unwrap();
1090        assert_eq!(node_status(&paths, "n-0002"), Status::Pending);
1091
1092        let out = cancel_run(&paths, None).unwrap();
1093        assert!(out.run_was_already_cancelled);
1094        assert_eq!(
1095            out.nodes_cancelled
1096                .iter()
1097                .map(NodeId::as_str)
1098                .collect::<Vec<_>>(),
1099            vec!["n-0002"],
1100            "only the straggler converges"
1101        );
1102        assert_eq!(
1103            out.nodes_already_terminal
1104                .iter()
1105                .map(NodeId::as_str)
1106                .collect::<Vec<_>>(),
1107            vec!["n-0001"]
1108        );
1109        assert_eq!(node_status(&paths, "n-0002"), Status::Cancelled);
1110    }
1111
1112    #[test]
1113    fn recancel_fully_converged_run_is_a_clean_noop() {
1114        let tmp = TempDir::new().unwrap();
1115        let paths = fresh_run(&tmp);
1116        bootstrap(&paths, 1);
1117        let _ = cancel_run(&paths, None).unwrap(); // first cancel converges everything
1118        let before = read_all_events(&paths.events()).unwrap().len();
1119
1120        let out = cancel_run(&paths, None).unwrap();
1121        assert!(out.run_was_already_cancelled);
1122        assert!(out.nodes_cancelled.is_empty(), "nothing left to converge");
1123        assert_eq!(
1124            out.nodes_already_terminal
1125                .iter()
1126                .map(NodeId::as_str)
1127                .collect::<Vec<_>>(),
1128            vec!["n-0001"]
1129        );
1130        assert_eq!(
1131            read_all_events(&paths.events()).unwrap().len(),
1132            before,
1133            "a fully-converged re-cancel appends nothing"
1134        );
1135    }
1136
1137    #[test]
1138    fn already_terminal_node_is_not_over_reported() {
1139        // The honesty guard: a node already settled (terminal) on entry is
1140        // reported under `nodes_already_terminal`, never `nodes_cancelled`,
1141        // even though it sits in nodes/ alongside a live node.
1142        let tmp = TempDir::new().unwrap();
1143        let paths = fresh_run(&tmp);
1144        bootstrap(&paths, 2);
1145        // n-0001 finishes on its own (Done) before the cancel.
1146        append_and_apply_event(
1147            &paths,
1148            "node.report",
1149            Some(&nid("n-0001")),
1150            None,
1151            json!({ "success": true }),
1152        )
1153        .unwrap();
1154
1155        let out = cancel_run(&paths, None).unwrap();
1156        assert_eq!(
1157            out.nodes_cancelled
1158                .iter()
1159                .map(NodeId::as_str)
1160                .collect::<Vec<_>>(),
1161            vec!["n-0002"]
1162        );
1163        assert_eq!(
1164            out.nodes_already_terminal
1165                .iter()
1166                .map(NodeId::as_str)
1167                .collect::<Vec<_>>(),
1168            vec!["n-0001"]
1169        );
1170        assert_eq!(
1171            node_status(&paths, "n-0001"),
1172            Status::Done,
1173            "Done node untouched"
1174        );
1175        assert_eq!(node_status(&paths, "n-0002"), Status::Cancelled);
1176    }
1177
1178    #[test]
1179    fn cancel_run_with_no_nodes_dir_settles_run_only() {
1180        let tmp = TempDir::new().unwrap();
1181        let paths = fresh_run(&tmp);
1182        append_and_apply_event(
1183            &paths,
1184            "run.created",
1185            None,
1186            None,
1187            json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
1188        )
1189        .unwrap();
1190
1191        let out = cancel_run(&paths, None).unwrap();
1192        assert!(!out.run_was_already_cancelled);
1193        assert_eq!(out.nodes_cancelled.len(), 0);
1194        assert_eq!(out.nodes_already_terminal.len(), 0);
1195        assert_eq!(
1196            crate::read_manifest(&paths).unwrap().status,
1197            Status::Cancelled
1198        );
1199    }
1200
1201    #[test]
1202    fn blank_note_falls_back_to_default_reason_and_does_not_brick_cancel() {
1203        // A `--note ""` (or whitespace-only) must NOT flow an empty `reason`
1204        // into the synthesized report — that would be rejected by the reducer
1205        // mid-loop and leave the run permanently un-cancellable. It normalizes
1206        // to the default reason and the cancel completes cleanly.
1207        for blank in ["", "   ", "\n\t"] {
1208            let tmp = TempDir::new().unwrap();
1209            let paths = fresh_run(&tmp);
1210            bootstrap(&paths, 1);
1211
1212            let out = cancel_run(&paths, Some(blank)).unwrap();
1213            assert_eq!(
1214                out.nodes_cancelled
1215                    .iter()
1216                    .map(NodeId::as_str)
1217                    .collect::<Vec<_>>(),
1218                vec!["n-0001"],
1219                "blank note {blank:?} still converges the live node"
1220            );
1221            assert_eq!(node_status(&paths, "n-0001"), Status::Cancelled);
1222            let report = crate::read_node(&paths, &NodeId::parse_str("n-0001").unwrap())
1223                .unwrap()
1224                .last_report
1225                .expect("cancel report recorded");
1226            assert_eq!(report["reason"], "cancelled by user");
1227        }
1228    }
1229
1230    #[test]
1231    fn nodes_are_converged_in_numeric_not_lexical_order() {
1232        // Past the digit-width boundary, lexical order would place n-10000
1233        // before n-9999. The numeric sort keeps the reported order intuitive.
1234        let tmp = TempDir::new().unwrap();
1235        let paths = fresh_run(&tmp);
1236        append_and_apply_event(
1237            &paths,
1238            "run.created",
1239            None,
1240            None,
1241            json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
1242        )
1243        .unwrap();
1244        for node in ["n-9999", "n-10000", "n-0001"] {
1245            append_and_apply_event(
1246                &paths,
1247                "node.created",
1248                Some(&nid(node)),
1249                None,
1250                json!({ "kind": "spinoff" }),
1251            )
1252            .unwrap();
1253        }
1254
1255        let out = cancel_run(&paths, None).unwrap();
1256        assert_eq!(
1257            out.nodes_cancelled
1258                .iter()
1259                .map(NodeId::as_str)
1260                .collect::<Vec<_>>(),
1261            vec!["n-0001", "n-9999", "n-10000"],
1262        );
1263    }
1264
1265    #[test]
1266    fn cancel_synthesizes_report_for_node_with_missing_projection() {
1267        // The crash window this fix closes: a `node.created` was appended+fsynced
1268        // to the log, but its projection write (`nodes/n-NNNN.json`) was
1269        // interrupted. A `nodes/*.json` scan would not see n-0002 and would
1270        // cancel the run while leaving a created-but-never-cancelled node a
1271        // future rebuild could resurrect as live. Enumerating from the event log
1272        // sees it and synthesizes the cancel report.
1273        let tmp = TempDir::new().unwrap();
1274        let paths = fresh_run(&tmp);
1275        bootstrap(&paths, 2);
1276        // Delete n-0002's projection file, leaving its `node.created` event in
1277        // the log — exactly the interrupted-fold state.
1278        let n2 = NodeId::parse_str("n-0002").unwrap();
1279        std::fs::remove_file(paths.node(&n2)).unwrap();
1280        assert!(
1281            read_node_opt(&paths, &n2).unwrap().is_none(),
1282            "projection gone"
1283        );
1284
1285        let out = cancel_run(&paths, Some("stop")).unwrap();
1286        // Both nodes are cancelled — the projection-present n-0001 AND the
1287        // projection-missing n-0002.
1288        assert_eq!(
1289            out.nodes_cancelled
1290                .iter()
1291                .map(NodeId::as_str)
1292                .collect::<Vec<_>>(),
1293            vec!["n-0001", "n-0002"],
1294            "the node with a missing projection is still cancelled"
1295        );
1296        assert_eq!(out.nodes_already_terminal.len(), 0);
1297        // The source-of-truth log now carries a terminal cancel report for the
1298        // node whose projection was missing — so a rebuild reconstructs it as
1299        // Cancelled, not live.
1300        assert_eq!(report_count(&paths, "n-0002"), 1);
1301        assert_eq!(
1302            crate::read_manifest(&paths).unwrap().status,
1303            Status::Cancelled
1304        );
1305    }
1306
1307    #[test]
1308    fn cancel_takes_the_run_lock_exactly_once() {
1309        // The single-lock honesty guarantee: the whole transaction (N node
1310        // reports + the run.status append) runs under ONE flock acquisition, not
1311        // one per appended event. Spy on `RunLock::acquire` to prove it.
1312        let tmp = TempDir::new().unwrap();
1313        let paths = fresh_run(&tmp);
1314        bootstrap(&paths, 5);
1315
1316        // Bootstrap itself takes the lock once per append; only the cancel call
1317        // is under measurement.
1318        ACQUIRE_COUNT.with(|c| c.set(0));
1319        let out = cancel_run(&paths, Some("stop")).unwrap();
1320        assert_eq!(out.nodes_cancelled.len(), 5);
1321        assert_eq!(
1322            ACQUIRE_COUNT.with(std::cell::Cell::get),
1323            1,
1324            "cancel must take the run lock exactly once, not once per node (N+1)"
1325        );
1326    }
1327
1328    #[test]
1329    fn cancel_does_not_duplicate_a_node_report_already_in_the_log() {
1330        // Crash-retry idempotency: a prior cancel appended+fsynced a node's
1331        // cancel `node.report` (carrying the deterministic key) but crashed
1332        // before folding its projection, so the node still reads live. A
1333        // re-cancel must NOT append a second logical-cancel event for it.
1334        let tmp = TempDir::new().unwrap();
1335        let paths = fresh_run(&tmp);
1336        bootstrap(&paths, 1); // run.created (seq 1) + node.created (seq 2)
1337
1338        // Durably append the cancel report WITH the deterministic key, but
1339        // without folding it — the node stays Pending (live), modeling the
1340        // fsynced-but-not-applied window.
1341        let node = nid("n-0001");
1342        let key = node_cancel_key(&paths.run_id, &node);
1343        RunLock::with_lock(&paths, |lock| {
1344            append_event_with_seq(
1345                lock,
1346                &paths,
1347                3,
1348                "node.report",
1349                Some(&node),
1350                Some(&key),
1351                json!({ "success": false, "cancelled": true, "reason": "x" }),
1352            )
1353        })
1354        .unwrap();
1355        assert_eq!(node_status(&paths, "n-0001"), Status::Pending);
1356        assert_eq!(report_count(&paths, "n-0001"), 1);
1357
1358        let out = cancel_run(&paths, None).unwrap();
1359        // The node converges (it is reported cancelled) but no duplicate report
1360        // is appended — the log still holds exactly one `node.report` for it.
1361        assert_eq!(
1362            out.nodes_cancelled
1363                .iter()
1364                .map(NodeId::as_str)
1365                .collect::<Vec<_>>(),
1366            vec!["n-0001"],
1367        );
1368        assert_eq!(
1369            report_count(&paths, "n-0001"),
1370            1,
1371            "the already-logged cancel report must not be duplicated"
1372        );
1373        // Convergence: the crash-stranded projection is folded from the
1374        // already-logged event, so the node reads Cancelled (not the stale
1375        // Pending) even though no new event was appended for it.
1376        assert_eq!(
1377            node_status(&paths, "n-0001"),
1378            Status::Cancelled,
1379            "the already-logged cancel must be re-folded, not just skipped"
1380        );
1381        assert_eq!(
1382            crate::read_manifest(&paths).unwrap().status,
1383            Status::Cancelled
1384        );
1385    }
1386
1387    #[test]
1388    fn cancel_does_not_duplicate_run_status_already_in_the_log() {
1389        // The run-status analogue: a prior cancel fsynced `run.status: cancelled`
1390        // (with its deterministic key) but crashed before folding the manifest,
1391        // so the manifest still reads non-terminal. A re-cancel must not append a
1392        // second `run.status: cancelled`.
1393        let tmp = TempDir::new().unwrap();
1394        let paths = fresh_run(&tmp);
1395        bootstrap(&paths, 0); // run.created only (seq 1)
1396
1397        let key = run_status_cancel_key(&paths.run_id);
1398        RunLock::with_lock(&paths, |lock| {
1399            append_event_with_seq(
1400                lock,
1401                &paths,
1402                2,
1403                "run.status",
1404                None,
1405                Some(&key),
1406                json!({ "status": "cancelled" }),
1407            )
1408        })
1409        .unwrap();
1410        // Manifest never folded the cancel, so it is not terminal here.
1411        assert_ne!(
1412            crate::read_manifest(&paths).unwrap().status,
1413            Status::Cancelled
1414        );
1415        let before = read_all_events(&paths.events()).unwrap().len();
1416
1417        let out = cancel_run(&paths, None).unwrap();
1418        assert!(!out.run_was_already_cancelled);
1419        assert_eq!(
1420            read_all_events(&paths.events()).unwrap().len(),
1421            before,
1422            "no duplicate run.status appended when one is already logged"
1423        );
1424        // Convergence: the manifest is folded from the already-logged
1425        // `run.status: cancelled` instead of being left stale.
1426        assert_eq!(
1427            crate::read_manifest(&paths).unwrap().status,
1428            Status::Cancelled,
1429            "the already-logged run.status must be re-folded, not just skipped"
1430        );
1431    }
1432
1433    #[test]
1434    fn cancel_skips_node_terminal_in_log_despite_stale_live_projection() {
1435        // cancel-liveness-from-log: a non-cancel terminal event (here a
1436        // `node.status` to a terminal value) was fsynced to the log but its
1437        // projection fold was crash-interrupted, so `nodes/n-0001.json` still
1438        // reads the stale live (Pending) status. The cancel must derive liveness
1439        // from the LOG and treat the node as already-terminal — never
1440        // synthesizing a cancel that would over-write the log's terminal and
1441        // diverge on a future rebuild (which replays node.status: done FIRST and
1442        // drops the later cancel).
1443        let tmp = TempDir::new().unwrap();
1444        let paths = fresh_run(&tmp);
1445        bootstrap(&paths, 2); // run.created(1) + node.created n-0001(2), n-0002(3)
1446
1447        // Raw-append (no fold) a terminal `node.status` for n-0001: the log
1448        // records it Done, but the projection stays the stale crash-window
1449        // Pending.
1450        let n1 = nid("n-0001");
1451        RunLock::with_lock(&paths, |lock| {
1452            append_event_with_seq(
1453                lock,
1454                &paths,
1455                4,
1456                "node.status",
1457                Some(&n1),
1458                None,
1459                json!({ "status": "done" }),
1460            )
1461        })
1462        .unwrap();
1463        assert_eq!(
1464            node_status(&paths, "n-0001"),
1465            Status::Pending,
1466            "projection is the stale, crash-stranded live status"
1467        );
1468
1469        let out = cancel_run(&paths, Some("stop")).unwrap();
1470        // n-0001 is settled by the log, NOT freshly cancelled; only the
1471        // genuinely live n-0002 is cancelled.
1472        assert_eq!(
1473            out.nodes_already_terminal
1474                .iter()
1475                .map(NodeId::as_str)
1476                .collect::<Vec<_>>(),
1477            vec!["n-0001"],
1478            "the log-terminal node is reported already-terminal, not cancelled"
1479        );
1480        assert_eq!(
1481            out.nodes_cancelled
1482                .iter()
1483                .map(NodeId::as_str)
1484                .collect::<Vec<_>>(),
1485            vec!["n-0002"],
1486        );
1487        // No cancel report was synthesized for n-0001: the log still holds zero
1488        // `node.report` lines for it, so a rebuild reconstructs it from the
1489        // `node.status: done` (Done), not a divergent Cancelled.
1490        assert_eq!(
1491            report_count(&paths, "n-0001"),
1492            0,
1493            "no cancel over-write was appended for the log-terminal node"
1494        );
1495    }
1496
1497    #[test]
1498    fn cancel_skips_node_with_unfolded_success_report_in_log() {
1499        // The issue's headline case: a `node.report { success: true }` fsynced
1500        // but not folded leaves a stale-live projection. Liveness from the log
1501        // settles the node as Done (already-terminal); the old projection-derived
1502        // check would have wrongly cancelled it over its already-logged success.
1503        let tmp = TempDir::new().unwrap();
1504        let paths = fresh_run(&tmp);
1505        bootstrap(&paths, 1); // run.created(1) + node.created n-0001(2)
1506        let n1 = nid("n-0001");
1507        RunLock::with_lock(&paths, |lock| {
1508            append_event_with_seq(
1509                lock,
1510                &paths,
1511                3,
1512                "node.report",
1513                Some(&n1),
1514                None,
1515                json!({ "success": true }),
1516            )
1517        })
1518        .unwrap();
1519        assert_eq!(
1520            node_status(&paths, "n-0001"),
1521            Status::Pending,
1522            "stale live projection (success report fsynced but not folded)"
1523        );
1524
1525        let out = cancel_run(&paths, None).unwrap();
1526        assert_eq!(
1527            out.nodes_already_terminal
1528                .iter()
1529                .map(NodeId::as_str)
1530                .collect::<Vec<_>>(),
1531            vec!["n-0001"],
1532        );
1533        assert!(
1534            out.nodes_cancelled.is_empty(),
1535            "a node the log shows Done must not be cancelled"
1536        );
1537        assert_eq!(
1538            report_count(&paths, "n-0001"),
1539            1,
1540            "only the original success report remains; no cancel was appended"
1541        );
1542    }
1543
1544    #[test]
1545    fn cancel_ledger_streams_large_report_payloads() {
1546        // The streaming ledger skims each line's envelope + a few small status
1547        // fields, never materializing the (here multi-KB) `node.report` `data`
1548        // payload. A node settled by such a report is still correctly seen as
1549        // terminal from the log, and a live sibling is still cancelled — proving
1550        // liveness is derived without holding whole reports in memory.
1551        let tmp = TempDir::new().unwrap();
1552        let paths = fresh_run(&tmp);
1553        bootstrap(&paths, 2);
1554        let big = "x".repeat(64 * 1024);
1555        append_and_apply_event(
1556            &paths,
1557            "node.report",
1558            Some(&nid("n-0001")),
1559            None,
1560            json!({ "success": true, "summary": big }),
1561        )
1562        .unwrap();
1563        assert_eq!(node_status(&paths, "n-0001"), Status::Done);
1564
1565        let out = cancel_run(&paths, Some("stop")).unwrap();
1566        assert_eq!(
1567            out.nodes_already_terminal
1568                .iter()
1569                .map(NodeId::as_str)
1570                .collect::<Vec<_>>(),
1571            vec!["n-0001"],
1572        );
1573        assert_eq!(
1574            out.nodes_cancelled
1575                .iter()
1576                .map(NodeId::as_str)
1577                .collect::<Vec<_>>(),
1578            vec!["n-0002"],
1579        );
1580    }
1581
1582    // --- per-node cancel (`cancel_node`) -----------------------------------
1583
1584    #[test]
1585    fn cancel_node_settles_one_node_and_leaves_the_run_and_siblings_live() {
1586        // The fan-out headline: cancelling one live child settles ONLY that node,
1587        // preserves it as Cancelled, and leaves the run + every sibling untouched
1588        // and non-terminal — the supervisor's rollup (not this call) terminalizes
1589        // the batch later.
1590        let tmp = TempDir::new().unwrap();
1591        let paths = fresh_run(&tmp);
1592        bootstrap(&paths, 3);
1593
1594        let out = cancel_node(&paths, &nid("n-0002"), Some("stuck")).unwrap();
1595        assert_eq!(out.node_id.as_str(), "n-0002");
1596        assert!(out.cancelled);
1597        assert!(!out.already_terminal);
1598        assert_eq!(out.rolled_up, None, "siblings live → run not rolled up");
1599
1600        assert_eq!(node_status(&paths, "n-0002"), Status::Cancelled);
1601        assert_eq!(node_status(&paths, "n-0001"), Status::Pending);
1602        assert_eq!(node_status(&paths, "n-0003"), Status::Pending);
1603        assert!(
1604            !crate::read_manifest(&paths).unwrap().status.is_terminal(),
1605            "no run.status is appended by a per-node cancel while siblings are live"
1606        );
1607        // The synthesized report carries the branch-preserving cancel shape.
1608        let report = crate::read_node(&paths, &nid("n-0002"))
1609            .unwrap()
1610            .last_report
1611            .expect("cancel report recorded");
1612        assert_eq!(report["cancelled"], true);
1613        assert_eq!(report["success"], false);
1614        assert_eq!(report["reason"], "stuck");
1615    }
1616
1617    #[test]
1618    fn cancel_node_unknown_id_is_node_not_found() {
1619        let tmp = TempDir::new().unwrap();
1620        let paths = fresh_run(&tmp);
1621        bootstrap(&paths, 1);
1622        let err = cancel_node(&paths, &nid("n-0009"), None).unwrap_err();
1623        assert!(
1624            matches!(err, Error::NodeNotFound { ref node_id } if node_id == "n-0009"),
1625            "got {err:?}"
1626        );
1627    }
1628
1629    #[test]
1630    fn cancel_node_on_already_terminal_node_is_idempotent_noop() {
1631        // A node that finished on its own (Done) is reported already-terminal,
1632        // never freshly cancelled, and no cancel report is appended over its
1633        // success.
1634        let tmp = TempDir::new().unwrap();
1635        let paths = fresh_run(&tmp);
1636        bootstrap(&paths, 2);
1637        append_and_apply_event(
1638            &paths,
1639            "node.report",
1640            Some(&nid("n-0001")),
1641            None,
1642            json!({ "success": true }),
1643        )
1644        .unwrap();
1645
1646        let out = cancel_node(&paths, &nid("n-0001"), None).unwrap();
1647        assert!(!out.cancelled);
1648        assert!(out.already_terminal);
1649        assert_eq!(node_status(&paths, "n-0001"), Status::Done, "untouched");
1650        assert_eq!(report_count(&paths, "n-0001"), 1, "no cancel over-write");
1651    }
1652
1653    #[test]
1654    fn cancel_node_twice_does_not_duplicate_the_report() {
1655        // Idempotent duplicate per-node cancel: the second call converges/no-ops
1656        // and never appends a second cancel `node.report`.
1657        let tmp = TempDir::new().unwrap();
1658        let paths = fresh_run(&tmp);
1659        bootstrap(&paths, 2);
1660
1661        let first = cancel_node(&paths, &nid("n-0001"), Some("x")).unwrap();
1662        assert!(first.cancelled);
1663        assert_eq!(report_count(&paths, "n-0001"), 1);
1664
1665        let second = cancel_node(&paths, &nid("n-0001"), Some("x")).unwrap();
1666        assert!(!second.cancelled);
1667        assert!(second.already_terminal);
1668        assert_eq!(
1669            report_count(&paths, "n-0001"),
1670            1,
1671            "a duplicate per-node cancel must not append a second report"
1672        );
1673        assert_eq!(node_status(&paths, "n-0001"), Status::Cancelled);
1674    }
1675
1676    #[test]
1677    fn cancel_node_converges_a_crash_stranded_prior_cancel_without_duplicating() {
1678        // Crash-retry: a prior cancel fsynced the node's cancel `node.report`
1679        // (with the deterministic key) but crashed before folding the projection,
1680        // so the node still reads live. A re-cancel re-folds the logged event
1681        // (node → Cancelled) without appending a second report.
1682        let tmp = TempDir::new().unwrap();
1683        let paths = fresh_run(&tmp);
1684        bootstrap(&paths, 1); // run.created(1) + node.created(2)
1685        let node = nid("n-0001");
1686        let key = node_cancel_key(&paths.run_id, &node);
1687        RunLock::with_lock(&paths, |lock| {
1688            append_event_with_seq(
1689                lock,
1690                &paths,
1691                3,
1692                "node.report",
1693                Some(&node),
1694                Some(&key),
1695                json!({ "success": false, "cancelled": true, "reason": "x" }),
1696            )
1697        })
1698        .unwrap();
1699        assert_eq!(node_status(&paths, "n-0001"), Status::Pending);
1700
1701        let out = cancel_node(&paths, &node, None).unwrap();
1702        assert!(out.cancelled, "the stranded cancel is converged");
1703        assert!(!out.already_terminal);
1704        assert_eq!(report_count(&paths, "n-0001"), 1, "no duplicate append");
1705        assert_eq!(node_status(&paths, "n-0001"), Status::Cancelled);
1706    }
1707
1708    #[test]
1709    fn cancel_node_resolves_a_node_with_a_missing_projection() {
1710        // The log — not the `nodes/*.json` scan — is authoritative for the node
1711        // set: a node whose projection write was crash-interrupted is still
1712        // cancellable (and its cancel report lands so a rebuild reconstructs it
1713        // Cancelled, not live).
1714        let tmp = TempDir::new().unwrap();
1715        let paths = fresh_run(&tmp);
1716        bootstrap(&paths, 2);
1717        let n2 = nid("n-0002");
1718        std::fs::remove_file(paths.node(&n2)).unwrap();
1719        assert!(read_node_opt(&paths, &n2).unwrap().is_none());
1720
1721        let out = cancel_node(&paths, &n2, Some("stop")).unwrap();
1722        assert!(out.cancelled);
1723        // The source-of-truth log now carries the terminal cancel report, so a
1724        // rebuild reconstructs the node Cancelled (its projection stays absent —
1725        // the reducer folds a report without resurrecting a deleted projection,
1726        // exactly as the whole-run cancel does).
1727        assert_eq!(report_count(&paths, "n-0002"), 1);
1728    }
1729
1730    #[test]
1731    fn cancel_node_blank_note_falls_back_to_default_reason() {
1732        let tmp = TempDir::new().unwrap();
1733        let paths = fresh_run(&tmp);
1734        bootstrap(&paths, 1);
1735        let out = cancel_node(&paths, &nid("n-0001"), Some("   ")).unwrap();
1736        assert!(out.cancelled);
1737        let report = crate::read_node(&paths, &nid("n-0001"))
1738            .unwrap()
1739            .last_report
1740            .expect("cancel report recorded");
1741        assert_eq!(report["reason"], "cancelled by user");
1742    }
1743
1744    #[test]
1745    fn cancel_node_takes_the_run_lock_exactly_once() {
1746        let tmp = TempDir::new().unwrap();
1747        let paths = fresh_run(&tmp);
1748        bootstrap(&paths, 3);
1749        ACQUIRE_COUNT.with(|c| c.set(0));
1750        let out = cancel_node(&paths, &nid("n-0002"), Some("x")).unwrap();
1751        assert!(out.cancelled);
1752        assert_eq!(
1753            ACQUIRE_COUNT.with(std::cell::Cell::get),
1754            1,
1755            "per-node cancel must take the run lock exactly once"
1756        );
1757    }
1758
1759    #[test]
1760    fn cancel_last_live_node_rolls_the_run_up_under_the_same_lock() {
1761        // Cancelling the final live node terminalizes the run HERE (llm-review
1762        // C1) — not deferred to a possibly-dead supervisor. Every node cancelled,
1763        // none failed → the run rolls up to Cancelled in the same transaction.
1764        let tmp = TempDir::new().unwrap();
1765        let paths = fresh_run(&tmp);
1766        bootstrap(&paths, 2);
1767
1768        // First cancel: n-0002 still live → run stays live.
1769        let first = cancel_node(&paths, &nid("n-0001"), Some("x")).unwrap();
1770        assert!(first.cancelled);
1771        assert_eq!(first.rolled_up, None);
1772        assert!(!crate::read_manifest(&paths).unwrap().status.is_terminal());
1773
1774        // Second cancel settles the last live node → the run rolls up.
1775        let out = cancel_node(&paths, &nid("n-0002"), Some("x")).unwrap();
1776        assert!(out.cancelled);
1777        assert_eq!(out.rolled_up, Some(Status::Cancelled));
1778        assert_eq!(node_status(&paths, "n-0001"), Status::Cancelled);
1779        assert_eq!(node_status(&paths, "n-0002"), Status::Cancelled);
1780        assert_eq!(
1781            crate::read_manifest(&paths).unwrap().status,
1782            Status::Cancelled,
1783            "the last per-node cancel terminalizes the run itself"
1784        );
1785    }
1786
1787    #[test]
1788    fn cancel_last_live_node_rolls_up_to_failed_when_a_sibling_failed() {
1789        // A genuine failure dominates the roll-up: cancelling the last live node
1790        // of a batch where a sibling already failed rolls the run up to Failed
1791        // (not Cancelled).
1792        let tmp = TempDir::new().unwrap();
1793        let paths = fresh_run(&tmp);
1794        bootstrap(&paths, 2);
1795        append_and_apply_event(
1796            &paths,
1797            "node.report",
1798            Some(&nid("n-0001")),
1799            None,
1800            json!({ "success": false }),
1801        )
1802        .unwrap();
1803        assert_eq!(node_status(&paths, "n-0001"), Status::Failed);
1804
1805        let out = cancel_node(&paths, &nid("n-0002"), Some("x")).unwrap();
1806        assert!(out.cancelled);
1807        assert_eq!(out.rolled_up, Some(Status::Failed));
1808        assert_eq!(crate::read_manifest(&paths).unwrap().status, Status::Failed);
1809    }
1810
1811    #[test]
1812    fn cancel_last_live_node_rolls_up_to_cancelled_on_done_plus_cancelled_mix() {
1813        // Some siblings merged (Done), the last is cancelled, none failed → the
1814        // batch rolls up to Cancelled (nothing failed, not a clean all-Done).
1815        let tmp = TempDir::new().unwrap();
1816        let paths = fresh_run(&tmp);
1817        bootstrap(&paths, 2);
1818        append_and_apply_event(
1819            &paths,
1820            "node.report",
1821            Some(&nid("n-0001")),
1822            None,
1823            json!({ "success": true }),
1824        )
1825        .unwrap();
1826        assert_eq!(node_status(&paths, "n-0001"), Status::Done);
1827
1828        let out = cancel_node(&paths, &nid("n-0002"), Some("x")).unwrap();
1829        assert!(out.cancelled);
1830        assert_eq!(out.rolled_up, Some(Status::Cancelled));
1831        assert_eq!(
1832            crate::read_manifest(&paths).unwrap().status,
1833            Status::Cancelled
1834        );
1835    }
1836
1837    #[test]
1838    fn cancel_node_refuses_a_done_run() {
1839        // Mirror `cancel_run`'s guard (llm-review C4): a Done/Failed run is
1840        // refused rather than appending a dead post-terminal cancel.
1841        let tmp = TempDir::new().unwrap();
1842        let paths = fresh_run(&tmp);
1843        bootstrap(&paths, 1);
1844        append_and_apply_event(
1845            &paths,
1846            "node.report",
1847            Some(&nid("n-0001")),
1848            None,
1849            json!({ "success": true }),
1850        )
1851        .unwrap();
1852        append_and_apply_event(
1853            &paths,
1854            "run.status",
1855            None,
1856            None,
1857            json!({ "status": "done" }),
1858        )
1859        .unwrap();
1860        let before = read_all_events(&paths.events()).unwrap().len();
1861
1862        let err = cancel_node(&paths, &nid("n-0001"), None).unwrap_err();
1863        assert!(
1864            matches!(
1865                err,
1866                Error::RunAlreadyTerminal {
1867                    status: Status::Done
1868                }
1869            ),
1870            "got {err:?}"
1871        );
1872        assert_eq!(
1873            read_all_events(&paths.events()).unwrap().len(),
1874            before,
1875            "a refused per-node cancel must not append any event"
1876        );
1877    }
1878
1879    #[test]
1880    fn cancel_node_last_node_shares_run_status_key_with_whole_run_cancel() {
1881        // The last-node roll-up uses the whole-run cancel's `run-status` key, so a
1882        // later `cancel_run` converges on the SAME logical run.status rather than
1883        // appending a duplicate.
1884        let tmp = TempDir::new().unwrap();
1885        let paths = fresh_run(&tmp);
1886        bootstrap(&paths, 1);
1887        let out = cancel_node(&paths, &nid("n-0001"), Some("x")).unwrap();
1888        assert_eq!(out.rolled_up, Some(Status::Cancelled));
1889        let run_status_events = read_all_events(&paths.events())
1890            .unwrap()
1891            .into_iter()
1892            .filter(|e| e.kind == "run.status")
1893            .count();
1894        assert_eq!(run_status_events, 1);
1895
1896        // A whole-run cancel now finds the run already cancelled and converges —
1897        // no second run.status.
1898        let cr = cancel_run(&paths, Some("x")).unwrap();
1899        assert!(cr.run_was_already_cancelled);
1900        let run_status_events = read_all_events(&paths.events())
1901            .unwrap()
1902            .into_iter()
1903            .filter(|e| e.kind == "run.status")
1904            .count();
1905        assert_eq!(
1906            run_status_events, 1,
1907            "cancel_run must not duplicate the run.status the last-node roll-up wrote"
1908        );
1909    }
1910}