Skip to main content

taskfleet_core/
events.rs

1//! Event append primitive + `seq` recovery (design.md §1.4, §4).
2
3use std::io::{BufRead, BufReader, Read, Seek, SeekFrom, Write};
4use std::path::{Path, PathBuf};
5
6use chrono::Utc;
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9
10use crate::atomic::{open_events_append, write_atomic};
11use crate::error::{Error, Result};
12use crate::lock::{LockedRun, RunLock};
13use crate::paths::RunPaths;
14use crate::projections::{derive_counters, read_manifest_opt, write_manifest};
15use crate::reducer::{commit_ops, reduce_event_to_ops};
16use crate::schema::{Event, NodeId};
17
18/// Backward-scan chunk size when looking for the previous newline.
19const SCAN_CHUNK: u64 = 64 * 1024;
20
21/// Read the last `seq` from `events.jsonl`, or `0` if empty/missing.
22///
23/// Tolerates:
24/// - lines larger than any fixed buffer (`node.report` payloads can be 10s of KB
25///   per `design.md` §1.4) — we scan backwards in chunks for the previous `\n`.
26/// - a crash-truncated final line lacking a trailing `\n` — that partial tail
27///   is discarded and recovery uses the last complete record.
28///
29/// Caller must already hold the run's [`RunLock`] for correctness against
30/// concurrent appenders.
31pub fn recover_last_seq(events_path: &Path) -> Result<u64> {
32    let mut f = match std::fs::File::open(events_path) {
33        Ok(f) => f,
34        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(0),
35        Err(e) => return Err(Error::io(events_path, e)),
36    };
37    let len = f.metadata().map_err(|e| Error::io(events_path, e))?.len();
38    if len == 0 {
39        return Ok(0);
40    }
41
42    // Require a newline-terminated final line; otherwise treat the last
43    // partial chunk as torn and recover from the previous complete line.
44    let mut tail_byte = [0u8; 1];
45    f.seek(SeekFrom::End(-1))
46        .map_err(|e| Error::io(events_path, e))?;
47    f.read_exact(&mut tail_byte)
48        .map_err(|e| Error::io(events_path, e))?;
49    let mut end = if tail_byte[0] == b'\n' {
50        len - 1
51    } else {
52        match find_prev_newline(&mut f, len, events_path)? {
53            Some(p) => p,
54            None => return Ok(0),
55        }
56    };
57
58    // `end` is the byte index of the trailing `\n` of the last complete
59    // record. Walk backward over complete lines, skipping any that are empty
60    // or whitespace-only — consecutive newlines or blank/whitespace lines (e.g.
61    // from external editing) shouldn't fool recovery into reading the wrong
62    // last record — and recover the seq from the last line bearing real bytes.
63    loop {
64        let line_start = match find_prev_newline(&mut f, end, events_path)? {
65            Some(p) => p + 1,
66            None => 0,
67        };
68        let line_len = end - line_start;
69        f.seek(SeekFrom::Start(line_start))
70            .map_err(|e| Error::io(events_path, e))?;
71        let mut line = vec![0u8; line_len as usize];
72        f.read_exact(&mut line)
73            .map_err(|e| Error::io(events_path, e))?;
74        // Any non-whitespace byte means a real record — parse it. Lines that
75        // are empty or hold only ASCII whitespace (a stray `\r`, `\t`, or
76        // spaces left by external editing) carry no record, so skip them and
77        // keep scanning back; serde tolerates whitespace surrounding a real
78        // envelope, so a genuine record with trailing spaces still parses.
79        if line.iter().any(|b| !b.is_ascii_whitespace()) {
80            return parse_seq(&line, events_path);
81        }
82        // Whitespace-only line: no record here. Step to the newline before it
83        // and keep scanning; reaching the start means the log holds no event.
84        if line_start == 0 {
85            return Ok(0);
86        }
87        end = line_start - 1;
88    }
89}
90
91/// The envelope fields recovered from the last complete line. Required fields
92/// mirror [`Event`]'s required shape, so `recover_last_seq` accepts a last line
93/// iff [`read_all_events`] would — the two readers agree on what the last
94/// record is. `data` / `idempotency_key` are skipped (serde ignores unknown
95/// fields) so a multi-KB `node.report` payload isn't re-materialized on the
96/// hot append path just to read `seq`.
97#[derive(Deserialize)]
98#[allow(dead_code)] // fields exist to force serde validation, not to be read
99struct SeqLine {
100    seq: u64,
101    ts: chrono::DateTime<chrono::Utc>,
102    kind: String,
103    run_id: crate::schema::RunId,
104    #[serde(default)]
105    node_id: Option<NodeId>,
106}
107
108fn parse_seq(line: &[u8], events_path: &Path) -> Result<u64> {
109    // The last complete line must be a full, valid event envelope — the same
110    // bar `read_all_events` applies to every line — so a `\n`-terminated line
111    // that parses as JSON but isn't a valid event (e.g. `{"seq":1}` missing
112    // `ts`/`run_id`) is event-log corruption, not a usable seq source. This
113    // keeps the three readers aligned on the last record.
114    let hdr: SeqLine = serde_json::from_slice(line).map_err(|e| Error::CorruptEventLog {
115        path: events_path.to_path_buf(),
116        reason: format!(
117            "last complete line is not a valid event: {} [{e}]",
118            excerpt(line)
119        ),
120    })?;
121    Ok(hdr.seq)
122}
123
124/// Find the byte offset of the last `\n` strictly before `before`. Returns
125/// `None` if no newline exists in `[0, before)`.
126fn find_prev_newline(
127    f: &mut std::fs::File,
128    before: u64,
129    events_path: &Path,
130) -> Result<Option<u64>> {
131    if before == 0 {
132        return Ok(None);
133    }
134    let mut pos = before;
135    loop {
136        let start = pos.saturating_sub(SCAN_CHUNK);
137        let len = pos - start;
138        f.seek(SeekFrom::Start(start))
139            .map_err(|e| Error::io(events_path, e))?;
140        let mut buf = vec![0u8; len as usize];
141        f.read_exact(&mut buf)
142            .map_err(|e| Error::io(events_path, e))?;
143        if let Some(i) = buf.iter().rposition(|b| *b == b'\n') {
144            return Ok(Some(start + i as u64));
145        }
146        if start == 0 {
147            return Ok(None);
148        }
149        pos = start;
150    }
151}
152
153/// Truncate a torn (newline-less) final line off `events.jsonl` so the next
154/// append never concatenates onto a partial record.
155///
156/// `recover_last_seq` only *ignores* a torn tail for seq purposes — it never
157/// removes the bytes. Without this, an append after a crash-truncated write
158/// would write its `\n`-terminated line directly onto the partial bytes,
159/// producing one malformed `…torn…{"seq":…}` line that every later reader
160/// (now sharing a strict torn-tail policy) hard-errors on. Cutting back to
161/// the last complete record here guarantees the file is always empty or
162/// `\n`-terminated before we append.
163///
164/// Caller must hold the run's [`RunLock`]. No-op when the file is absent,
165/// empty, or already `\n`-terminated (the common, clean case — one `stat` +
166/// one-byte read, no rewrite).
167fn truncate_torn_tail(events_path: &Path) -> Result<()> {
168    let mut opts = std::fs::OpenOptions::new();
169    opts.read(true).write(true);
170    // `O_NOFOLLOW`: refuse to rewrite the tail through a symlinked event log.
171    crate::paths::nofollow(&mut opts);
172    let mut f = match opts.open(events_path) {
173        Ok(f) => f,
174        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
175        Err(e) => return Err(Error::io(events_path, e)),
176    };
177    let len = f.metadata().map_err(|e| Error::io(events_path, e))?.len();
178    if len == 0 {
179        return Ok(());
180    }
181    let mut tail = [0u8; 1];
182    f.seek(SeekFrom::End(-1))
183        .map_err(|e| Error::io(events_path, e))?;
184    f.read_exact(&mut tail)
185        .map_err(|e| Error::io(events_path, e))?;
186    if tail[0] == b'\n' {
187        return Ok(());
188    }
189    // Torn final line: cut back to just past the last complete record's
190    // trailing newline, or to empty when no complete record exists.
191    let keep = match find_prev_newline(&mut f, len, events_path)? {
192        Some(nl) => nl + 1,
193        None => 0,
194    };
195    f.set_len(keep).map_err(|e| Error::io(events_path, e))?;
196    f.sync_all().map_err(|e| Error::io(events_path, e))?;
197    // Surface the recovery so an operator inspecting the run knows a
198    // crash-torn tail was discarded (and how many bytes), rather than the
199    // truncation happening invisibly under the lock.
200    tracing::warn!(
201        target: "taskfleet_core::events",
202        path = %events_path.display(),
203        discarded_bytes = len - keep,
204        kept_bytes = keep,
205        "truncated crash-torn final line off events.jsonl before append"
206    );
207    Ok(())
208}
209
210/// Append one event with a caller-supplied `seq`. The `_witness: &LockedRun`
211/// is compile-time proof the caller holds the run's exclusive [`RunLock`] for
212/// the duration of this call; the caller is still responsible for ensuring
213/// `seq` is monotonic. Misuse can corrupt the event log.
214///
215/// Test-only (`#[cfg(test)]`): a raw, no-reducer, caller-managed-`seq`
216/// primitive used by the crate's fixtures and the flock stress test to craft
217/// event logs with explicit seqs. Production mutation goes through
218/// [`append_and_apply_event`]; projection rebuild (future) replays via
219/// [`crate::reducer`], so neither needs this.
220#[cfg(test)]
221pub(crate) fn append_event_with_seq(
222    _witness: &LockedRun<'_>,
223    paths: &RunPaths,
224    seq: u64,
225    kind: &str,
226    node_id: Option<&NodeId>,
227    idempotency_key: Option<&str>,
228    data: Value,
229) -> Result<()> {
230    write_event_line(paths, seq, kind, node_id, idempotency_key, data)
231}
232
233#[cfg(test)]
234fn write_event_line(
235    paths: &RunPaths,
236    seq: u64,
237    kind: &str,
238    node_id: Option<&NodeId>,
239    idempotency_key: Option<&str>,
240    data: Value,
241) -> Result<()> {
242    let ev = Event {
243        ts: Utc::now(),
244        seq,
245        kind: kind.to_string(),
246        run_id: paths.run_id.clone(),
247        node_id: node_id.cloned(),
248        idempotency_key: idempotency_key.map(str::to_string),
249        data,
250    };
251    let events_path = paths.events();
252    let mut line = serde_json::to_vec(&ev).map_err(|e| Error::json(events_path.clone(), e))?;
253    line.push(b'\n');
254    let mut f = open_events_append(&events_path)?;
255    f.write_all(&line)
256        .map_err(|e| Error::io(events_path.clone(), e))?;
257    f.sync_all().map_err(|e| Error::io(events_path, e))?;
258    Ok(())
259}
260
261/// Outcome of an [`append_and_apply_event`] call.
262///
263/// `seq` is the value a caller surfaces to a user: the freshly appended
264/// event's `seq`, or — on an idempotent replay — the `seq` of the
265/// pre-existing matching event. A reducer no-op (e.g. an event dropped by
266/// the terminal-state guard) is still a success at this layer: `seq` names
267/// the appended event regardless of whether the reducer changed anything.
268///
269/// There is intentionally no `derived_event_ids` field. This API mutates
270/// exactly one event; the supervisor's report consumption, which emits a
271/// *batch* of derived discussion/spinoff events under one held lock, uses
272/// [`append_and_apply_unlocked`] instead (the sanctioned lock-held
273/// composition path) and tracks its own emitted ids.
274#[derive(Debug, Serialize)]
275pub struct AppendResult {
276    /// `seq` of the appended event, or of the prior event on an idempotent
277    /// replay.
278    pub seq: u64,
279    /// True when `idempotency_key` matched a prior event so nothing new was
280    /// appended or applied; `seq`/`prior` then describe that prior event.
281    pub idempotent_replay: bool,
282    /// True when the reducer produced at least one projection write for THIS
283    /// append — i.e. the event actually changed state, rather than folding to a
284    /// no-op (an unknown/audit kind, or an event dropped by a `*.created` /
285    /// terminal-state guard). Lets a caller distinguish "the reducer applied my
286    /// event" from "it was a dead event" WITHOUT re-reading the projection and
287    /// pattern-matching a field (issue `reducer-adopt-explicit-merge`).
288    ///
289    /// This is a report of what the reducer did on THIS call, NOT a durable
290    /// "is teardown pending?" signal: it is `false` both on an idempotent replay
291    /// AND on a fresh append the reducer no-op'd (e.g. re-submitting the exact
292    /// report already adopted). Callers making a DURABLE decision (does the run
293    /// still need a teardown actor?) must read projection state, not this flag —
294    /// see `run merge`'s `ensure_report_consumer`, which deliberately does NOT gate
295    /// its reattach on `applied` (that was a crash-retry leak caught in review).
296    pub applied: bool,
297    /// On an idempotent replay, the prior event's recorded `node_id` and
298    /// `data`, so a caller can reject a key reused with a conflicting
299    /// request (Stripe-style). `None` on a fresh append.
300    #[serde(skip_serializing_if = "Option::is_none")]
301    pub prior: Option<PriorEvent>,
302}
303
304/// The one canonical mutation entry point: append a single event to
305/// `events.jsonl` *and* fold it into the projection files via the reducer,
306/// all under the run's `flock`, with idempotency-key dedup.
307///
308/// On success, every `events.jsonl` line is folded into `manifest.json` /
309/// `nodes/*.json` / `discussions/*.json` / `spinoffs/*.json` before the lock
310/// is released, so a read CLI run a millisecond later never sees a stale
311/// projection. This is *not* a crash-atomic transaction: the event is fsynced
312/// before the reducer runs, so a crash (or an I/O error from `apply_event`)
313/// after the append but before the projection write leaves the log ahead of
314/// the projections — recoverable only by a future `rebuild_projections`. The
315/// log is the source of truth; projections are a derived cache.
316///
317/// The append is transactional against reducer *validation*: the event is
318/// first reduced through [`reduce_event_to_ops`](crate::reducer) under the
319/// lock — the single plan-then-commit path that both validates and computes
320/// the projection writes — and only a validating event is appended (and
321/// fsynced) and then committed by the reducer. A reducer-rejected event (a
322/// `CorruptEventLog` for a malformed payload) errors *before* any bytes are
323/// written, so the log never gains a poison line that a future replay /
324/// `rebuild_projections` would choke on.
325/// (A pre-existing torn tail may still be truncated before validation runs —
326/// those bytes are uncommitted by definition; see [`recover_last_seq`].)
327///
328/// When `idempotency_key` is `Some` and a prior event with the same `kind` +
329/// key already exists ([`find_prior_with_key`](crate::events)), nothing is appended or
330/// applied: the result carries the prior event's `seq`, `idempotent_replay:
331/// true`, and `prior: Some(..)` so the caller can detect a key reused with a
332/// conflicting payload. With `idempotency_key: None` no scan runs.
333///
334/// Callers that must compose several writes — or a read-modify-write
335/// transaction (read a projection, decide, then append) — under one lock
336/// window hold the lock themselves and use [`append_and_apply_unlocked`],
337/// the sanctioned lock-held composition path. Re-entering this function
338/// while already holding the lock would deadlock: `flock` blocks when a
339/// second open of the lock file from the same process tries `LOCK_EX`.
340pub fn append_and_apply_event(
341    paths: &RunPaths,
342    kind: &str,
343    node_id: Option<&NodeId>,
344    idempotency_key: Option<&str>,
345    data: Value,
346) -> Result<AppendResult> {
347    RunLock::with_lock(paths, |lock| {
348        // Catch the projections up to the event log before either the
349        // idempotency lookup or a fresh append. This is the recovery half of
350        // append+apply atomicity: any unapplied tail left by a prior crash is
351        // folded here, under the same lock, so an idempotent replay returns
352        // only once the prior event's projection is durably committed
353        // (`applied_seq >= prior.seq`) — never a stale "found, but not applied"
354        // result. A clean run with no tail makes this a cheap no-op.
355        replay_unapplied_unlocked(lock, paths)?;
356        // Idempotency lookup + append share this one lock window so a
357        // concurrent retry can't see "no prior event" and double-append.
358        if let Some(key) = idempotency_key {
359            if let Some(prior) = find_prior_with_key(lock, paths, kind, key)? {
360                return Ok(AppendResult {
361                    seq: prior.seq,
362                    idempotent_replay: true,
363                    // Nothing was applied by THIS call — the prior event (already
364                    // folded) carried any state change.
365                    applied: false,
366                    prior: Some(prior),
367                });
368            }
369        }
370        let (seq, applied) =
371            append_and_apply_reporting(lock, paths, kind, node_id, idempotency_key, data)?;
372        Ok(AppendResult {
373            seq,
374            idempotent_replay: false,
375            applied,
376            prior: None,
377        })
378    })
379}
380
381/// Catch projections up to the durable event log under an already-held
382/// exclusive lock, without appending an event. Mutation commands that must
383/// decide from current projections use this before their read/authorize step.
384pub fn replay_unapplied_unlocked(_witness: &LockedRun<'_>, paths: &RunPaths) -> Result<()> {
385    let events_path = paths.checked_events()?;
386    truncate_torn_tail(&events_path)?;
387    replay_unapplied(paths, &events_path)
388}
389
390/// Append one event and fold it into projections. The `_witness: &LockedRun`
391/// is compile-time proof the caller already holds the run's exclusive
392/// [`RunLock`] — obtained from [`RunLock::with_lock`] or [`RunLock::witness`],
393/// so this entry point cannot be reached without the lock. The **sanctioned
394/// lock-held composition path**: use it to fold extra logic (an idempotency-key
395/// lookup, a status precondition) or several writes (the supervisor's
396/// derived discussion/spinoff batch) into one locked critical section.
397/// Calling [`append_and_apply_event`] from within a held lock would
398/// deadlock because `flock` blocks when a second open of the lock file from
399/// the same process tries to acquire `LOCK_EX`.
400///
401/// # The witness is mandatory
402///
403/// Without a `&LockedRun` proof the lock is held, this does not compile — there
404/// is no way to skip the parameter, and [`LockedRun`] cannot be constructed
405/// outside this crate (its field is private), so the only source is a held
406/// [`RunLock`]:
407///
408/// ```compile_fail
409/// use taskfleet_core::{append_and_apply_unlocked, RunPaths};
410/// # fn demo(paths: &RunPaths) {
411/// // No witness passed — the first argument must be a `&LockedRun`, which a
412/// // caller can only obtain by actually holding the run's exclusive lock.
413/// let _ = append_and_apply_unlocked(paths, "run.status", None, None, serde_json::json!({}));
414/// # }
415/// ```
416pub fn append_and_apply_unlocked(
417    witness: &LockedRun<'_>,
418    paths: &RunPaths,
419    kind: &str,
420    node_id: Option<&NodeId>,
421    idempotency_key: Option<&str>,
422    data: Value,
423) -> Result<u64> {
424    append_and_apply_reporting(witness, paths, kind, node_id, idempotency_key, data)
425        .map(|(seq, _)| seq)
426}
427
428/// As [`append_and_apply_unlocked`], but also reports whether the reducer APPLIED
429/// (produced ≥1 projection op) vs folded to a no-op — the `bool` feeding
430/// [`AppendResult::applied`]. Kept private so the public composition primitive
431/// stays `-> u64` for its 15+ callers (none of which need the applied bit); only
432/// [`append_and_apply_event`] threads it out. See [`AppendResult::applied`] for
433/// why callers want it (issue `reducer-adopt-explicit-merge`).
434fn append_and_apply_reporting(
435    witness: &LockedRun<'_>,
436    paths: &RunPaths,
437    kind: &str,
438    node_id: Option<&NodeId>,
439    idempotency_key: Option<&str>,
440    data: Value,
441) -> Result<(u64, bool)> {
442    // Direct lock-held callers get the same catch-up guarantee as the ordinary
443    // append wrapper before computing this event against projections.
444    replay_unapplied_unlocked(witness, paths)?;
445    let events_path = paths.checked_events()?;
446    let last = recover_last_seq(&events_path)?;
447    let seq = last + 1;
448    let ev = Event {
449        ts: Utc::now(),
450        seq,
451        kind: kind.to_string(),
452        run_id: paths.run_id.clone(),
453        node_id: node_id.cloned(),
454        idempotency_key: idempotency_key.map(str::to_string),
455        data,
456    };
457    // Transactional gate, plan-then-commit: reduce the event against current
458    // projection state BEFORE the durable append. `reduce_event_to_ops` both
459    // validates and computes the exact projection writes to make; a reducer-
460    // rejected event errors here and is never written, so a later replay /
461    // rebuild can't trip on a poison line. The planned ops are then committed
462    // *after* the fsynced append — nothing mutates the projections between the
463    // plan and the commit (the append only touches `events.jsonl`), so the
464    // planned writes are still valid. One reduce pass serves both the gate and
465    // the apply, so there is no validate/apply branch pair to drift apart.
466    let ops = reduce_event_to_ops(paths, &ev)?;
467    // Whether the reducer changed state for this event — reported to the caller
468    // via `AppendResult::applied`. Captured before `commit_ops` consumes `ops`.
469    let applied = !ops.is_empty();
470    let mut line = serde_json::to_vec(&ev).map_err(|e| Error::json(events_path.clone(), e))?;
471    line.push(b'\n');
472    let mut f = open_events_append(&events_path)?;
473    f.write_all(&line)
474        .map_err(|e| Error::io(events_path.clone(), e))?;
475    f.sync_all().map_err(|e| Error::io(events_path, e))?;
476    commit_ops(paths, ops)?;
477    // Advance the watermark only after every projection this event touched is
478    // durably committed. A crash before this point leaves `applied_seq < seq`,
479    // and the next lock acquisition replays the event (idempotently — the
480    // reducer's existence/terminal guards make a re-fold a no-op) before
481    // advancing. So the watermark can only ever lag the projections, never lead
482    // them — the projection a reader sees is always at least as new as
483    // `applied_seq` claims.
484    advance_applied_seq(paths, seq)?;
485    Ok((seq, applied))
486}
487
488/// The three observable outcomes of an [`append_and_apply_idempotent`] call —
489/// the shared `--idempotency-key` contract that `event create`, `discussion
490/// resolve`, and future keyed verbs (`spinoff approve|reject`, `run create`,
491/// `node report`) all answer to, lifted out of each CLI's private log scan.
492///
493/// The discriminator is whether a prior event with the same `kind` + key
494/// already exists, and — if so — whether the call's `(node_id, data)` identity
495/// matches that prior event:
496///
497/// - [`AppendOutcome::Appended`] — no prior event carried this key: a fresh
498///   event was appended and folded into the projections. `seq` is its sequence.
499/// - [`AppendOutcome::IdempotentReplay`] — a prior event carried this key **and**
500///   the same `node_id` + `data`: a true retry. Nothing was appended; the
501///   `prior` event (its `seq` / `node_id` / `data`) is returned so the caller
502///   can surface the original sequence.
503/// - [`AppendOutcome::Conflict`] — a prior event carried this key but with a
504///   **different** `node_id` or `data`: the key was reused for a different
505///   request (a client bug, Stripe-style). Nothing was appended; `prior` is
506///   returned so the caller can build a precise conflict error (e.g. diff the
507///   payload vs. the node id).
508#[derive(Debug)]
509pub enum AppendOutcome {
510    /// A fresh event was appended and applied; `seq` is its sequence number.
511    Appended {
512        /// The appended event's `seq`.
513        seq: u64,
514    },
515    /// The key matched a prior event with identical `node_id` + `data`. No new
516    /// event was written; `prior.seq` is the original sequence to surface.
517    IdempotentReplay {
518        /// The pre-existing matching event (its `seq`, `node_id`, and `data`).
519        prior: PriorEvent,
520    },
521    /// The key matched a prior event whose `node_id` or `data` differs from this
522    /// request. No new event was written; the caller should reject the reuse.
523    Conflict {
524        /// The pre-existing event recorded under the same key, for the caller's
525        /// conflict diagnostics (`prior.seq` is the original sequence).
526        prior: PriorEvent,
527    },
528}
529
530/// Append one keyed event idempotently: scan for a prior event with the same
531/// `kind` + `key`, and either replay it, reject a conflicting reuse, or append
532/// fresh — the centralized `--idempotency-key` primitive (issue
533/// `core-idempotency-api`).
534///
535/// This is the **sanctioned lock-held composition path** for keyed appends: the
536/// `_witness: &LockedRun` proves the caller already holds the run's exclusive
537/// [`RunLock`] (from [`RunLock::with_lock`] or [`RunLock::witness`]), so the
538/// scan and the append share one lock window and a concurrent retry can never
539/// see "no prior event" and double-append. Calling it composes with the
540/// applied-seq watermark and the path-traversal defense exactly as
541/// [`append_and_apply_unlocked`] does — it catches the projections up to the log
542/// (`truncate_torn_tail` + `replay_unapplied`) before scanning, guards the run
543/// root + event log via `RunPaths::checked_events`, and routes the fresh
544/// append through `append_and_apply_unlocked`.
545///
546/// `build` lazily produces the event's `data` payload given the sequence the
547/// fresh event *would* receive. It is a **pure** constructor: it is invoked once
548/// to materialize the candidate payload (to compare against a prior event, or to
549/// write a fresh one) and must not encode caller-side domain preconditions — a
550/// verb whose append is gated on projection state (e.g. `discussion resolve`'s
551/// already-resolved / no-op decision) keeps that logic in its own locked body
552/// and uses [`find_prior_with_key`] directly. The `u64` lets a payload embed its
553/// own `seq`; a payload that does so is not replay-stable and should not be used
554/// with idempotency.
555///
556/// The key must be non-empty: an empty key is rejected with
557/// [`Error::EmptyIdempotencyKey`] before any scan, since `""` would collapse
558/// every keyless append into one dedup slot.
559///
560/// # Examples
561///
562/// ```no_run
563/// use taskfleet_core::{append_and_apply_idempotent, AppendOutcome, RunLock, RunPaths};
564/// use serde_json::json;
565///
566/// # fn demo(paths: &RunPaths) -> taskfleet_core::Result<()> {
567/// let outcome = RunLock::with_lock(paths, |lock| {
568///     append_and_apply_idempotent(
569///         paths,
570///         lock,
571///         "node.status",
572///         None,            // no target node
573///         "retry-key-42",  // the caller's idempotency key (non-empty)
574///         |_seq| Ok(json!({ "status": "running" })),
575///     )
576/// })?;
577/// match outcome {
578///     AppendOutcome::Appended { seq } => println!("appended at seq {seq}"),
579///     AppendOutcome::IdempotentReplay { prior } => println!("replayed seq {}", prior.seq),
580///     AppendOutcome::Conflict { prior } => println!("key reused; prior seq {}", prior.seq),
581/// }
582/// # Ok(())
583/// # }
584/// ```
585pub fn append_and_apply_idempotent<F>(
586    paths: &RunPaths,
587    witness: &LockedRun<'_>,
588    kind: &str,
589    node_id: Option<&NodeId>,
590    key: &str,
591    build: F,
592) -> Result<AppendOutcome>
593where
594    F: FnOnce(u64) -> Result<Value>,
595{
596    if key.is_empty() {
597        return Err(Error::EmptyIdempotencyKey);
598    }
599    // Catch the projections up to the log before scanning, mirroring
600    // `append_and_apply_unlocked`'s recovery half: an idempotent replay must
601    // only report once the prior event's projection is durably committed, never
602    // a stale "found, but not applied" result. A clean run makes this a no-op.
603    let events_path = paths.checked_events()?;
604    truncate_torn_tail(&events_path)?;
605    replay_unapplied(paths, &events_path)?;
606
607    // The sequence a fresh append *would* take. Computed once, after catch-up,
608    // so `build`'s payload sees the same seq `append_and_apply_unlocked` will
609    // assign under this still-held lock.
610    let next_seq = recover_last_seq(&events_path)? + 1;
611    let data = build(next_seq)?;
612
613    if let Some(prior) = find_prior_with_key(witness, paths, kind, key)? {
614        // A prior event carries this key. It is a true replay only when the
615        // full request identity — the envelope `node_id` *and* the `data`
616        // payload — matches; any divergence is a key reused for a different
617        // request and must surface as a conflict, never a silent no-op.
618        let same_node = prior.node_id.as_deref() == node_id.map(NodeId::as_str);
619        if same_node && prior.data == data {
620            return Ok(AppendOutcome::IdempotentReplay { prior });
621        }
622        return Ok(AppendOutcome::Conflict { prior });
623    }
624
625    let seq = append_and_apply_unlocked(witness, paths, kind, node_id, Some(key), data)?;
626    Ok(AppendOutcome::Appended { seq })
627}
628
629/// Replay every unapplied tail event — those with `seq > manifest.applied_seq`
630/// — into the projections, advancing the watermark after each, so the
631/// projection cache is caught up to `events.jsonl` before any new append.
632///
633/// This is the recovery half of the append+apply atomicity guarantee. A writer
634/// that crashed after fsyncing an event row but before fsyncing its projection
635/// (or before advancing `applied_seq`) leaves `applied_seq < last_seq`; the
636/// next lock acquisition heals it here. The reducer is idempotent — every
637/// `*.created` reducer short-circuits when its projection already exists, and
638/// every status/report reducer is a no-op once the target is terminal — so
639/// re-folding an event whose projection *did* land changes nothing. The
640/// manifest's denormalized counters can't desync across this replay either:
641/// they are not folded incrementally but re-derived from projection state by
642/// [`advance_applied_seq`] after each event, so a re-fold simply recomputes the
643/// same totals.
644///
645/// No manifest yet (pre-`run.created`) means there is no watermark to anchor
646/// and nothing durable to catch up, so this returns immediately until the
647/// manifest exists. A legacy manifest reads as `applied_seq = 0` (serde
648/// default), so the first call re-folds the entire log; that is intentional
649/// and safe — see [`crate::schema::Manifest::applied_seq`].
650///
651/// # Corrupt-line tolerance
652///
653/// A line that does not parse as an [`Event`] is skipped, not hard-errored —
654/// the same definition of "corrupt" the quarantine path uses, and the same
655/// tolerance the pre-watermark append path had (it only ever parsed the *last*
656/// line via [`recover_last_seq`]). Bricking every append on an interior poison
657/// line would, among other things, make it impossible to even *record* the
658/// supervisor's `event_log_skipped_line` diagnostic about that very line.
659/// Healing such a line is the supervisor's quarantine job, not the writer's.
660///
661/// A *parse-valid* event whose payload is semantically corrupt is skipped the
662/// same way (with a `warn`), rather than hard-erroring. The dangerous subclass
663/// is an event carrying an embedded id (`child_run_id`, `child_node_id`) that
664/// fails its strict `parse_str` and would
665/// otherwise be joined onto a path — the reducer's independent second line of
666/// defense against a corrupt log, a restored backup, or a future writer that
667/// bypasses the CLI validators (issue `reducer-path-traversal-defense`). Such
668/// an event is a *valid `Event` envelope* (only its `data` is bad), so the
669/// supervisor's [`quarantine_corrupt_lines`] — which only excises lines that
670/// fail the strict envelope parse — can never heal it; hard-erroring here would
671/// brick every future append on that line with no automated recovery path.
672/// Skipping it converges the projection to the largest safe subset and never
673/// joins a tainted id onto a path (the typed-id constructors already make
674/// traversal structurally impossible — a `"../escape"` id never parses into a
675/// [`RunId`](crate::RunId) / [`NodeId`](crate::NodeId), so it can never reach
676/// `nodes/<id>.json`). The append
677/// *gate* stays fail-closed: [`reduce_event_to_ops`] rejects such an event
678/// before it is ever written, so a sanctioned log never reaches this branch and
679/// re-reducing real events on replay is a clean idempotent no-op. A genuine I/O
680/// fault (from the commit or watermark write) still propagates.
681///
682/// Because a sanctioned log is appended in `seq` order under the lock, file
683/// order equals `seq` order for real events; the only out-of-order bytes are
684/// skipped junk, so advancing the watermark to each applied event's `seq` never
685/// jumps over an unfolded real event.
686///
687/// Caller must hold the run's [`RunLock`] and must have already truncated any
688/// torn tail, so the final line is either complete or absent.
689fn replay_unapplied(paths: &RunPaths, events_path: &Path) -> Result<()> {
690    let applied = match read_manifest_opt(paths)? {
691        Some(m) => m.applied_seq,
692        None => return Ok(()),
693    };
694    // Cheap fast path for the overwhelmingly common clean case: the watermark
695    // already covers the log, so there is nothing to replay and no full scan.
696    if applied >= recover_last_seq(events_path)? {
697        return Ok(());
698    }
699    let f = match std::fs::File::open(events_path) {
700        Ok(f) => f,
701        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
702        Err(e) => return Err(Error::io(events_path, e)),
703    };
704    let mut reader = PhysicalLineReader::new(BufReader::new(f));
705    while let Some(line) = reader.next_line().map_err(|e| Error::io(events_path, e))? {
706        // A torn final line is an uncommitted partial write — stop, exactly as
707        // every other reader does.
708        if !line.complete {
709            break;
710        }
711        if line.content.is_empty() {
712            continue;
713        }
714        // Skip a parse-failing line (external junk by the quarantine
715        // definition); apply every event past the watermark in order.
716        let ev: Event = match serde_json::from_slice(line.content) {
717            Ok(ev) => ev,
718            Err(_) => continue,
719        };
720        if ev.seq <= applied {
721            continue;
722        }
723        // Plan the projection writes. A parse-valid but domain-corrupt event —
724        // most dangerously one whose embedded id fails its strict `parse_str`
725        // and would otherwise be joined onto a path — surfaces here as
726        // `CorruptEventLog`. Quarantine cannot excise it (it is a valid
727        // envelope), so we skip it with a warn rather than aborting the whole
728        // catch-up replay; the watermark is not advanced for a skipped event.
729        // See this function's "Corrupt-line tolerance" doc. I/O faults from the
730        // commit/watermark write below still propagate.
731        let ops = match reduce_event_to_ops(paths, &ev) {
732            Ok(ops) => ops,
733            Err(Error::CorruptEventLog { reason, .. }) => {
734                tracing::warn!(
735                    target: "taskfleet_core::events",
736                    path = %events_path.display(),
737                    seq = ev.seq,
738                    kind = %ev.kind,
739                    reason = %reason,
740                    "skipping corrupt event during replay (unsafe id or malformed payload); projection not advanced for it"
741                );
742                continue;
743            }
744            Err(e) => return Err(e),
745        };
746        commit_ops(paths, ops)?;
747        advance_applied_seq(paths, ev.seq)?;
748    }
749    Ok(())
750}
751
752/// Advance `manifest.applied_seq` to `seq` and fsync the manifest (atomic
753/// temp-file + rename), recording that every projection touched by event `seq`
754/// is durably committed.
755///
756/// A no-op when no manifest exists yet, or when the watermark already covers
757/// `seq` — so re-folding an already-applied event (during replay) doesn't churn
758/// the manifest. The reducer for the event may itself have just rewritten the
759/// manifest (e.g. a status transition); reading it back here preserves those
760/// fields while moving only the watermark forward. Caller holds the [`RunLock`].
761///
762/// This is also the single point that persists the manifest's denormalized
763/// `node_count` counter. It is **derived**, not incremented: [`derive_counters`]
764/// recomputes it from the
765/// projection directories — which, because the caller commits an event's
766/// projection ops *before* calling this, already reflect event `seq`. Pinning
767/// the counters to the watermark advance is what makes them undriftable: even
768/// when a crash-replay re-folds an event whose reducer short-circuits to zero
769/// ops (its projection already landed before the crash), this still runs and
770/// re-derives the true counts, healing any counter the old incremental path
771/// would have stranded. See [`derive_counters`] and issue
772/// `manifest-counter-desync`.
773fn advance_applied_seq(paths: &RunPaths, seq: u64) -> Result<()> {
774    if let Some(mut m) = read_manifest_opt(paths)? {
775        if m.applied_seq < seq {
776            let counters = derive_counters(paths)?;
777            m.node_count = counters.node_count;
778            m.applied_seq = seq;
779            write_manifest(paths, &m)?;
780        }
781    }
782    Ok(())
783}
784
785/// One physical line surfaced by [`PhysicalLineReader`]: its content with
786/// any trailing terminator stripped, plus enough framing for the torn-tail
787/// policy (whether it was newline-terminated) and for error context (byte
788/// offset + 1-based line number).
789struct PhysicalLine<'a> {
790    /// Line content with a single trailing terminator (`\n`, optionally
791    /// preceded by `\r`) removed. Interior/leading bytes are untouched.
792    content: &'a [u8],
793    /// `false` only for a final line lacking a trailing `\n` — a torn,
794    /// in-flight append. `true` for every newline-terminated line. Because a
795    /// non-terminated line can only be the last bytes in the file, this is
796    /// `false` for at most one line, and only ever the last one.
797    complete: bool,
798    /// 1-based line number, for `CorruptEventLog` context.
799    lineno: u64,
800}
801
802/// The single physical-line reader behind both [`read_all_events`] and
803/// [`find_prior_with_key`], so the read paths can never disagree about the
804/// torn-tail policy (design.md §1.4; torn-line-policy-consistency).
805///
806/// Bytes are read with [`BufRead::read_until`] (not `read_line`/`lines()`)
807/// for two reasons: it keeps the trailing `\n` so a torn final line is
808/// distinguishable from a newline-terminated interior one, and it reads raw
809/// bytes so a torn tail that cuts a multi-byte UTF-8 sequence is tolerated as
810/// a partial write rather than surfacing as an I/O error. A *newline-
811/// terminated* line with invalid UTF-8 still reaches the caller's parse,
812/// which classifies it as `CorruptEventLog`.
813///
814/// `next_line` lends a slice into an internal buffer, so a caller holds at
815/// most one line at a time — the streaming (lending-iterator) pattern, which
816/// keeps the per-line allocation cost to a single reused buffer.
817struct PhysicalLineReader<R: BufRead> {
818    reader: R,
819    buf: Vec<u8>,
820    lineno: u64,
821    done: bool,
822}
823
824impl<R: BufRead> PhysicalLineReader<R> {
825    fn new(reader: R) -> Self {
826        Self {
827            reader,
828            buf: Vec::new(),
829            lineno: 0,
830            done: false,
831        }
832    }
833
834    /// Yield the next physical line, or `None` at end of file. I/O errors are
835    /// surfaced raw so the caller can attach the log path.
836    fn next_line(&mut self) -> std::io::Result<Option<PhysicalLine<'_>>> {
837        if self.done {
838            return Ok(None);
839        }
840        self.buf.clear();
841        let n = self.reader.read_until(b'\n', &mut self.buf)?;
842        if n == 0 {
843            self.done = true;
844            return Ok(None);
845        }
846        self.lineno += 1;
847        let complete = self.buf.last() == Some(&b'\n');
848        // A non-terminated line is necessarily the final bytes of the file;
849        // stop after handing it back so the torn-tail policy only ever sees
850        // it last.
851        if !complete {
852            self.done = true;
853        }
854        let len = trim_line_end(&self.buf).len();
855        Ok(Some(PhysicalLine {
856            content: &self.buf[..len],
857            complete,
858            lineno: self.lineno,
859        }))
860    }
861}
862
863/// Stream `events.jsonl` line by line, deserializing each complete line into a
864/// caller-chosen envelope probe `T` and invoking `visit(probe, raw_line)`.
865///
866/// This is the streaming counterpart to [`read_all_events`]: it shares the exact
867/// [`PhysicalLineReader`] torn-tail / [`Error::CorruptEventLog`] policy (a torn
868/// final line lacking a trailing `\n` is dropped *without* parsing even if its
869/// bytes are valid JSON; any newline-terminated unparseable line is interior
870/// corruption surfaced as [`Error::CorruptEventLog`]) but never materializes the
871/// whole log — the caller accumulates only what it needs into its own state.
872///
873/// `T` deserializes only the envelope fields it declares; serde ignores the
874/// rest, so a multi-KB `node.report` `data` payload is scanned but never
875/// allocated. The raw line bytes are *lent* to `visit` (a streaming
876/// lending-iterator borrow into the reader's reused buffer), so the closure can
877/// re-parse the full payload for the rare line it must materialize without the
878/// reader holding more than one line at a time.
879///
880/// A missing log is an empty stream (`Ok(())` with no calls). Caller must hold
881/// the run's [`RunLock`]; the scan is read-only over an append-only file.
882pub(crate) fn for_each_event_probe<T, F>(events_path: &Path, mut visit: F) -> Result<()>
883where
884    T: serde::de::DeserializeOwned,
885    F: FnMut(T, &[u8]) -> Result<()>,
886{
887    let f = match std::fs::File::open(events_path) {
888        Ok(f) => f,
889        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
890        Err(e) => return Err(Error::io(events_path, e)),
891    };
892    let mut reader = PhysicalLineReader::new(BufReader::new(f));
893    while let Some(line) = reader.next_line().map_err(|e| Error::io(events_path, e))? {
894        // Torn final line (no trailing newline): uncommitted partial write,
895        // discarded without parsing — mirrors `recover_last_seq`.
896        if !line.complete {
897            break;
898        }
899        if line.content.is_empty() {
900            continue;
901        }
902        let probe: T =
903            serde_json::from_slice(line.content).map_err(|e| Error::CorruptEventLog {
904                path: events_path.to_path_buf(),
905                reason: format!(
906                    "line {} is not a valid event: {} [{e}]",
907                    line.lineno,
908                    excerpt(line.content)
909                ),
910            })?;
911        visit(probe, line.content)?;
912    }
913    Ok(())
914}
915
916/// Read every event from `events.jsonl`. Used by tests and reducer replays.
917///
918/// # Torn-line policy
919///
920/// Built on the shared [`for_each_event_probe`](crate::events) (hence
921/// [`PhysicalLineReader`](crate::events)), so it matches
922/// [`find_prior_with_key`](crate::events) and [`recover_last_seq`] exactly: a
923/// torn final line lacking a trailing `\n` is an in-flight partial write,
924/// dropped *without* parsing even if its bytes happen to be valid JSON. Any
925/// newline-terminated line that fails to parse is interior corruption and
926/// surfaces as [`Error::CorruptEventLog`] — not a transient JSON fault — so a
927/// replay rejects a poisoned log loudly instead of silently dropping a line.
928pub fn read_all_events(events_path: &Path) -> Result<Vec<Event>> {
929    let mut out = Vec::new();
930    for_each_event_probe::<Event, _>(events_path, |ev, _raw| {
931        out.push(ev);
932        Ok(())
933    })?;
934    Ok(out)
935}
936
937/// Outcome of a [`quarantine_corrupt_lines`] call that removed at least one
938/// poison line. `backup_path` is the renamed copy of the original log (kept
939/// verbatim for operator forensics / hand-repair); `removed_byte_offsets`
940/// are the start offsets, in that original, of every newline-terminated line
941/// that failed to parse as an [`Event`] and was excised from the recovered
942/// `events.jsonl`.
943#[derive(Debug, Clone, Serialize)]
944pub struct Quarantine {
945    /// Path to the timestamped `.bak` holding the original poisoned log.
946    pub backup_path: PathBuf,
947    /// Byte offsets (in the original log) of every excised corrupt line.
948    pub removed_byte_offsets: Vec<u64>,
949}
950
951/// Heal a poisoned `events.jsonl` by excising its corrupt physical lines.
952///
953/// P2 made the supervisor *skip* a corrupt JSONL line in memory and keep
954/// tailing, but the bytes stayed on disk forever — so every fresh strict
955/// reader ([`read_all_events`] / a future `rebuild_projections`) still
956/// hard-errors on them, and the skip diagnostic is unreachable to a strict
957/// replay (the corrupt line aborts the read before it). This is the durable
958/// repair: under the run's [`RunLock`], the original log is renamed to
959/// `events.jsonl.corrupt-<ts>.bak` and a recovered `events.jsonl` is written
960/// in its place containing every line *except* the corrupt ones.
961///
962/// "Corrupt" means exactly what the strict readers reject: a
963/// newline-terminated, non-empty line that does not parse as a full [`Event`]
964/// envelope. Empty lines and a torn (newline-less) final line are retained
965/// verbatim — the readers already tolerate both, so excising them would be a
966/// behavior change, not a repair.
967///
968/// Returns `Ok(None)` when the log is missing or already clean (no rename, no
969/// rewrite — the common case is cheap: one read, no corrupt line found).
970/// Returns `Ok(Some(_))` with the backup path and removed offsets when at
971/// least one line was excised. Caller is expected to surface the outcome
972/// (e.g. a `supervisor.event_log_quarantined` diagnostic) and, for a live
973/// tail, restart its read cursor at offset 0 since every byte offset shifts.
974///
975/// `backup_ts` is supplied by the caller (kept out of core so the rename is
976/// deterministic in tests); a filename-safe basic-ISO stamp like
977/// `20260628T120000Z` is the intended form.
978///
979/// # Operator recovery
980///
981/// The excised bytes are never destroyed — they survive verbatim in the
982/// `events.jsonl.corrupt-<ts>.bak` sibling (named by the emitted
983/// `supervisor.event_log_quarantined { backup_path }` diagnostic). To recover
984/// a line the automated repair dropped: open the `.bak`, inspect the line(s)
985/// at the reported `removed_byte_offsets`, hand-fix any salvageable JSON, and —
986/// if you want the record back — stop the run's supervisor, append the
987/// corrected line to the live `events.jsonl` (or replace the file wholesale
988/// from a fixed copy of the backup), then restart the supervisor. The healed
989/// log is the source of truth; projections rebuild from it.
990pub fn quarantine_corrupt_lines(paths: &RunPaths, backup_ts: &str) -> Result<Option<Quarantine>> {
991    RunLock::with_lock(paths, |lock| {
992        quarantine_corrupt_lines_unlocked(lock, paths, backup_ts)
993    })
994}
995
996/// As [`quarantine_corrupt_lines`] but takes a `&LockedRun` witness proving the
997/// caller already holds the run's exclusive [`RunLock`] — the sanctioned
998/// lock-held composition path, mirroring [`append_and_apply_unlocked`].
999/// Re-entering [`quarantine_corrupt_lines`] under a held lock would deadlock on
1000/// the second `flock` open.
1001pub fn quarantine_corrupt_lines_unlocked(
1002    _witness: &LockedRun<'_>,
1003    paths: &RunPaths,
1004    backup_ts: &str,
1005) -> Result<Option<Quarantine>> {
1006    // Guard the run root + event log against symlink redirection before the
1007    // rename/rewrite, exactly as the append path does.
1008    let events_path = paths.checked_events()?;
1009    let raw = match std::fs::read(&events_path) {
1010        Ok(b) => b,
1011        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
1012        Err(e) => return Err(Error::io(&events_path, e)),
1013    };
1014
1015    // Walk physical lines, keeping the raw bytes (terminator included) of every
1016    // retained line so the recovered file is byte-identical save for the
1017    // excised corruption. A line is corrupt iff it is newline-terminated,
1018    // non-empty, and fails the same strict `Event` parse `read_all_events`
1019    // applies — so the recovered log is guaranteed to pass a strict replay.
1020    let mut recovered: Vec<u8> = Vec::with_capacity(raw.len());
1021    let mut removed_byte_offsets: Vec<u64> = Vec::new();
1022    let mut offset: u64 = 0;
1023    let mut i = 0usize;
1024    while i < raw.len() {
1025        let (line_end, complete) = match raw[i..].iter().position(|b| *b == b'\n') {
1026            Some(p) => (i + p + 1, true), // include the trailing '\n'
1027            None => (raw.len(), false),   // torn final line, no '\n'
1028        };
1029        let raw_line = &raw[i..line_end];
1030        let content = trim_line_end(raw_line);
1031        let corrupt =
1032            complete && !content.is_empty() && serde_json::from_slice::<Event>(content).is_err();
1033        if corrupt {
1034            removed_byte_offsets.push(offset);
1035        } else {
1036            recovered.extend_from_slice(raw_line);
1037        }
1038        offset += raw_line.len() as u64;
1039        i = line_end;
1040    }
1041
1042    if removed_byte_offsets.is_empty() {
1043        return Ok(None);
1044    }
1045
1046    // Rename the poisoned log aside (forensics), then atomically drop the
1047    // recovered log in its place. Order matters: the rename frees the path for
1048    // `write_atomic`'s tempfile+rename and preserves the original even if the
1049    // rewrite then fails.
1050    let backup_path = backup_path_for(&events_path, backup_ts);
1051    std::fs::rename(&events_path, &backup_path).map_err(|e| Error::io(&backup_path, e))?;
1052    write_atomic(&events_path, &recovered)?;
1053    Ok(Some(Quarantine {
1054        backup_path,
1055        removed_byte_offsets,
1056    }))
1057}
1058
1059/// Build the `events.jsonl.corrupt-<ts>.bak` sibling path for a quarantine
1060/// backup, preserving the original file name as a prefix.
1061fn backup_path_for(events_path: &Path, ts: &str) -> PathBuf {
1062    let mut name = events_path
1063        .file_name()
1064        .map(std::ffi::OsStr::to_os_string)
1065        .unwrap_or_default();
1066    name.push(format!(".corrupt-{ts}.bak"));
1067    events_path.with_file_name(name)
1068}
1069
1070/// A prior event located by [`find_prior_with_key`](crate::events). Carries enough to let
1071/// an idempotent-retry caller both return the recorded `seq` and verify the
1072/// retry payload matches what was originally written.
1073#[derive(Debug, Clone, PartialEq, Serialize)]
1074pub struct PriorEvent {
1075    /// The recorded `seq` of the matching event.
1076    pub seq: u64,
1077    /// The event's top-level `node_id`, if any.
1078    pub node_id: Option<String>,
1079    /// The event's `data` payload.
1080    pub data: Value,
1081}
1082
1083/// Fields skimmed from every line to test for a match without ever
1084/// allocating the (potentially large) `data` payload. `seq` is optional and
1085/// used only for best-effort error context — it is never a match key, so a
1086/// line missing it must not change whether a `kind` + `idempotency_key`
1087/// match is found.
1088#[derive(Deserialize)]
1089struct ProbeFields {
1090    #[serde(default)]
1091    seq: Option<u64>,
1092    kind: String,
1093    idempotency_key: Option<String>,
1094}
1095
1096/// Fields pulled from the one matching line, including the full payload.
1097#[derive(Deserialize)]
1098struct FullEventForReplay {
1099    seq: u64,
1100    node_id: Option<String>,
1101    data: Value,
1102}
1103
1104/// Maximum number of bytes from a malformed line to surface (escaped) in an
1105/// [`Error::CorruptEventLog`] reason.
1106const CORRUPT_LINE_EXCERPT_BYTES: usize = 100;
1107
1108/// Stream-scan `events.jsonl` for the first event with matching `kind` and
1109/// `idempotency_key`, returning a typed [`PriorEvent`] (or `None` when the
1110/// log is missing or holds no such event).
1111///
1112/// The skim parses each line's envelope (`kind` / `idempotency_key` / `seq`)
1113/// but never materializes `data` for non-matching lines; the full payload
1114/// (`node_id` plus `data`) is deserialized only for the one matching line.
1115/// JSON parsing still scans every byte of every line, so the scan is linear
1116/// in total log bytes under the lock — there is no payload-skipping shortcut.
1117///
1118/// # Torn-line policy
1119///
1120/// [`recover_last_seq`] tolerates a crash-truncated *final* line that lacks
1121/// a trailing newline and discards it regardless of whether its bytes
1122/// happen to form valid JSON. This scanner mirrors that exactly: a final
1123/// line with no trailing `\n` is treated as an in-flight partial write and
1124/// ignored — *before* any parse attempt — so the read (dedup) and write
1125/// (recovery) paths never disagree about whether that tail is committed.
1126///
1127/// Any *interior* line that fails to parse (it is newline-terminated, so a
1128/// later line follows) is a data-integrity fault, so it returns
1129/// [`Error::CorruptEventLog`] rather than silently skipping a line that
1130/// might carry the very key being looked up, which would let the caller
1131/// double-append. This is strictly *more* conservative than
1132/// `recover_last_seq` (which only inspects the last complete line) — a
1133/// deliberate choice for the dedup read.
1134///
1135/// Bytes are read with [`std::io::BufRead::read_until`] rather than
1136/// `read_line` so a torn tail that cuts a multi-byte UTF-8 sequence is
1137/// tolerated as a partial write (matching `recover_last_seq`) instead of
1138/// surfacing as an I/O error; a *newline-terminated* line containing
1139/// invalid UTF-8 is reported as `CorruptEventLog`, not I/O.
1140///
1141/// The `_witness: &LockedRun` is compile-time proof the caller holds the run's
1142/// exclusive [`RunLock`] — the scan is read-only, but it is only meaningful
1143/// fused with an append under one lock window (otherwise a concurrent retry can
1144/// see "no prior event" and double-append). The witness gates the public surface
1145/// so a caller cannot run the scan-then-append race: it must already hold the
1146/// lock to scan, and the same held lock covers the append it threads into
1147/// [`append_and_apply_unlocked`]. [`append_and_apply_idempotent`] fuses the two
1148/// for the common case; a caller that must interleave domain logic between the
1149/// scan and the append (e.g. `discussion resolve`'s already-resolved / no-op
1150/// precedence) calls this primitive directly under its own held lock.
1151pub fn find_prior_with_key(
1152    _witness: &LockedRun<'_>,
1153    paths: &RunPaths,
1154    kind: &str,
1155    idempotency_key: &str,
1156) -> Result<Option<PriorEvent>> {
1157    // Guard the run root + event log before reading: the idempotency scan
1158    // opens `events.jsonl` ahead of the append, so it must refuse a symlinked
1159    // log too rather than read through it.
1160    let events_path = paths.checked_events()?;
1161    let f = match std::fs::File::open(&events_path) {
1162        Ok(f) => f,
1163        Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(None),
1164        Err(e) => return Err(Error::io(&events_path, e)),
1165    };
1166    let mut reader = PhysicalLineReader::new(BufReader::new(f));
1167    // `seq` of the last successfully-parsed line, for best-effort error
1168    // context pointing at where corruption begins.
1169    let mut last_good_seq: u64 = 0;
1170    while let Some(line) = reader.next_line().map_err(|e| Error::io(&events_path, e))? {
1171        // Mirror `recover_last_seq`: a final line lacking a trailing newline
1172        // is an uncommitted partial write, discarded WITHOUT parsing — even
1173        // if its bytes form valid JSON. Parsing it could otherwise return a
1174        // "match" for an event recovery considers unwritten, double-counting
1175        // the seq or skipping a real append.
1176        if !line.complete {
1177            break;
1178        }
1179        if line.content.is_empty() {
1180            continue;
1181        }
1182        let probe: ProbeFields =
1183            serde_json::from_slice(line.content).map_err(|e| Error::CorruptEventLog {
1184                path: events_path.clone(),
1185                reason: format!(
1186                    "line {} is not a valid event envelope (last good seq {last_good_seq}): \
1187                 {} [{e}]",
1188                    line.lineno,
1189                    excerpt(line.content),
1190                ),
1191            })?;
1192        if let Some(seq) = probe.seq {
1193            last_good_seq = seq;
1194        }
1195        if probe.kind != kind || probe.idempotency_key.as_deref() != Some(idempotency_key) {
1196            continue;
1197        }
1198        let full: FullEventForReplay =
1199            serde_json::from_slice(line.content).map_err(|e| Error::CorruptEventLog {
1200                path: events_path.clone(),
1201                reason: format!(
1202                    "line {} matched idempotency key but is not a replayable event: {} [{e}]",
1203                    line.lineno,
1204                    excerpt(line.content),
1205                ),
1206            })?;
1207        return Ok(Some(PriorEvent {
1208            seq: full.seq,
1209            node_id: full.node_id,
1210            data: full.data,
1211        }));
1212    }
1213    Ok(None)
1214}
1215
1216/// Strip a single trailing line terminator (`\n`, optionally preceded by
1217/// `\r`) from a raw line. Unlike `trim_end_matches`, this removes exactly
1218/// one terminator so interior/leading bytes are never altered.
1219fn trim_line_end(buf: &[u8]) -> &[u8] {
1220    let mut end = buf.len();
1221    if end > 0 && buf[end - 1] == b'\n' {
1222        end -= 1;
1223        if end > 0 && buf[end - 1] == b'\r' {
1224            end -= 1;
1225        }
1226    }
1227    &buf[..end]
1228}
1229
1230/// Render a bounded, escaped prefix of a malformed log line for inclusion
1231/// in an error message. Bytes are lossily decoded (a torn multi-byte tail
1232/// becomes the replacement char) and control characters are escaped so an
1233/// excerpt can't inject newlines or ANSI sequences into CLI output.
1234pub(crate) fn excerpt(line: &[u8]) -> String {
1235    let shown = &line[..line.len().min(CORRUPT_LINE_EXCERPT_BYTES)];
1236    let mut out: String = String::from_utf8_lossy(shown).escape_debug().to_string();
1237    if line.len() > CORRUPT_LINE_EXCERPT_BYTES {
1238        out.push('…');
1239    }
1240    out
1241}
1242
1243#[cfg(test)]
1244mod tests {
1245    use super::*;
1246    use crate::RunPaths;
1247    use serde_json::json;
1248    use tempfile::TempDir;
1249
1250    #[test]
1251    fn envelope_run_id_comes_from_paths_not_directory_basename() {
1252        // The whole point of storing run_id: even when the on-disk directory
1253        // name disagrees with the run id (symlinked/non-canonical root, the
1254        // original `root.file_name()` bug), the envelope must carry the stored
1255        // run_id verbatim — never the basename.
1256        let tmp = TempDir::new().unwrap();
1257        let dir = tmp.path().join("not-a-ulid-basename");
1258        std::fs::create_dir_all(&dir).unwrap();
1259        let run_id = "01jxsnap000000000000000000";
1260        let paths = RunPaths::new(dir, run_id).unwrap();
1261
1262        let r = append_and_apply_event(&paths, "run.status", None, None, serde_json::json!({}))
1263            .unwrap();
1264        assert_eq!(r.seq, 1);
1265
1266        let events = read_all_events(&paths.events()).unwrap();
1267        assert_eq!(events.len(), 1);
1268        assert_eq!(events[0].run_id.as_str(), run_id);
1269    }
1270
1271    #[cfg(unix)]
1272    #[test]
1273    fn append_rejects_a_symlinked_event_log() {
1274        // `events.jsonl` is the run's source of truth and highest-leverage
1275        // write — a symlinked log must be refused, not appended through.
1276        use crate::Error;
1277        use std::os::unix::fs::symlink;
1278        let tmp = TempDir::new().unwrap();
1279        let paths = fresh_run(&tmp);
1280        let target = tmp.path().join("evil-events.jsonl");
1281        symlink(&target, paths.events()).unwrap();
1282        let err = append_and_apply_event(&paths, "run.status", None, None, json!({})).unwrap_err();
1283        assert!(
1284            matches!(err, Error::SymlinkStateFile { name: "events", .. }),
1285            "got {err:?}"
1286        );
1287        // The forged append never reached the symlink target.
1288        assert!(!target.exists());
1289    }
1290
1291    /// Build a fresh, empty run directory with a valid `RunPaths` whose
1292    /// `run_id` matches the envelope the reducer will fold.
1293    fn fresh_run(tmp: &TempDir) -> RunPaths {
1294        let run_id = "01jxsnap000000000000000000";
1295        let dir = tmp.path().join(run_id);
1296        std::fs::create_dir_all(&dir).unwrap();
1297        RunPaths::new(dir, run_id).unwrap()
1298    }
1299
1300    /// Parse a `NodeId` for a test append call (the typed envelope id).
1301    fn nid(s: &str) -> NodeId {
1302        NodeId::parse_str(s).unwrap()
1303    }
1304
1305    /// Drive a run to a live node so reducer-affecting events have a target.
1306    fn bootstrap_live_node(paths: &RunPaths) {
1307        append_and_apply_event(
1308            paths,
1309            "run.created",
1310            None,
1311            None,
1312            serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "fix" }),
1313        )
1314        .unwrap();
1315        append_and_apply_event(
1316            paths,
1317            "node.created",
1318            Some(&nid("n-0001")),
1319            None,
1320            serde_json::json!({ "kind": "spinoff" }),
1321        )
1322        .unwrap();
1323    }
1324
1325    #[test]
1326    fn append_and_apply_event_success_path_appends_and_folds() {
1327        let tmp = TempDir::new().unwrap();
1328        let paths = fresh_run(&tmp);
1329
1330        let r = append_and_apply_event(
1331            &paths,
1332            "run.created",
1333            None,
1334            None,
1335            serde_json::json!({ "kind": "spinoff", "lifecycle": "autonomous", "title": "t" }),
1336        )
1337        .unwrap();
1338        assert_eq!(r.seq, 1);
1339        assert!(!r.idempotent_replay);
1340        assert!(r.prior.is_none());
1341
1342        // The reducer ran under the same lock: the manifest projection exists.
1343        let m = crate::read_manifest(&paths).unwrap();
1344        assert_eq!(m.run_id.as_str(), paths.run_id.as_str());
1345    }
1346
1347    #[test]
1348    fn append_and_apply_idempotent_appended_path_returns_fresh_seq() {
1349        let tmp = TempDir::new().unwrap();
1350        let paths = fresh_run(&tmp);
1351        bootstrap_live_node(&paths); // seq 1 run.created, seq 2 node.created
1352
1353        let before = read_all_events(&paths.events()).unwrap().len();
1354        let data = json!({ "status": "running" });
1355        let outcome = RunLock::with_lock(&paths, |lock| {
1356            append_and_apply_idempotent(
1357                &paths,
1358                lock,
1359                "node.status",
1360                Some(&nid("n-0001")),
1361                "k1",
1362                |_seq| Ok(data.clone()),
1363            )
1364        })
1365        .unwrap();
1366        match outcome {
1367            AppendOutcome::Appended { seq } => {
1368                assert_eq!(seq, 3, "fresh append takes the next seq");
1369            }
1370            other => panic!("expected Appended, got {other:?}"),
1371        }
1372        assert_eq!(
1373            read_all_events(&paths.events()).unwrap().len(),
1374            before + 1,
1375            "a fresh key appends exactly one event"
1376        );
1377    }
1378
1379    #[test]
1380    fn append_and_apply_idempotent_replay_returns_prior_without_appending() {
1381        let tmp = TempDir::new().unwrap();
1382        let paths = fresh_run(&tmp);
1383        bootstrap_live_node(&paths);
1384        let node = nid("n-0001");
1385        let data = json!({ "status": "running" });
1386
1387        let first = RunLock::with_lock(&paths, |lock| {
1388            append_and_apply_idempotent(&paths, lock, "node.status", Some(&node), "k1", |_seq| {
1389                Ok(data.clone())
1390            })
1391        })
1392        .unwrap();
1393        let first_seq = match first {
1394            AppendOutcome::Appended { seq } => seq,
1395            other => panic!("expected Appended, got {other:?}"),
1396        };
1397        let after_first = read_all_events(&paths.events()).unwrap().len();
1398
1399        // Same kind + key + node + data → a true replay: nothing appended, the
1400        // prior event (its seq + data) is returned.
1401        let replay = RunLock::with_lock(&paths, |lock| {
1402            append_and_apply_idempotent(&paths, lock, "node.status", Some(&node), "k1", |_seq| {
1403                Ok(data.clone())
1404            })
1405        })
1406        .unwrap();
1407        match replay {
1408            AppendOutcome::IdempotentReplay { prior } => {
1409                assert_eq!(prior.seq, first_seq);
1410                assert_eq!(prior.node_id.as_deref(), Some("n-0001"));
1411                assert_eq!(prior.data, data);
1412            }
1413            other => panic!("expected IdempotentReplay, got {other:?}"),
1414        }
1415        assert_eq!(
1416            read_all_events(&paths.events()).unwrap().len(),
1417            after_first,
1418            "a replay must not append a new event"
1419        );
1420    }
1421
1422    #[test]
1423    fn append_and_apply_idempotent_conflict_on_different_data() {
1424        let tmp = TempDir::new().unwrap();
1425        let paths = fresh_run(&tmp);
1426        bootstrap_live_node(&paths);
1427        let node = nid("n-0001");
1428
1429        let first = RunLock::with_lock(&paths, |lock| {
1430            append_and_apply_idempotent(&paths, lock, "node.status", Some(&node), "k1", |_seq| {
1431                Ok(json!({ "status": "running" }))
1432            })
1433        })
1434        .unwrap();
1435        let first_seq = match first {
1436            AppendOutcome::Appended { seq } => seq,
1437            other => panic!("expected Appended, got {other:?}"),
1438        };
1439        let after_first = read_all_events(&paths.events()).unwrap().len();
1440
1441        // Same key, DIFFERENT payload → conflict, carrying the prior event's seq;
1442        // nothing new is appended.
1443        let conflict = RunLock::with_lock(&paths, |lock| {
1444            append_and_apply_idempotent(&paths, lock, "node.status", Some(&node), "k1", |_seq| {
1445                Ok(json!({ "status": "done" }))
1446            })
1447        })
1448        .unwrap();
1449        match conflict {
1450            AppendOutcome::Conflict { prior } => {
1451                assert_eq!(prior.seq, first_seq);
1452                assert_eq!(prior.data, json!({ "status": "running" }));
1453            }
1454            other => panic!("expected Conflict, got {other:?}"),
1455        }
1456        assert_eq!(
1457            read_all_events(&paths.events()).unwrap().len(),
1458            after_first,
1459            "a conflict must not append a new event"
1460        );
1461    }
1462
1463    #[test]
1464    fn append_and_apply_idempotent_conflict_on_different_node_id() {
1465        // Same key + same data but a different envelope node is still a reused
1466        // key for a different request → conflict, not a silent replay.
1467        let tmp = TempDir::new().unwrap();
1468        let paths = fresh_run(&tmp);
1469        bootstrap_live_node(&paths);
1470        // A second live node so the conflicting append targets a real node.
1471        append_and_apply_event(
1472            &paths,
1473            "node.created",
1474            Some(&nid("n-0002")),
1475            None,
1476            json!({ "kind": "spinoff" }),
1477        )
1478        .unwrap();
1479        let data = json!({ "status": "running" });
1480
1481        RunLock::with_lock(&paths, |lock| {
1482            append_and_apply_idempotent(
1483                &paths,
1484                lock,
1485                "node.status",
1486                Some(&nid("n-0001")),
1487                "k1",
1488                |_seq| Ok(data.clone()),
1489            )
1490        })
1491        .unwrap();
1492
1493        let conflict = RunLock::with_lock(&paths, |lock| {
1494            append_and_apply_idempotent(
1495                &paths,
1496                lock,
1497                "node.status",
1498                Some(&nid("n-0002")),
1499                "k1",
1500                |_seq| Ok(data.clone()),
1501            )
1502        })
1503        .unwrap();
1504        assert!(
1505            matches!(conflict, AppendOutcome::Conflict { prior } if prior.node_id.as_deref() == Some("n-0001")),
1506            "a node-id mismatch under the same key is a conflict"
1507        );
1508    }
1509
1510    #[test]
1511    fn append_and_apply_idempotent_rejects_empty_key() {
1512        let tmp = TempDir::new().unwrap();
1513        let paths = fresh_run(&tmp);
1514        bootstrap_live_node(&paths);
1515        let err = RunLock::with_lock(&paths, |lock| {
1516            append_and_apply_idempotent(
1517                &paths,
1518                lock,
1519                "node.status",
1520                Some(&nid("n-0001")),
1521                "",
1522                |_seq| Ok(json!({ "status": "running" })),
1523            )
1524        })
1525        .unwrap_err();
1526        assert!(matches!(err, Error::EmptyIdempotencyKey), "got {err:?}");
1527    }
1528
1529    #[test]
1530    fn append_and_apply_event_idempotent_replay_returns_prior_without_appending() {
1531        let tmp = TempDir::new().unwrap();
1532        let paths = fresh_run(&tmp);
1533        bootstrap_live_node(&paths);
1534
1535        let data = serde_json::json!({ "status": "running" });
1536        let first = append_and_apply_event(
1537            &paths,
1538            "node.status",
1539            Some(&nid("n-0001")),
1540            Some("k1"),
1541            data.clone(),
1542        )
1543        .unwrap();
1544        assert!(!first.idempotent_replay);
1545        let before = read_all_events(&paths.events()).unwrap().len();
1546
1547        // Same kind + key: a replay returns the prior event and appends nothing.
1548        let replay = append_and_apply_event(
1549            &paths,
1550            "node.status",
1551            Some(&nid("n-0001")),
1552            Some("k1"),
1553            data.clone(),
1554        )
1555        .unwrap();
1556        assert!(replay.idempotent_replay);
1557        assert!(
1558            !replay.applied,
1559            "an idempotent replay applies nothing this call (applied: false)"
1560        );
1561        assert_eq!(replay.seq, first.seq);
1562        let prior = replay.prior.expect("replay carries the prior event");
1563        assert_eq!(prior.node_id.as_deref(), Some("n-0001"));
1564        assert_eq!(prior.data, data);
1565        assert_eq!(
1566            read_all_events(&paths.events()).unwrap().len(),
1567            before,
1568            "replay must not append a new line"
1569        );
1570    }
1571
1572    #[test]
1573    fn append_and_apply_event_reducer_noop_is_still_a_success() {
1574        let tmp = TempDir::new().unwrap();
1575        let paths = fresh_run(&tmp);
1576        bootstrap_live_node(&paths);
1577
1578        // Settle the node terminal. A real state change → `applied: true`.
1579        let n0001 = nid("n-0001");
1580        let settle = append_and_apply_event(
1581            &paths,
1582            "node.report",
1583            Some(&n0001),
1584            None,
1585            serde_json::json!({ "success": true }),
1586        )
1587        .unwrap();
1588        assert!(
1589            settle.applied,
1590            "a report that terminalizes a live node applied a projection op"
1591        );
1592        assert_eq!(
1593            crate::read_node(&paths, &n0001).unwrap().status,
1594            crate::schema::Status::Done
1595        );
1596
1597        // A later status event is dropped by the terminal-state guard, but the
1598        // append still happened: the result names the appended event's seq and
1599        // is not a replay. The node stays Done. `applied` is FALSE — the reducer
1600        // planned zero ops (issue `reducer-adopt-explicit-merge`).
1601        let before = read_all_events(&paths.events()).unwrap().len();
1602        let r = append_and_apply_event(
1603            &paths,
1604            "node.status",
1605            Some(&n0001),
1606            None,
1607            serde_json::json!({ "status": "running" }),
1608        )
1609        .unwrap();
1610        assert!(!r.idempotent_replay);
1611        assert!(
1612            !r.applied,
1613            "a dead event dropped by the terminal guard reports applied: false"
1614        );
1615        assert_eq!(r.seq as usize, before + 1);
1616        assert_eq!(
1617            read_all_events(&paths.events()).unwrap().len(),
1618            before + 1,
1619            "the event is appended even when the reducer no-ops"
1620        );
1621        assert_eq!(
1622            crate::read_node(&paths, &n0001).unwrap().status,
1623            crate::schema::Status::Done,
1624            "terminal status is frozen"
1625        );
1626    }
1627
1628    #[test]
1629    fn bootstrap_advances_the_watermark_past_every_appended_event() {
1630        // Baseline for the replay tests: the normal append path keeps the
1631        // watermark pinned to the last appended seq, so `applied_seq == last`
1632        // whenever the log is clean.
1633        let tmp = TempDir::new().unwrap();
1634        let paths = fresh_run(&tmp);
1635        bootstrap_live_node(&paths); // seq 1 run.created, seq 2 node.created
1636        assert_eq!(
1637            crate::read_manifest(&paths).unwrap().applied_seq,
1638            2,
1639            "watermark tracks the last appended event"
1640        );
1641    }
1642
1643    #[test]
1644    fn append_replays_unapplied_tail_before_appending() {
1645        use crate::schema::Status;
1646        // Failure scenario 1: a reducer crash after the event-row fsync but
1647        // before the projection/watermark write leaves the log ahead of the
1648        // projections. The next lock acquisition must replay that tail.
1649        let tmp = TempDir::new().unwrap();
1650        let paths = fresh_run(&tmp);
1651        bootstrap_live_node(&paths); // applied_seq == 2, node n-0001 Pending
1652        let n0001 = nid("n-0001");
1653
1654        // Append a tail event (seq 3) WITHOUT running the reducer — exactly the
1655        // on-disk state a crash between the row fsync and the projection write
1656        // would leave behind. The raw append still needs the witness (lock held).
1657        RunLock::with_lock(&paths, |lock| {
1658            append_event_with_seq(
1659                lock,
1660                &paths,
1661                3,
1662                "node.status",
1663                Some(&n0001),
1664                None,
1665                json!({ "status": "running" }),
1666            )
1667        })
1668        .unwrap();
1669        assert_eq!(
1670            crate::read_node(&paths, &n0001).unwrap().status,
1671            Status::Pending,
1672            "the tail event's projection has not landed yet"
1673        );
1674        assert_eq!(crate::read_manifest(&paths).unwrap().applied_seq, 2);
1675
1676        // Any new append acquires the lock and replays seq 3 first, so the new
1677        // event takes seq 4 and the stale projection is healed.
1678        let r = append_and_apply_event(
1679            &paths,
1680            "run.status",
1681            None,
1682            None,
1683            json!({ "status": "running" }),
1684        )
1685        .unwrap();
1686        assert_eq!(r.seq, 4, "the new event follows the replayed tail");
1687        assert_eq!(
1688            crate::read_node(&paths, &n0001).unwrap().status,
1689            Status::Running,
1690            "the previously-unapplied tail event is now folded"
1691        );
1692        assert_eq!(
1693            crate::read_manifest(&paths).unwrap().applied_seq,
1694            4,
1695            "the watermark now covers the whole log"
1696        );
1697    }
1698
1699    #[test]
1700    fn legacy_manifest_without_applied_seq_migrates_on_next_write() {
1701        use crate::schema::Status;
1702        // A `manifest.json` written before `applied_seq` existed must read back
1703        // as 0 (serde default) and self-migrate on the next write via an
1704        // idempotent full replay — without double-counting counters or
1705        // resurrecting a terminal node (failure scenario 2's no-double-count
1706        // guarantee, exercised over the whole log).
1707        let tmp = TempDir::new().unwrap();
1708        let paths = fresh_run(&tmp);
1709        bootstrap_live_node(&paths);
1710        let n0001 = nid("n-0001");
1711        append_and_apply_event(
1712            &paths,
1713            "node.report",
1714            Some(&n0001),
1715            None,
1716            json!({ "success": true }),
1717        )
1718        .unwrap(); // seq 3 → node Done, applied_seq == 3, node_count == 1
1719
1720        // Rewrite the manifest WITHOUT an `applied_seq` field, mimicking a
1721        // pre-watermark binary's output.
1722        let mut mv: serde_json::Value =
1723            serde_json::from_slice(&std::fs::read(paths.manifest()).unwrap()).unwrap();
1724        assert!(mv.as_object_mut().unwrap().remove("applied_seq").is_some());
1725        std::fs::write(paths.manifest(), serde_json::to_vec_pretty(&mv).unwrap()).unwrap();
1726        assert_eq!(
1727            crate::read_manifest(&paths).unwrap().applied_seq,
1728            0,
1729            "a legacy manifest reads as applied_seq 0"
1730        );
1731
1732        // The next write triggers a full idempotent replay of seq 1..=3 (all
1733        // no-ops) and advances the watermark to last_seq.
1734        append_and_apply_event(
1735            &paths,
1736            "run.status",
1737            None,
1738            None,
1739            json!({ "status": "running" }),
1740        )
1741        .unwrap(); // seq 4
1742        let m = crate::read_manifest(&paths).unwrap();
1743        assert_eq!(m.applied_seq, 4, "watermark caught up to the log");
1744        assert_eq!(
1745            m.node_count, 1,
1746            "full replay did not double-count node_count"
1747        );
1748        assert_eq!(
1749            crate::read_node(&paths, &n0001).unwrap().status,
1750            Status::Done,
1751            "replaying its history did not resurrect the terminal node"
1752        );
1753    }
1754
1755    #[test]
1756    fn replay_skips_events_with_unsafe_ids_and_never_escapes_run_dir() {
1757        // Issue `reducer-path-traversal-defense`: the reducer must independently
1758        // defend against ids read from `events.jsonl` that bypass the CLI
1759        // validators — a corrupt log, a restored backup, or a future writer.
1760        // We craft a log straight onto disk (skipping the append gate) holding
1761        // two poison `child.spawned` lines (a traversal-laden and an empty
1762        // `child_run_id`) and one good one, then drive a catch-up replay and
1763        // assert: the poison events are skipped (not fatal) and the good event
1764        // still applies (its child ref lands on the parent node).
1765        let tmp = TempDir::new().unwrap();
1766        let paths = fresh_run(&tmp);
1767        bootstrap_live_node(&paths); // applied_seq == 2, node n-0001 live
1768
1769        // seq 3 — a traversal-laden `child_run_id`; seq 4 — an empty one. Both
1770        // fail their strict `parse_str`, so `reduce_event_to_ops` rejects them.
1771        // seq 5 — a well-formed pair that must be applied despite the poison
1772        // lines preceding it. All three raw appends share one held lock.
1773        let child_run = "02jxsnap000000000000000000";
1774        RunLock::with_lock(&paths, |lock| {
1775            append_event_with_seq(
1776                lock,
1777                &paths,
1778                3,
1779                "child.spawned",
1780                Some(&nid("n-0001")),
1781                None,
1782                json!({ "child_run_id": "../escape", "child_node_id": "n-0001" }),
1783            )?;
1784            append_event_with_seq(
1785                lock,
1786                &paths,
1787                4,
1788                "child.spawned",
1789                Some(&nid("n-0001")),
1790                None,
1791                json!({ "child_run_id": "", "child_node_id": "n-0001" }),
1792            )?;
1793            append_event_with_seq(
1794                lock,
1795                &paths,
1796                5,
1797                "child.spawned",
1798                Some(&nid("n-0001")),
1799                None,
1800                json!({ "child_run_id": child_run, "child_node_id": "n-0001" }),
1801            )
1802        })
1803        .unwrap();
1804
1805        // The poison lines must NOT abort the replay (the regression this fixes:
1806        // a `..`-laden id is a valid envelope quarantine can't excise, so a hard
1807        // error here would brick every future append on the run).
1808        replay_unapplied(&paths, &paths.events()).expect("poison lines skipped, not fatal");
1809
1810        // The good child.spawned landed: the parent node carries exactly one
1811        // child ref, and the two poison ids added nothing.
1812        let parent = crate::read_node(&paths, &nid("n-0001")).unwrap();
1813        assert_eq!(
1814            parent.children.len(),
1815            1,
1816            "only the good child ref was applied; poison ids added none"
1817        );
1818        assert_eq!(parent.children[0].run_id.as_str(), child_run);
1819        assert!(
1820            !paths.root.join("escape").exists(),
1821            "traversal id must never have been joined onto a path"
1822        );
1823
1824        // The watermark jumped past the skipped seqs to the applied good event.
1825        let m = crate::read_manifest(&paths).unwrap();
1826        assert_eq!(
1827            m.applied_seq, 5,
1828            "watermark advanced past the skipped poison"
1829        );
1830    }
1831
1832    #[test]
1833    fn node_count_desync_heals_on_replay() {
1834        // Faithful reproduction of issue `manifest-counter-desync`: a crash left
1835        // the node projection on disk but lost the follow-on manifest write (the
1836        // counter bump + watermark advance). Before the fix, the replay
1837        // short-circuited on the already-existing node and the stale counter
1838        // stuck forever; now the counter is re-derived at the watermark advance.
1839        let tmp = TempDir::new().unwrap();
1840        let paths = fresh_run(&tmp);
1841        bootstrap_live_node(&paths); // node n-0001 on disk, node_count == 1, applied_seq == 2
1842
1843        // Rewind the manifest to the exact mid-crash state: the node file
1844        // exists, but the manifest still shows the pre-node counter and a
1845        // watermark that sits before the `node.created` at seq 2.
1846        let mut m = crate::read_manifest(&paths).unwrap();
1847        assert_eq!(m.node_count, 1, "precondition: bootstrap counted the node");
1848        m.node_count = 0;
1849        m.applied_seq = 1;
1850        write_manifest(&paths, &m).unwrap();
1851
1852        // The next append acquires the lock, replays seq 2 (node already exists,
1853        // so the reducer plans zero ops), and re-derives the counter when it
1854        // advances the watermark past seq 2.
1855        append_and_apply_event(
1856            &paths,
1857            "run.status",
1858            None,
1859            None,
1860            json!({ "status": "running" }),
1861        )
1862        .unwrap();
1863
1864        let healed = crate::read_manifest(&paths).unwrap();
1865        assert_eq!(
1866            healed.node_count, 1,
1867            "node_count converged to the true projection count"
1868        );
1869        assert!(healed.applied_seq >= 2, "watermark caught up past the node");
1870    }
1871
1872    #[test]
1873    fn full_replay_does_not_double_count_any_counter() {
1874        // Idempotence across a full from-scratch replay: re-folding every event
1875        // must re-derive the same total, never accumulate. Covers `node_count`.
1876        let tmp = TempDir::new().unwrap();
1877        let paths = fresh_run(&tmp);
1878        bootstrap_live_node(&paths);
1879        append_and_apply_event(
1880            &paths,
1881            "node.created",
1882            Some(&nid("n-0002")),
1883            None,
1884            json!({ "kind": "spinoff" }),
1885        )
1886        .unwrap();
1887        let before = crate::read_manifest(&paths).unwrap();
1888        assert_eq!(before.node_count, 2, "precondition: two nodes");
1889
1890        // Reset the watermark to force a full idempotent replay of the whole log
1891        // on the next append (the legacy-migration path), and deliberately
1892        // corrupt the counter so a heal is observable.
1893        let mut m = before;
1894        m.applied_seq = 0;
1895        m.node_count = 99;
1896        write_manifest(&paths, &m).unwrap();
1897        append_and_apply_event(
1898            &paths,
1899            "run.status",
1900            None,
1901            None,
1902            json!({ "status": "running" }),
1903        )
1904        .unwrap();
1905
1906        let after = crate::read_manifest(&paths).unwrap();
1907        assert_eq!(
1908            after.node_count, 2,
1909            "counter re-derived to the true total — no double-count across full replay"
1910        );
1911    }
1912
1913    #[test]
1914    fn idempotent_replay_catches_up_projection_before_returning() {
1915        use crate::projections::write_manifest;
1916        use crate::schema::Status;
1917        use crate::write_node;
1918        // Requirement 3: an idempotency-key replay must ensure the projection is
1919        // caught up (`applied_seq >= prior.seq`) before returning the prior
1920        // envelope — never a "found, but not yet applied" result.
1921        let tmp = TempDir::new().unwrap();
1922        let paths = fresh_run(&tmp);
1923        bootstrap_live_node(&paths);
1924        let n0001 = nid("n-0001");
1925
1926        // A keyed event lands and folds normally...
1927        let first = append_and_apply_event(
1928            &paths,
1929            "node.status",
1930            Some(&n0001),
1931            Some("k1"),
1932            json!({ "status": "running" }),
1933        )
1934        .unwrap(); // seq 3
1935        assert!(!first.idempotent_replay);
1936
1937        // ...then simulate a crash that lost the fold: rewind the watermark
1938        // below seq 3 and revert the node to its pre-event Pending state.
1939        let mut m = crate::read_manifest(&paths).unwrap();
1940        m.applied_seq = 2;
1941        write_manifest(&paths, &m).unwrap();
1942        let mut n = crate::read_node(&paths, &n0001).unwrap();
1943        n.status = Status::Pending;
1944        write_node(&paths, &n).unwrap();
1945
1946        // The idempotent retry returns the prior seq AND catches the projection
1947        // up first.
1948        let replay = append_and_apply_event(
1949            &paths,
1950            "node.status",
1951            Some(&n0001),
1952            Some("k1"),
1953            json!({ "status": "running" }),
1954        )
1955        .unwrap();
1956        assert!(replay.idempotent_replay);
1957        assert_eq!(replay.seq, first.seq);
1958        assert!(
1959            crate::read_manifest(&paths).unwrap().applied_seq >= first.seq,
1960            "watermark caught up before the replay returned"
1961        );
1962        assert_eq!(
1963            crate::read_node(&paths, &n0001).unwrap().status,
1964            Status::Running,
1965            "the prior event's projection is durable before returning"
1966        );
1967    }
1968
1969    /// Build a `RunPaths` over a fresh tempdir and write `bytes` verbatim to
1970    /// `events.jsonl` — verbatim so a test can craft torn-line boundaries
1971    /// (a missing trailing `\n`) that the append path never produces.
1972    fn paths_with_events(tmp: &TempDir, bytes: &[u8]) -> RunPaths {
1973        let dir = tmp.path().join("run");
1974        std::fs::create_dir_all(&dir).unwrap();
1975        let paths = RunPaths::new(dir, "01jxsnap000000000000000000").unwrap();
1976        std::fs::write(paths.events(), bytes).unwrap();
1977        paths
1978    }
1979
1980    /// Run [`find_prior_with_key`] under a freshly-acquired exclusive lock —
1981    /// the witness it now requires. The scan is read-only, so taking the lock
1982    /// just to mint the witness is exactly what a real caller does.
1983    fn scan(paths: &RunPaths, kind: &str, key: &str) -> Result<Option<PriorEvent>> {
1984        RunLock::with_lock(paths, |w| find_prior_with_key(w, paths, kind, key))
1985    }
1986
1987    #[test]
1988    fn find_prior_with_key_missing_log_is_none() {
1989        let tmp = TempDir::new().unwrap();
1990        let dir = tmp.path().join("run");
1991        std::fs::create_dir_all(&dir).unwrap();
1992        let paths = RunPaths::new(dir, "01jxsnap000000000000000000").unwrap();
1993        // No events.jsonl written at all.
1994        let got = scan(&paths, "node.report", "k1").unwrap();
1995        assert!(got.is_none());
1996    }
1997
1998    #[test]
1999    fn find_prior_with_key_finds_the_matching_line() {
2000        let tmp = TempDir::new().unwrap();
2001        let log = concat!(
2002            r#"{"seq":1,"kind":"node.status","idempotency_key":"k0","node_id":"n-1","data":{}}"#,
2003            "\n",
2004            r#"{"seq":2,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{"ok":true}}"#,
2005            "\n",
2006        );
2007        let paths = paths_with_events(&tmp, log.as_bytes());
2008        let got = scan(&paths, "node.report", "k1").unwrap().expect("match");
2009        assert_eq!(got.seq, 2);
2010        assert_eq!(got.node_id.as_deref(), Some("n-1"));
2011        assert_eq!(got.data, serde_json::json!({"ok": true}));
2012    }
2013
2014    #[test]
2015    fn find_prior_with_key_no_match_is_none() {
2016        let tmp = TempDir::new().unwrap();
2017        let log = concat!(
2018            r#"{"seq":1,"kind":"node.report","idempotency_key":"other","node_id":"n-1","data":{}}"#,
2019            "\n",
2020        );
2021        let paths = paths_with_events(&tmp, log.as_bytes());
2022        assert!(scan(&paths, "node.report", "k1").unwrap().is_none());
2023    }
2024
2025    #[test]
2026    fn find_prior_with_key_tolerates_torn_final_line() {
2027        // A complete record, then a crash-truncated final line with NO
2028        // trailing newline — exactly what `recover_last_seq` tolerates.
2029        // The scan must still return the earlier match and never error.
2030        let tmp = TempDir::new().unwrap();
2031        let mut log = String::new();
2032        log.push_str(
2033            r#"{"seq":1,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{"ok":true}}"#,
2034        );
2035        log.push('\n');
2036        log.push_str(r#"{"seq":2,"kind":"node.rep"#); // torn mid-write, no newline
2037        let paths = paths_with_events(&tmp, log.as_bytes());
2038
2039        let got = scan(&paths, "node.report", "k1")
2040            .unwrap()
2041            .expect("match before the torn tail");
2042        assert_eq!(got.seq, 1);
2043
2044        // A torn final line with no matching key ahead of it returns None,
2045        // not an error.
2046        let tmp2 = TempDir::new().unwrap();
2047        let paths2 = paths_with_events(&tmp2, br#"{"seq":1,"kind":"node.rep"#);
2048        assert!(scan(&paths2, "node.report", "k1").unwrap().is_none());
2049    }
2050
2051    #[test]
2052    fn find_prior_with_key_ignores_valid_json_final_line_without_newline() {
2053        // The dangerous case: a crash landed a COMPLETE, valid-JSON event
2054        // but the trailing newline never flushed. `recover_last_seq`
2055        // discards any newline-less tail, so it considers this event
2056        // unwritten (returns 0). The dedup scan MUST agree and return None
2057        // — otherwise it would report "already appended", the caller skips
2058        // the append, and the event is lost / the seq double-counts.
2059        let tmp = TempDir::new().unwrap();
2060        let line =
2061            br#"{"seq":1,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{}}"#;
2062        let paths = paths_with_events(&tmp, line);
2063        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 0);
2064        assert!(
2065            scan(&paths, "node.report", "k1").unwrap().is_none(),
2066            "torn tail must be ignored even when it parses as valid JSON"
2067        );
2068    }
2069
2070    #[test]
2071    fn find_prior_with_key_skips_nonmatching_line_missing_seq() {
2072        // `seq` is not a match key, so a NON-matching envelope that happens
2073        // to lack `seq` must be skimmed past, not treated as corruption that
2074        // aborts the scan before a later match. (The pre-lift scanner's
2075        // probe didn't require `seq`; making it required would have been a
2076        // regression that hid a real key behind an unrelated seq-less line.)
2077        let tmp = TempDir::new().unwrap();
2078        let log = concat!(
2079            r#"{"kind":"node.status","idempotency_key":"other","node_id":"n-1","data":{}}"#,
2080            "\n",
2081            r#"{"seq":2,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{"ok":true}}"#,
2082            "\n",
2083        );
2084        let paths = paths_with_events(&tmp, log.as_bytes());
2085        let got = scan(&paths, "node.report", "k1")
2086            .unwrap()
2087            .expect("match after a seq-less non-matching line");
2088        assert_eq!(got.seq, 2);
2089        assert_eq!(got.node_id.as_deref(), Some("n-1"));
2090    }
2091
2092    #[test]
2093    fn find_prior_with_key_matched_line_bad_payload_is_corrupt_log() {
2094        // A line that skims fine (kind + key match) but whose full payload
2095        // is malformed (`node_id` is a number, not a string) is event-log
2096        // corruption — it must surface as CorruptEventLog (exit 1), not a
2097        // generic JSON/io error (exit 2).
2098        let tmp = TempDir::new().unwrap();
2099        let log = concat!(
2100            r#"{"seq":1,"kind":"node.report","idempotency_key":"k1","node_id":42,"data":{}}"#,
2101            "\n",
2102        );
2103        let paths = paths_with_events(&tmp, log.as_bytes());
2104        let err = scan(&paths, "node.report", "k1").unwrap_err();
2105        assert!(
2106            matches!(err, Error::CorruptEventLog { .. }),
2107            "expected CorruptEventLog, got {err:?}"
2108        );
2109    }
2110
2111    #[test]
2112    fn find_prior_with_key_handles_crlf_line_endings() {
2113        let tmp = TempDir::new().unwrap();
2114        let log = concat!(
2115            r#"{"seq":1,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{}}"#,
2116            "\r\n",
2117        );
2118        let paths = paths_with_events(&tmp, log.as_bytes());
2119        let got = scan(&paths, "node.report", "k1")
2120            .unwrap()
2121            .expect("CRLF-terminated match");
2122        assert_eq!(got.seq, 1);
2123    }
2124
2125    #[test]
2126    fn find_prior_with_key_tolerates_partial_utf8_torn_tail() {
2127        // A crash can cut a multi-byte UTF-8 sequence mid-character. With
2128        // byte-oriented reading this torn (newline-less) tail is tolerated
2129        // like any other partial write, not surfaced as an I/O error.
2130        let tmp = TempDir::new().unwrap();
2131        let mut log = Vec::new();
2132        log.extend_from_slice(
2133            br#"{"seq":1,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{}}"#,
2134        );
2135        log.push(b'\n');
2136        log.extend_from_slice(&[0xF0, 0x9F]); // start of a 4-byte char, truncated
2137        let paths = paths_with_events(&tmp, &log);
2138        let got = scan(&paths, "node.report", "k1")
2139            .unwrap()
2140            .expect("match before the partial-UTF8 tail");
2141        assert_eq!(got.seq, 1);
2142    }
2143
2144    #[test]
2145    fn recover_last_seq_newline_terminated_garbage_is_corrupt_log() {
2146        // Consistency guard with find_prior_with_key: a newline-terminated
2147        // final line that isn't valid JSON is CorruptEventLog from BOTH
2148        // readers, so the CLI maps both to the same corrupt-event-log exit.
2149        let tmp = TempDir::new().unwrap();
2150        let paths = paths_with_events(&tmp, b"{not json at all\n");
2151        let err = recover_last_seq(&paths.events()).unwrap_err();
2152        assert!(
2153            matches!(err, Error::CorruptEventLog { .. }),
2154            "expected CorruptEventLog, got {err:?}"
2155        );
2156    }
2157
2158    #[test]
2159    fn rejected_event_is_not_appended() {
2160        // The transactional fix: a reducer-rejected event must error BEFORE
2161        // any durable write, so events.jsonl never gains a poison line.
2162        let tmp = TempDir::new().unwrap();
2163        let paths = fresh_run(&tmp);
2164        bootstrap_live_node(&paths);
2165        let before = read_all_events(&paths.events()).unwrap().len();
2166
2167        // `node.report` with neither success nor cancelled → reducer rejects.
2168        let err =
2169            append_and_apply_event(&paths, "node.report", Some(&nid("n-0001")), None, json!({}))
2170                .unwrap_err();
2171        assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
2172
2173        assert_eq!(
2174            read_all_events(&paths.events()).unwrap().len(),
2175            before,
2176            "a rejected event must not be appended"
2177        );
2178        // The log is still clean and re-readable (no poison line stranded it).
2179        assert!(recover_last_seq(&paths.events()).is_ok());
2180        let next = append_and_apply_event(
2181            &paths,
2182            "node.report",
2183            Some(&nid("n-0001")),
2184            None,
2185            json!({ "success": true }),
2186        )
2187        .unwrap();
2188        assert_eq!(
2189            next.seq as usize,
2190            before + 1,
2191            "the next valid append reuses the seq the rejected event never consumed"
2192        );
2193    }
2194
2195    #[test]
2196    fn validate_event_agrees_with_apply_event() {
2197        // Drift guard: `validate_event` (the pre-append gate) must return Err
2198        // in EXACTLY the cases `apply_event` would, for the same state — else
2199        // it would refuse a harmless no-op or let a poison line through.
2200        use crate::reducer::{apply_event, validate_event};
2201
2202        fn ev(paths: &RunPaths, kind: &str, node_id: Option<&str>, data: Value) -> Event {
2203            Event {
2204                ts: Utc::now(),
2205                seq: 999,
2206                kind: kind.to_string(),
2207                run_id: paths.run_id.clone(),
2208                node_id: node_id.map(|s| crate::schema::NodeId::parse_str(s).unwrap()),
2209                idempotency_key: None,
2210                data,
2211            }
2212        }
2213        // validate is read-only, so running it first leaves apply's pre-state
2214        // intact; we compare the two verdicts on the same fresh run.
2215        fn agree(paths: &RunPaths, e: &Event, label: &str) {
2216            let v = validate_event(paths, e).is_err();
2217            let a = apply_event(paths, e).is_err();
2218            assert_eq!(v, a, "{label}: validate_err={v} apply_err={a}");
2219        }
2220
2221        // Live node: bad report rejected; good report accepted; missing
2222        // node_id rejected; bad status rejected.
2223        {
2224            let tmp = TempDir::new().unwrap();
2225            let paths = fresh_run(&tmp);
2226            bootstrap_live_node(&paths);
2227            agree(
2228                &paths,
2229                &ev(&paths, "node.report", Some("n-0001"), json!({})),
2230                "report-bare",
2231            );
2232        }
2233        {
2234            let tmp = TempDir::new().unwrap();
2235            let paths = fresh_run(&tmp);
2236            bootstrap_live_node(&paths);
2237            agree(
2238                &paths,
2239                &ev(
2240                    &paths,
2241                    "node.report",
2242                    Some("n-0001"),
2243                    json!({ "success": true }),
2244                ),
2245                "report-good",
2246            );
2247        }
2248        {
2249            let tmp = TempDir::new().unwrap();
2250            let paths = fresh_run(&tmp);
2251            bootstrap_live_node(&paths);
2252            agree(
2253                &paths,
2254                &ev(&paths, "node.report", None, json!({})),
2255                "report-no-node-id",
2256            );
2257        }
2258        {
2259            let tmp = TempDir::new().unwrap();
2260            let paths = fresh_run(&tmp);
2261            bootstrap_live_node(&paths);
2262            agree(
2263                &paths,
2264                &ev(&paths, "node.status", Some("n-0001"), json!({})),
2265                "status-missing",
2266            );
2267        }
2268        // Terminal node: a malformed report is a clean no-op (guard before
2269        // validate) — both must accept it.
2270        {
2271            let tmp = TempDir::new().unwrap();
2272            let paths = fresh_run(&tmp);
2273            bootstrap_live_node(&paths);
2274            append_and_apply_event(
2275                &paths,
2276                "node.report",
2277                Some(&nid("n-0001")),
2278                None,
2279                json!({ "success": true }),
2280            )
2281            .unwrap();
2282            agree(
2283                &paths,
2284                &ev(&paths, "node.report", Some("n-0001"), json!({})),
2285                "report-bare-on-terminal",
2286            );
2287        }
2288        // Missing node: a status with no `status` field is a no-op.
2289        {
2290            let tmp = TempDir::new().unwrap();
2291            let paths = fresh_run(&tmp);
2292            agree(
2293                &paths,
2294                &ev(&paths, "node.status", Some("n-0001"), json!({})),
2295                "status-missing-node",
2296            );
2297        }
2298        // Existing manifest: a run.status with no `status` is rejected.
2299        {
2300            let tmp = TempDir::new().unwrap();
2301            let paths = fresh_run(&tmp);
2302            bootstrap_live_node(&paths);
2303            agree(
2304                &paths,
2305                &ev(&paths, "run.status", None, json!({})),
2306                "run-status-missing",
2307            );
2308        }
2309        // Open discussion: a resolve without `resolution` is rejected.
2310        {
2311            let tmp = TempDir::new().unwrap();
2312            let paths = fresh_run(&tmp);
2313            bootstrap_live_node(&paths);
2314            append_and_apply_event(
2315                &paths,
2316                "discussion.opened",
2317                Some(&nid("n-0001")),
2318                None,
2319                json!({ "discussion_id": "d-abcdefghij", "topic": "t", "node_id": "n-0001" }),
2320            )
2321            .unwrap();
2322            agree(
2323                &paths,
2324                &ev(
2325                    &paths,
2326                    "discussion.resolved",
2327                    None,
2328                    json!({ "discussion_id": "d-abcdefghij" }),
2329                ),
2330                "resolve-missing-resolution",
2331            );
2332        }
2333        // node.created: new node missing `kind` rejected; replay over an
2334        // existing node with bad payload is a no-op (existence short-circuit).
2335        {
2336            let tmp = TempDir::new().unwrap();
2337            let paths = fresh_run(&tmp);
2338            agree(
2339                &paths,
2340                &ev(&paths, "node.created", Some("n-0002"), json!({})),
2341                "node-created-missing-kind",
2342            );
2343        }
2344        {
2345            let tmp = TempDir::new().unwrap();
2346            let paths = fresh_run(&tmp);
2347            bootstrap_live_node(&paths);
2348            agree(
2349                &paths,
2350                &ev(&paths, "node.created", Some("n-0001"), json!({})),
2351                "node-created-replay-bad-payload",
2352            );
2353        }
2354        // discussion.opened missing `topic`.
2355        {
2356            let tmp = TempDir::new().unwrap();
2357            let paths = fresh_run(&tmp);
2358            bootstrap_live_node(&paths);
2359            agree(
2360                &paths,
2361                &ev(
2362                    &paths,
2363                    "discussion.opened",
2364                    Some("n-0001"),
2365                    json!({ "discussion_id": "d-abcdefghij", "node_id": "n-0001" }),
2366                ),
2367                "discussion-opened-missing-topic",
2368            );
2369        }
2370        // spinoff.proposed missing `proposed_title`; spinoff.{approved,rejected}
2371        // with an unparseable proposal id.
2372        {
2373            let tmp = TempDir::new().unwrap();
2374            let paths = fresh_run(&tmp);
2375            bootstrap_live_node(&paths);
2376            agree(
2377                &paths,
2378                &ev(
2379                    &paths,
2380                    "spinoff.proposed",
2381                    Some("n-0001"),
2382                    json!({ "proposal_id": "p-abcdefghij", "proposed_kind": "spinoff", "node_id": "n-0001" }),
2383                ),
2384                "spinoff-proposed-missing-title",
2385            );
2386        }
2387        {
2388            let tmp = TempDir::new().unwrap();
2389            let paths = fresh_run(&tmp);
2390            agree(
2391                &paths,
2392                &ev(
2393                    &paths,
2394                    "spinoff.approved",
2395                    None,
2396                    json!({ "proposal_id": "not a valid id" }),
2397                ),
2398                "spinoff-approved-bad-id",
2399            );
2400            agree(
2401                &paths,
2402                &ev(
2403                    &paths,
2404                    "spinoff.rejected",
2405                    None,
2406                    json!({ "proposal_id": "not a valid id" }),
2407                ),
2408                "spinoff-rejected-bad-id",
2409            );
2410        }
2411        // child.spawned: missing/invalid child_run_id.
2412        {
2413            let tmp = TempDir::new().unwrap();
2414            let paths = fresh_run(&tmp);
2415            agree(
2416                &paths,
2417                &ev(&paths, "child.spawned", Some("n-0001"), json!({})),
2418                "child-spawned-missing-child-run-id",
2419            );
2420            agree(
2421                &paths,
2422                &ev(
2423                    &paths,
2424                    "child.spawned",
2425                    Some("n-0001"),
2426                    json!({ "child_run_id": "bad" }),
2427                ),
2428                "child-spawned-bad-child-run-id",
2429            );
2430        }
2431        // Cross-run envelope and unknown kind.
2432        {
2433            let tmp = TempDir::new().unwrap();
2434            let paths = fresh_run(&tmp);
2435            let mut foreign = ev(&paths, "run.status", None, json!({ "status": "running" }));
2436            foreign.run_id = crate::schema::RunId::parse_str("02jxsnap000000000000000000").unwrap();
2437            agree(&paths, &foreign, "cross-run");
2438            agree(
2439                &paths,
2440                &ev(&paths, "totally.unknown", None, json!({})),
2441                "unknown-kind",
2442            );
2443        }
2444    }
2445
2446    #[test]
2447    fn read_all_events_drops_torn_final_line() {
2448        // The bug this fixes: `read_all_events` used to silently ACCEPT a
2449        // valid-JSON final line lacking a trailing newline — a line
2450        // `recover_last_seq` discards as an uncommitted partial write. Now it
2451        // shares the torn-tail policy: the torn final line is dropped without
2452        // error, and the reader agrees with `recover_last_seq`.
2453        let tmp = TempDir::new().unwrap();
2454        let mut log = String::new();
2455        log.push_str(
2456            r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2457        );
2458        log.push('\n');
2459        // A COMPLETE, valid-JSON event whose trailing newline never flushed.
2460        log.push_str(
2461            r#"{"ts":"2026-06-12T00:00:00Z","seq":2,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2462        );
2463        let paths = paths_with_events(&tmp, log.as_bytes());
2464
2465        let events = read_all_events(&paths.events()).unwrap();
2466        assert_eq!(
2467            events.iter().map(|e| e.seq).collect::<Vec<_>>(),
2468            vec![1],
2469            "torn final line must be dropped, not parsed"
2470        );
2471        // And it agrees with the recovery path.
2472        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 1);
2473    }
2474
2475    #[test]
2476    fn recover_last_seq_rejects_seq_only_last_line() {
2477        // A `\n`-terminated last line that is valid JSON with a `seq` but is
2478        // NOT a valid event envelope (missing ts/kind/run_id) must be rejected
2479        // by recover_last_seq, matching read_all_events — otherwise an append
2480        // would continue past a line replay can never fold.
2481        let tmp = TempDir::new().unwrap();
2482        let paths = paths_with_events(&tmp, b"{\"seq\":99}\n");
2483        let err = recover_last_seq(&paths.events()).unwrap_err();
2484        assert!(matches!(err, Error::CorruptEventLog { .. }), "got {err:?}");
2485        // And the forward reader agrees.
2486        assert!(matches!(
2487            read_all_events(&paths.events()).unwrap_err(),
2488            Error::CorruptEventLog { .. }
2489        ));
2490    }
2491
2492    #[test]
2493    fn recover_last_seq_skips_multiple_trailing_blank_lines() {
2494        // External editing can leave several trailing blank lines. The forward
2495        // reader skips them; seq recovery must walk back over all of them to
2496        // the last real record (not just one), so the two readers agree.
2497        let tmp = TempDir::new().unwrap();
2498        let mut log = String::new();
2499        log.push_str(
2500            r#"{"ts":"2026-06-12T00:00:00Z","seq":7,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2501        );
2502        log.push_str("\n\n\n\n");
2503        let paths = paths_with_events(&tmp, log.as_bytes());
2504        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 7);
2505        let events = read_all_events(&paths.events()).unwrap();
2506        assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![7]);
2507    }
2508
2509    #[test]
2510    fn recover_last_seq_skips_trailing_whitespace_only_lines() {
2511        // External editing can leave trailing lines holding only spaces, tabs,
2512        // or stray CRs. Recovery must walk back over every whitespace-only line
2513        // to the last real record, not stop at (and fail to parse) the blanks.
2514        let tmp = TempDir::new().unwrap();
2515        let mut log = String::new();
2516        log.push_str(
2517            r#"{"ts":"2026-06-12T00:00:00Z","seq":5,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2518        );
2519        log.push_str("\n  \n\t\n \r\n");
2520        let paths = paths_with_events(&tmp, log.as_bytes());
2521        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 5);
2522    }
2523
2524    #[test]
2525    fn recover_last_seq_all_whitespace_file_is_zero() {
2526        // A log holding only blank/whitespace lines carries no event — recovery
2527        // returns the zero-event sentinel rather than erroring on the blanks.
2528        let tmp = TempDir::new().unwrap();
2529        let paths = paths_with_events(&tmp, b"\n  \n\t\n \r\n");
2530        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 0);
2531    }
2532
2533    #[test]
2534    fn recover_last_seq_single_newline_terminated_record_is_regression_guard() {
2535        // The common, healthy case: one record with a single trailing newline
2536        // must still recover its seq unchanged after the blank-line tolerance.
2537        let tmp = TempDir::new().unwrap();
2538        let paths = paths_with_events(
2539            &tmp,
2540            concat!(
2541                r#"{"ts":"2026-06-12T00:00:00Z","seq":5,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2542                "\n",
2543            )
2544            .as_bytes(),
2545        );
2546        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 5);
2547    }
2548
2549    #[test]
2550    fn read_all_events_rejects_corrupt_middle_line() {
2551        // A newline-terminated garbage line FOLLOWED by another line is
2552        // interior corruption — a hard `CorruptEventLog`, never a silent skip.
2553        let tmp = TempDir::new().unwrap();
2554        let log = concat!(
2555            r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2556            "\n",
2557            "{not valid json at all\n",
2558            r#"{"ts":"2026-06-12T00:00:00Z","seq":3,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2559            "\n",
2560        );
2561        let paths = paths_with_events(&tmp, log.as_bytes());
2562        let err = read_all_events(&paths.events()).unwrap_err();
2563        match err {
2564            Error::CorruptEventLog { reason, .. } => {
2565                assert!(reason.contains("line 2"), "reason was: {reason}");
2566            }
2567            other => panic!("expected CorruptEventLog, got {other:?}"),
2568        }
2569    }
2570
2571    #[test]
2572    fn append_truncates_torn_tail_before_writing() {
2573        // A crash left a valid record then a torn (newline-less) partial
2574        // write. The next append must truncate the torn bytes BEFORE writing,
2575        // so the log never gains a `…torn…{"seq":N}` malformed line.
2576        let tmp = TempDir::new().unwrap();
2577        let mut bytes = Vec::new();
2578        bytes.extend_from_slice(
2579            br#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2580        );
2581        bytes.push(b'\n');
2582        bytes.extend_from_slice(br#"{"seq":2,"kind":"TORN_PARTIAL_NEVER_FLUSHED"#); // no newline
2583        let paths = paths_with_events(&tmp, &bytes);
2584
2585        // The torn tail is ignored for seq recovery (last complete seq = 1).
2586        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 1);
2587
2588        // `marker` is an unknown kind → reducer no-op, so the append succeeds
2589        // without any projection prerequisites.
2590        let r = append_and_apply_event(&paths, "marker", None, None, serde_json::json!({"x": 1}))
2591            .unwrap();
2592        assert_eq!(r.seq, 2, "seq continues from the last complete record");
2593
2594        let raw = std::fs::read(paths.events()).unwrap();
2595        assert!(
2596            raw.ends_with(b"\n"),
2597            "log must be newline-terminated after a clean append"
2598        );
2599        assert!(
2600            !String::from_utf8_lossy(&raw).contains("TORN_PARTIAL_NEVER_FLUSHED"),
2601            "the torn tail must be truncated away before the append"
2602        );
2603        let events = read_all_events(&paths.events()).unwrap();
2604        assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1, 2]);
2605    }
2606
2607    #[test]
2608    fn append_truncates_all_torn_file_to_empty_then_writes_seq_1() {
2609        // The whole file is one torn (newline-less) partial write — no complete
2610        // record exists. truncate_torn_tail must cut it to empty, and the next
2611        // append starts a fresh seq 1.
2612        let tmp = TempDir::new().unwrap();
2613        let paths = paths_with_events(&tmp, br#"{"seq":1,"kind":"marker"#);
2614        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 0);
2615
2616        let r = append_and_apply_event(&paths, "marker", None, None, json!({})).unwrap();
2617        assert_eq!(r.seq, 1);
2618        let events = read_all_events(&paths.events()).unwrap();
2619        assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1]);
2620    }
2621
2622    #[test]
2623    fn truncate_torn_tail_cuts_partial_line_at_last_newline() {
2624        // The headline case (issue torn-write-truncate-tail): a complete
2625        // record followed by a torn (newline-less) partial write. Recovery
2626        // must cut the file back to the byte immediately after the last
2627        // complete record's trailing `\n` — the partial bytes are gone.
2628        let tmp = TempDir::new().unwrap();
2629        let complete = r#"{"ts":"2026-06-12T00:00:00Z","seq":5,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#;
2630        let mut bytes = Vec::new();
2631        bytes.extend_from_slice(complete.as_bytes());
2632        bytes.push(b'\n');
2633        let keep = bytes.len() as u64; // offset just past seq-5's newline
2634        bytes.extend_from_slice(br#"{"seq":6,"par"#); // torn mid-line, no newline
2635        let paths = paths_with_events(&tmp, &bytes);
2636
2637        truncate_torn_tail(&paths.events()).unwrap();
2638
2639        let raw = std::fs::read(paths.events()).unwrap();
2640        assert_eq!(
2641            raw.len() as u64,
2642            keep,
2643            "file must end at the offset after seq-5's newline"
2644        );
2645        assert!(raw.ends_with(b"\n"), "file is newline-terminated after cut");
2646        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 5);
2647    }
2648
2649    #[test]
2650    fn truncate_torn_tail_clean_file_is_noop() {
2651        // A file already ending in `\n` is the clean, common case: recovery
2652        // must leave every byte untouched (no rewrite, no length change).
2653        let tmp = TempDir::new().unwrap();
2654        let log = concat!(
2655            r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2656            "\n",
2657        );
2658        let paths = paths_with_events(&tmp, log.as_bytes());
2659
2660        truncate_torn_tail(&paths.events()).unwrap();
2661
2662        assert_eq!(
2663            std::fs::read(paths.events()).unwrap(),
2664            log.as_bytes(),
2665            "a clean, newline-terminated log must be left byte-for-byte intact"
2666        );
2667    }
2668
2669    #[test]
2670    fn truncate_torn_tail_zero_length_file_is_noop() {
2671        // An empty log has no tail to cut: recovery is a no-op and the file
2672        // stays empty.
2673        let tmp = TempDir::new().unwrap();
2674        let paths = paths_with_events(&tmp, b"");
2675        truncate_torn_tail(&paths.events()).unwrap();
2676        assert_eq!(std::fs::read(paths.events()).unwrap(), b"");
2677    }
2678
2679    #[test]
2680    fn truncate_torn_tail_missing_file_is_noop() {
2681        // No `events.jsonl` at all (a run that never appended): recovery must
2682        // not create the file or error.
2683        let tmp = TempDir::new().unwrap();
2684        let dir = tmp.path().join("run");
2685        std::fs::create_dir_all(&dir).unwrap();
2686        let paths = RunPaths::new(dir, "01jxsnap000000000000000000").unwrap();
2687        truncate_torn_tail(&paths.events()).unwrap();
2688        assert!(!paths.events().exists());
2689    }
2690
2691    #[test]
2692    fn truncate_torn_tail_single_complete_row_is_noop() {
2693        // Exactly one complete `\n`-terminated record and nothing else: the
2694        // last byte is already a newline, so there is no tail to cut.
2695        let tmp = TempDir::new().unwrap();
2696        let log = concat!(
2697            r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2698            "\n",
2699        );
2700        let paths = paths_with_events(&tmp, log.as_bytes());
2701        truncate_torn_tail(&paths.events()).unwrap();
2702        assert_eq!(std::fs::read(paths.events()).unwrap(), log.as_bytes());
2703    }
2704
2705    #[test]
2706    fn truncate_torn_tail_single_partial_row_truncates_to_zero() {
2707        // The whole file is one torn (newline-less) partial write with no
2708        // complete record ahead of it: there is nothing to keep, so recovery
2709        // truncates the file to zero length.
2710        let tmp = TempDir::new().unwrap();
2711        let paths = paths_with_events(&tmp, br#"{"seq":1,"kind":"marker"#);
2712        truncate_torn_tail(&paths.events()).unwrap();
2713        assert_eq!(
2714            std::fs::read(paths.events()).unwrap(),
2715            b"",
2716            "a file holding only a partial row must be cut to empty"
2717        );
2718        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 0);
2719    }
2720
2721    #[test]
2722    fn quarantine_excises_corrupt_middle_line_and_recovers() {
2723        // A valid record, a newline-terminated garbage line, then another
2724        // valid record. Quarantine must rename the original aside, write a
2725        // recovered log holding only the two valid lines, and report the bad
2726        // line's byte offset.
2727        let tmp = TempDir::new().unwrap();
2728        let good1 = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#;
2729        let bad = "{not valid json at all";
2730        let good3 = r#"{"ts":"2026-06-12T00:00:00Z","seq":3,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#;
2731        let log = format!("{good1}\n{bad}\n{good3}\n");
2732        let paths = paths_with_events(&tmp, log.as_bytes());
2733
2734        // Strict replay chokes on the poison line beforehand.
2735        assert!(matches!(
2736            read_all_events(&paths.events()).unwrap_err(),
2737            Error::CorruptEventLog { .. }
2738        ));
2739
2740        let q = quarantine_corrupt_lines(&paths, "20260612T000000Z")
2741            .unwrap()
2742            .expect("a corrupt line was excised");
2743        // The bad line started at the byte after `good1\n`.
2744        assert_eq!(q.removed_byte_offsets, vec![(good1.len() + 1) as u64]);
2745        assert_eq!(
2746            q.backup_path.file_name().unwrap().to_str().unwrap(),
2747            "events.jsonl.corrupt-20260612T000000Z.bak"
2748        );
2749
2750        // The backup is the verbatim original; the recovered log now replays
2751        // strictly with only the two valid records.
2752        assert_eq!(std::fs::read(&q.backup_path).unwrap(), log.as_bytes());
2753        let events = read_all_events(&paths.events()).unwrap();
2754        assert_eq!(events.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1, 3]);
2755    }
2756
2757    #[test]
2758    fn quarantine_clean_log_is_noop() {
2759        // A log with no corruption must not be renamed or rewritten.
2760        let tmp = TempDir::new().unwrap();
2761        let log = concat!(
2762            r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#,
2763            "\n",
2764        );
2765        let paths = paths_with_events(&tmp, log.as_bytes());
2766        assert!(quarantine_corrupt_lines(&paths, "20260612T000000Z")
2767            .unwrap()
2768            .is_none());
2769        // No backup created; original untouched.
2770        assert_eq!(std::fs::read(paths.events()).unwrap(), log.as_bytes());
2771        let bak = paths
2772            .events()
2773            .with_file_name("events.jsonl.corrupt-20260612T000000Z.bak");
2774        assert!(!bak.exists());
2775    }
2776
2777    #[test]
2778    fn quarantine_missing_log_is_none() {
2779        let tmp = TempDir::new().unwrap();
2780        let dir = tmp.path().join("run");
2781        std::fs::create_dir_all(&dir).unwrap();
2782        let paths = RunPaths::new(dir, "01jxsnap000000000000000000").unwrap();
2783        assert!(quarantine_corrupt_lines(&paths, "20260612T000000Z")
2784            .unwrap()
2785            .is_none());
2786    }
2787
2788    #[test]
2789    fn quarantine_preserves_torn_tail_and_excises_only_corruption() {
2790        // A valid record, a corrupt newline-terminated line, then a torn
2791        // (newline-less) final line. Only the corrupt middle line is excised;
2792        // the torn tail is retained verbatim (the readers tolerate it as an
2793        // in-flight partial write — excising it would change behavior).
2794        let tmp = TempDir::new().unwrap();
2795        let good = r#"{"ts":"2026-06-12T00:00:00Z","seq":1,"kind":"marker","run_id":"01jxsnap000000000000000000","data":{}}"#;
2796        let bad = "{garbage";
2797        let torn = r#"{"seq":2,"kind":"node.rep"#; // mid-write, no newline
2798        let mut log = Vec::new();
2799        log.extend_from_slice(format!("{good}\n{bad}\n{torn}").as_bytes());
2800        let paths = paths_with_events(&tmp, &log);
2801
2802        let q = quarantine_corrupt_lines(&paths, "20260612T000000Z")
2803            .unwrap()
2804            .expect("the corrupt middle line was excised");
2805        assert_eq!(q.removed_byte_offsets, vec![(good.len() + 1) as u64]);
2806
2807        let recovered = std::fs::read(paths.events()).unwrap();
2808        assert_eq!(recovered, format!("{good}\n{torn}").as_bytes());
2809        // The torn tail still recovers the last complete seq as 1.
2810        assert_eq!(recover_last_seq(&paths.events()).unwrap(), 1);
2811    }
2812
2813    #[test]
2814    fn find_prior_with_key_rejects_torn_middle_line() {
2815        // A newline-terminated garbage line FOLLOWED by another line: this
2816        // is interior corruption, not an in-flight tail. It must be a hard
2817        // error, never a silent skip — a skipped line could carry the very
2818        // key being looked up and let the caller double-append.
2819        let tmp = TempDir::new().unwrap();
2820        let log = concat!(
2821            r#"{"seq":1,"kind":"node.report","idempotency_key":"k0","node_id":"n-1","data":{}}"#,
2822            "\n",
2823            "{not valid json at all\n",
2824            r#"{"seq":3,"kind":"node.report","idempotency_key":"k1","node_id":"n-1","data":{}}"#,
2825            "\n",
2826        );
2827        let paths = paths_with_events(&tmp, log.as_bytes());
2828        let err = scan(&paths, "node.report", "k1").unwrap_err();
2829        match err {
2830            Error::CorruptEventLog { reason, .. } => {
2831                assert!(reason.contains("line 2"), "reason was: {reason}");
2832                assert!(reason.contains("last good seq 1"), "reason was: {reason}");
2833            }
2834            other => panic!("expected CorruptEventLog, got {other:?}"),
2835        }
2836    }
2837}