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}