Skip to main content

treeship_core/session/
event_log.rs

1//! Append-only, file-backed event log for session events.
2//!
3//! Events are stored as newline-delimited JSON (JSONL) in
4//! `.treeship/sessions/<session_id>/events.jsonl`.
5//!
6//! Concurrency model: `append()` is safe to call from multiple processes
7//! concurrently. Each call attempts to acquire an exclusive advisory lock
8//! (via `fs2::FileExt::try_lock_exclusive` -- backed by `flock(2)` on Unix
9//! and `LockFileEx` on Windows) on a sidecar `events.jsonl.lock` file in a
10//! ~500ms bounded retry loop. Under the lock, a counter sidecar
11//! `events.jsonl.count` is the authoritative source for the next
12//! `sequence_no`. The per-process AtomicU64 is retained as a hot-path
13//! optimization for non-contended use, but its value is overwritten by the
14//! on-disk counter after every locked append.
15//!
16//! Counter sidecar format (16 bytes):
17//!   - bytes 0..8:  count (u64 LE) -- number of events written to events.jsonl
18//!   - bytes 8..16: byte_size (u64 LE) -- size of events.jsonl when count was recorded
19//!
20//! The byte_size field is the crash detector. If a peer wrote events.jsonl
21//! but crashed before fsyncing the counter (or vice versa), the size on disk
22//! and the size in the counter disagree. On any mismatch we fall back to an
23//! O(N) line count and rewrite the counter -- one paid scan, then back to
24//! O(1) on every subsequent append.
25//!
26//! This bounds steady-state append cost at constant: read 16 bytes, write
27//! one JSONL line, write 16 bytes. The previous implementation re-streamed
28//! the entire events.jsonl on every append, which made hooks O(N) in
29//! session length and dominated PostToolUse latency on long sessions.
30//!
31//! Fail-closed semantics: the writer ALWAYS acquires the exclusive flock
32//! before reading the counter and writing the event. We use fs2's blocking
33//! `lock_exclusive()` (flock(2) without LOCK_NB on Unix, LockFileEx without
34//! LOCKFILE_FAIL_IMMEDIATELY on Windows) so contended writers queue rather
35//! than race. An earlier "best-effort, fall through on contention" path
36//! existed here -- it could produce duplicate `sequence_no` under hook
37//! contention (P0, audit lane F), so it was removed. Hook callers already
38//! sit under Claude Code's 60s hook timeout, which bounds wall-clock
39//! exposure to a wedged peer. Trading a bounded wait for guaranteed
40//! injective sequence numbers is the right call when receipts are the
41//! trust artifact.
42//!
43//! Lock file permissions are 0o600 (owner-only) on Unix, applied at file
44//! creation via `OpenOptionsExt::mode` and re-tightened on every open if
45//! a previous run left the file with looser perms.
46
47use std::io::{BufRead, Write};
48use std::path::{Path, PathBuf};
49use std::sync::atomic::{AtomicU64, Ordering};
50
51#[cfg(not(target_family = "wasm"))]
52use fs2::FileExt;
53
54use crate::session::event::SessionEvent;
55
56/// Error from event log operations.
57#[derive(Debug)]
58pub enum EventLogError {
59    Io(std::io::Error),
60    Json(serde_json::Error),
61}
62
63impl std::fmt::Display for EventLogError {
64    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
65        match self {
66            Self::Io(e) => write!(f, "event log io: {e}"),
67            Self::Json(e) => write!(f, "event log json: {e}"),
68        }
69    }
70}
71
72impl std::error::Error for EventLogError {}
73impl From<std::io::Error> for EventLogError {
74    fn from(e: std::io::Error) -> Self {
75        Self::Io(e)
76    }
77}
78impl From<serde_json::Error> for EventLogError {
79    fn from(e: serde_json::Error) -> Self {
80        Self::Json(e)
81    }
82}
83
84/// An append-only event log backed by a JSONL file.
85pub struct EventLog {
86    path: PathBuf,
87    sequence: AtomicU64,
88}
89
90impl EventLog {
91    /// Open or create an event log for the given session directory.
92    ///
93    /// The session directory is typically `.treeship/sessions/<session_id>/`.
94    /// If the directory does not exist, it will be created.
95    ///
96    /// Initialization reads the counter sidecar in O(1) when present and
97    /// consistent with events.jsonl's byte size; falls back to an O(N) line
98    /// count (and rewrites the sidecar) when the sidecar is missing,
99    /// short-read, or stale from a crashed previous appender.
100    pub fn open(session_dir: &Path) -> Result<Self, EventLogError> {
101        std::fs::create_dir_all(session_dir)?;
102        let path = session_dir.join("events.jsonl");
103        // Read-only. `open` used to call `read_counter_or_recount`, which
104        // REWRITES the counter sidecar when it finds it stale or missing --
105        // an unlocked read-modify-write of the same shared state
106        // `append_locked` takes an exclusive flock to protect.
107        //
108        // In the model that matters, every hook invocation constructs a fresh
109        // handle, so `open` sits inside the contention window: one process
110        // could clobber the counter with a value computed before another
111        // process's append landed. The size check in
112        // `read_counter_consistent` makes that self-healing rather than
113        // fatal, but a constructor writing shared state behind the lock's
114        // back is a discipline violation regardless of whether a given
115        // interleaving is survivable. Issue #275.
116        //
117        // Locking `open` instead would serialize every handle construction
118        // on an exclusive flock just to seed a hint. Not writing is both
119        // cheaper and more obviously correct: on native, this value is only
120        // a hint. `append_locked` re-reads the counter INSIDE the lock and
121        // overwrites `sequence` (see the store at the end of that function),
122        // and counter repair still happens there, where it is serialized.
123        let count = read_counter_hint(&path);
124        Ok(Self {
125            path,
126            sequence: AtomicU64::new(count),
127        })
128    }
129
130    /// Append a single event to the log.
131    ///
132    /// The event's `sequence_no` is set automatically. Under contention from
133    /// multiple writer processes, the sequence number is re-derived from the
134    /// on-disk line count under an exclusive flock so two parallel writers
135    /// never collide.
136    pub fn append(&self, event: &mut SessionEvent) -> Result<(), EventLogError> {
137        self.append_locked(event)
138    }
139
140    /// Cross-process safe append: acquires an exclusive advisory lock on a
141    /// sidecar `.lock` file, re-counts events.jsonl lines, assigns sequence_no,
142    /// writes the new event, then releases the lock on drop.
143    ///
144    /// Locking is BLOCKING (`fs2::FileExt::lock_exclusive`, i.e. flock(2)
145    /// without LOCK_NB). A previous implementation polled `try_lock_exclusive`
146    /// for ~500ms and then fell through to an UNLOCKED write -- which under
147    /// hook contention (multiple PostToolUse invocations racing) could
148    /// assign duplicate `sequence_no` values to different events. That
149    /// broke the injective-sequence invariant receipts depend on, even
150    /// though local merkle verification still passed. Audit lane F (P0).
151    ///
152    /// The trade-off: a wedged peer holding the lock will now stall this
153    /// caller until the peer releases (e.g. by crashing -- flock is
154    /// released by the kernel when the holder's FD closes). Hook
155    /// invocations are bounded by Claude Code's hook timeout (60s by
156    /// default), so the worst-case wedge surfaces as a hook failure
157    /// rather than a duplicate sequence_no. That mirrors how
158    /// `journal/mod.rs` treats the journal append lock as a hard
159    /// correctness barrier rather than a soft hint.
160    ///
161    /// The locked region covers BOTH the counter read AND the file write,
162    /// so two concurrent writers cannot read the same count and append
163    /// twice with that count.
164    ///
165    /// Lock file is created mode 0o600 (owner-only) so the sidecar can
166    /// never be opened by other users on a shared machine.
167    ///
168    /// Skipped on WASM (no fs, no concurrency).
169    #[cfg(not(target_family = "wasm"))]
170    fn append_locked(&self, event: &mut SessionEvent) -> Result<(), EventLogError> {
171        // Sidecar lock file: contention here doesn't block readers of events.jsonl.
172        let lock_path = self.path.with_extension("jsonl.lock");
173
174        // Open or create the lock file. On Unix we set 0o600 explicitly so
175        // the sidecar isn't group/world readable; the umask-derived default
176        // would otherwise be permissive on some setups.
177        let lock_file = open_lock_file(&lock_path)?;
178
179        // Blocking flock. Returns when the lock is held exclusively, or
180        // propagates a real I/O error (filesystem that doesn't support
181        // flock, FD revoked, etc.). EINTR retry isn't needed -- fs2 wraps
182        // the syscall and retries internally on POSIX.
183        FileExt::lock_exclusive(&lock_file)?;
184
185        // From here until `lock_file` is dropped (end of function), we
186        // hold the exclusive flock. The counter read + event write + counter
187        // update MUST stay inside this block; any early return must still
188        // drop `lock_file`, which Rust guarantees by RAII.
189        let result = (|| -> Result<(), EventLogError> {
190            // Read sequence_no from the counter sidecar in O(1) when
191            // consistent with events.jsonl size. Stale or missing counters
192            // force a one-time O(N) rescan that also rewrites the counter,
193            // so subsequent appends return to O(1). Only the on-disk state
194            // (counter + size check) is authoritative when multiple
195            // processes are appending; the per-process AtomicU64 is a
196            // stale hint.
197            let count = read_counter_or_recount(&self.path)?;
198            event.sequence_no = count;
199
200            let mut line = serde_json::to_vec(event)?;
201            line.push(b'\n');
202
203            let mut file = std::fs::OpenOptions::new()
204                .create(true)
205                .append(true)
206                .open(&self.path)?;
207            file.write_all(&line)?;
208            file.flush()?;
209
210            // Update the counter sidecar with the new count and the new
211            // events.jsonl size, so the next append can short-circuit the
212            // line scan. Failure to update the counter is non-fatal: the
213            // next reader will detect the size mismatch and recount.
214            let new_size = file.metadata().map(|m| m.len()).unwrap_or(0);
215            let _ = write_counter(&self.path, count + 1, new_size);
216
217            // Keep the in-process AtomicU64 in sync so non-contended callers
218            // see the right value via event_count() without re-reading.
219            self.sequence.store(count + 1, Ordering::SeqCst);
220            Ok(())
221        })();
222
223        // Explicit unlock matches the journal precedent (journal/mod.rs
224        // also calls `unlock` before dropping). Drop alone would release
225        // the flock via close(2), but being explicit makes the lock
226        // window obvious to readers of this function.
227        let _ = FileExt::unlock(&lock_file);
228        result
229    }
230
231    /// WASM build: no filesystem locks available, no concurrent writers.
232    /// Falls back to the simple AtomicU64 path.
233    #[cfg(target_family = "wasm")]
234    fn append_locked(&self, event: &mut SessionEvent) -> Result<(), EventLogError> {
235        event.sequence_no = self.sequence.fetch_add(1, Ordering::SeqCst);
236
237        let mut line = serde_json::to_vec(event)?;
238        line.push(b'\n');
239
240        let mut file = std::fs::OpenOptions::new()
241            .create(true)
242            .append(true)
243            .open(&self.path)?;
244        file.write_all(&line)?;
245        file.flush()?;
246
247        Ok(())
248    }
249
250    /// Read all events from the log.
251    ///
252    /// Per-line tolerant: a single bad line -- malformed JSON (unknown
253    /// event type, missing field, truncated write) OR a non-UTF-8 byte
254    /// (partial write, corruption) -- is logged to stderr, counted as a
255    /// skip, and stepped over, not propagated as an error. The caller --
256    /// session close, in particular -- composes a receipt from whatever
257    /// events parse, instead of dropping every event when any one is bad.
258    /// (The read only returns Err for a whole-file failure such as the
259    /// events file being unopenable; a per-line decode error is a skip.)
260    ///
261    /// Why this matters: events.jsonl is append-only and written by
262    /// hooks, daemons, SDKs, and bridges from multiple processes. A
263    /// single bad event from one buggy emitter would otherwise nuke
264    /// the entire receipt's side_effects / agent_graph / timeline.
265    /// Real-world repro: a hook that emitted events with an unknown
266    /// `type` field caused side_effects.files_written to come back
267    /// empty even though the rest of the events in the log were valid
268    /// agent.wrote_file events the aggregator would have happily
269    /// processed.
270    pub fn read_all(&self) -> Result<Vec<SessionEvent>, EventLogError> {
271        // Drop the skipped count for callers that don't carry it through.
272        // Receipt composition uses read_all_with_stats to record the
273        // count in-band on the sealed receipt -- see Codex finding #8.
274        self.read_all_with_stats().map(|(events, _skipped)| events)
275    }
276
277    /// Same as `read_all` but returns the count of malformed lines that
278    /// were skipped during parsing alongside the valid events.
279    ///
280    /// Codex adversarial review finding #8: skipping malformed events
281    /// on stderr only is silent data loss from the verifier's
282    /// perspective. The receipt gets sealed under a merkle root that
283    /// represents only the events that successfully parsed -- a
284    /// downstream consumer cannot tell whether the receipt is complete
285    /// or whether N events were silently dropped.
286    ///
287    /// session::close calls this and stores the count on
288    /// `receipt.proofs.event_log_skipped`. `treeship package verify`
289    /// surfaces it as a WARN when nonzero so the receipt's
290    /// completeness signal is visible without breaking byte-identical
291    /// re-verification of pre-existing receipts.
292    pub fn read_all_with_stats(&self) -> Result<(Vec<SessionEvent>, usize), EventLogError> {
293        if !self.path.exists() {
294            return Ok((Vec::new(), 0));
295        }
296        let file = std::fs::File::open(&self.path)?;
297        let reader = std::io::BufReader::new(file);
298        let mut events = Vec::new();
299        let mut skipped = 0usize;
300        for (idx, line) in reader.lines().enumerate() {
301            // A per-LINE read error (a non-UTF-8 byte from a partial write or
302            // corruption) must be counted as a skip, not propagated — the `?`
303            // here previously aborted the WHOLE read, and `session::close`
304            // maps that Err to (empty, 0), sealing an empty receipt with a
305            // zero skip-count exactly when the log is most damaged. That
306            // silently bypasses the completeness signal this function exists
307            // to carry. Treat the bad line like a malformed JSON line: skip
308            // it, count it, keep the rest.
309            let line = match line {
310                Ok(l) => l,
311                Err(e) => {
312                    skipped += 1;
313                    eprintln!(
314                        "[treeship] event_log: skipping unreadable line {} in {}: {}",
315                        idx + 1,
316                        self.path.display(),
317                        e,
318                    );
319                    continue;
320                }
321            };
322            if line.trim().is_empty() {
323                continue;
324            }
325            match serde_json::from_str::<SessionEvent>(&line) {
326                Ok(event) => events.push(event),
327                Err(e) => {
328                    skipped += 1;
329                    eprintln!(
330                        "[treeship] event_log: skipping malformed line {} in {}: {}",
331                        idx + 1,
332                        self.path.display(),
333                        e,
334                    );
335                }
336            }
337        }
338        if skipped > 0 {
339            eprintln!(
340                "[treeship] event_log: {} malformed line(s) skipped while reading {} (kept {} valid event(s))",
341                skipped,
342                self.path.display(),
343                events.len(),
344            );
345        }
346        Ok((events, skipped))
347    }
348
349    /// Return the current event count.
350    pub fn event_count(&self) -> u64 {
351        self.sequence.load(Ordering::SeqCst)
352    }
353
354    /// Return the path to the JSONL file.
355    pub fn path(&self) -> &Path {
356        &self.path
357    }
358}
359
360/// Open the sidecar lock file with owner-only permissions (0o600 on Unix).
361///
362/// On Unix the mode is set atomically via `OpenOptionsExt::mode` for newly
363/// created files. For files that already exist (e.g. left over from a
364/// prior crash or an upgrade from a pre-0.9.3 CLI that didn't tighten
365/// perms), we additionally re-chmod to 0o600 after open IF the file is
366/// owned by the current user. This is best-effort: if the chmod fails
367/// (file owned by another user, read-only filesystem, etc.) we proceed
368/// silently rather than refuse to open the lock -- the lock semantics
369/// don't depend on the perms being tight, only the privacy of the
370/// sidecar's existence does.
371///
372/// On Windows the mode concept doesn't apply; ACLs default to inheriting
373/// the parent dir's permissions, which for `.treeship/sessions/<id>/`
374/// should already be scoped to the owning user.
375#[cfg(all(not(target_family = "wasm"), unix))]
376fn open_lock_file(path: &Path) -> Result<std::fs::File, std::io::Error> {
377    use std::os::unix::fs::{MetadataExt, OpenOptionsExt, PermissionsExt};
378    use std::os::unix::io::AsRawFd;
379
380    let file = std::fs::OpenOptions::new()
381        .create(true)
382        // Explicitly NOT truncating: this is a flock target, and its contents
383        // are irrelevant, but truncating would race a concurrent holder.
384        .truncate(false)
385        .read(true)
386        .write(true)
387        .mode(0o600)
388        .open(path)?;
389
390    // Re-tighten if a pre-existing file has loose perms. Use `fchmod` on the
391    // open file descriptor rather than `set_permissions(path, ...)` to
392    // eliminate the TOCTOU window -- between metadata() and a path-based
393    // chmod, an attacker could swap the file. `fchmod` operates on the
394    // already-opened inode, so the target is pinned.
395    //
396    // Only act when the file is owned by us (uid match via geteuid). If
397    // fchmod fails (NFS mount with restricted metadata writes, or some
398    // filesystems without full POSIX perm support), emit a one-line
399    // stderr warning so an operator has visibility. The lock still works;
400    // only the privacy of the sidecar's existence is affected.
401    if let Ok(meta) = file.metadata() {
402        let mode = meta.permissions().mode() & 0o777;
403        let owned_by_us = meta.uid() == nix_uid();
404        if owned_by_us && mode != 0o600 {
405            let fd = file.as_raw_fd();
406            // `libc_fchmod` is a safe wrapper; the unsafe extern lives inside
407            // it. The caller-side obligation it documents still holds here:
408            // fd is valid (we just opened it) and 0o600 is a well-formed mode.
409            let rc = libc_fchmod(fd, 0o600);
410            if rc != 0 {
411                let err = std::io::Error::last_os_error();
412                eprintln!(
413                    "[treeship] warning: could not tighten lock file perms on {} \
414                     to 0o600 (current: 0o{:o}). Error: {}. Lock still functions; \
415                     only the privacy of the sidecar is affected. Common cause: \
416                     NFS mount or filesystem without full POSIX perm support.",
417                    path.display(),
418                    mode,
419                    err
420                );
421            }
422        }
423    }
424
425    Ok(file)
426}
427
428/// Thin FFI wrapper around libc::fchmod. Declared here so event_log.rs
429/// doesn't need a direct libc crate dep -- the symbol is available in
430/// every Unix libc binary.
431#[cfg(all(not(target_family = "wasm"), unix))]
432fn libc_fchmod(fd: i32, mode: u32) -> i32 {
433    // SAFETY: posix-standard FFI signature; `fd` validity and `mode`
434    // bounds are enforced by the caller.
435    unsafe extern "C" {
436        fn fchmod(fd: i32, mode: u32) -> i32;
437    }
438    unsafe { fchmod(fd, mode) }
439}
440
441/// Lightweight wrapper around `geteuid` so we can compare to file ownership
442/// without pulling in the `nix` crate. Uses `libc` directly (already a
443/// transitive dep via several upstream crates).
444#[cfg(all(not(target_family = "wasm"), unix))]
445fn nix_uid() -> u32 {
446    // SAFETY: geteuid is async-signal-safe and never fails per POSIX.
447    unsafe extern "C" {
448        fn geteuid() -> u32;
449    }
450    unsafe { geteuid() }
451}
452
453#[cfg(all(not(target_family = "wasm"), not(unix)))]
454fn open_lock_file(path: &Path) -> Result<std::fs::File, std::io::Error> {
455    std::fs::OpenOptions::new()
456        .create(true)
457        .read(true)
458        .write(true)
459        .open(path)
460}
461
462/// Path of the counter sidecar for a given events.jsonl path.
463fn counter_path(events_path: &Path) -> PathBuf {
464    events_path.with_extension("jsonl.count")
465}
466
467/// Read the counter sidecar if it exists and is consistent with events.jsonl.
468///
469/// Returns `Some(count)` when the sidecar's recorded byte_size matches the
470/// current events.jsonl size, and `None` otherwise (missing sidecar, short
471/// read, parse failure, or size mismatch from a crashed previous appender).
472#[cfg(not(target_family = "wasm"))]
473fn read_counter_consistent(events_path: &Path) -> Option<u64> {
474    let counter = counter_path(events_path);
475    let bytes = std::fs::read(&counter).ok()?;
476    if bytes.len() != 16 {
477        return None;
478    }
479    let count = u64::from_le_bytes(bytes[0..8].try_into().ok()?);
480    let recorded_size = u64::from_le_bytes(bytes[8..16].try_into().ok()?);
481
482    // events.jsonl may not exist yet -- counter records (0, 0) for that case.
483    let actual_size = match std::fs::metadata(events_path) {
484        Ok(m) => m.len(),
485        Err(e) if e.kind() == std::io::ErrorKind::NotFound => 0,
486        Err(_) => return None,
487    };
488    if actual_size != recorded_size {
489        return None;
490    }
491    Some(count)
492}
493
494/// Read the counter via the sidecar (O(1)) or fall back to an O(N) line
495/// scan, rewriting the sidecar on the way out. This is the recovery path
496/// after a crash that left the counter and events.jsonl out of sync.
497#[cfg(not(target_family = "wasm"))]
498/// Best-effort event count for seeding the in-process hint, with NO writes.
499///
500/// Used by `open`. Never repairs the counter sidecar: repair belongs inside
501/// the append lock, where it cannot race a concurrent writer. A wrong hint is
502/// harmless on native -- `append_locked` re-reads authoritatively under the
503/// lock before assigning any `sequence_no`.
504#[cfg(not(target_family = "wasm"))]
505fn read_counter_hint(events_path: &Path) -> u64 {
506    if let Some(count) = read_counter_consistent(events_path) {
507        return count;
508    }
509    let Ok(f) = std::fs::File::open(events_path) else {
510        return 0;
511    };
512    std::io::BufReader::new(f)
513        .lines()
514        .filter(|l| l.is_ok())
515        .count() as u64
516}
517
518/// WASM: no fs, no concurrent writers, so the in-memory AtomicU64 is the
519/// whole story. There is no wasm `read_counter_or_recount` counterpart --
520/// its only caller is `append_locked`, which is itself native-only.
521#[cfg(target_family = "wasm")]
522fn read_counter_hint(_events_path: &Path) -> u64 {
523    0
524}
525
526#[cfg(not(target_family = "wasm"))]
527fn read_counter_or_recount(events_path: &Path) -> Result<u64, EventLogError> {
528    if let Some(count) = read_counter_consistent(events_path) {
529        return Ok(count);
530    }
531    let count = if events_path.exists() {
532        let f = std::fs::File::open(events_path)?;
533        let r = std::io::BufReader::new(f);
534        r.lines().filter(|l| l.is_ok()).count() as u64
535    } else {
536        0
537    };
538    let size = std::fs::metadata(events_path).map(|m| m.len()).unwrap_or(0);
539    let _ = write_counter(events_path, count, size);
540    Ok(count)
541}
542
543/// Atomically replace the counter sidecar with the new (count, byte_size).
544///
545/// Writes to a temp file in the same directory and renames into place so a
546/// reader either sees the old 16 bytes or the new 16 bytes, never a partial
547/// write. The 0o600 perm matches the lock file -- the counter doesn't leak
548/// secrets but its existence is a session signal worth scoping to the owner.
549#[cfg(not(target_family = "wasm"))]
550fn write_counter(events_path: &Path, count: u64, byte_size: u64) -> Result<(), std::io::Error> {
551    use std::io::Write as _;
552    let counter = counter_path(events_path);
553    let dir = counter.parent().ok_or_else(|| {
554        std::io::Error::new(
555            std::io::ErrorKind::InvalidInput,
556            "counter path has no parent",
557        )
558    })?;
559    std::fs::create_dir_all(dir)?;
560
561    let mut buf = [0u8; 16];
562    buf[0..8].copy_from_slice(&count.to_le_bytes());
563    buf[8..16].copy_from_slice(&byte_size.to_le_bytes());
564
565    let tmp = counter.with_extension("count.tmp");
566    {
567        let mut f = open_counter_tmp(&tmp)?;
568        f.write_all(&buf)?;
569        f.sync_all()?;
570    }
571    std::fs::rename(&tmp, &counter)?;
572    Ok(())
573}
574
575#[cfg(all(not(target_family = "wasm"), unix))]
576fn open_counter_tmp(path: &Path) -> Result<std::fs::File, std::io::Error> {
577    use std::os::unix::fs::OpenOptionsExt;
578    std::fs::OpenOptions::new()
579        .create(true)
580        .write(true)
581        .truncate(true)
582        .mode(0o600)
583        .open(path)
584}
585
586#[cfg(all(not(target_family = "wasm"), not(unix)))]
587fn open_counter_tmp(path: &Path) -> Result<std::fs::File, std::io::Error> {
588    std::fs::OpenOptions::new()
589        .create(true)
590        .write(true)
591        .truncate(true)
592        .open(path)
593}
594
595#[cfg(test)]
596mod tests {
597    use super::*;
598    use crate::session::event::*;
599
600    pub(super) fn make_event(session_id: &str, event_type: EventType) -> SessionEvent {
601        SessionEvent {
602            session_id: session_id.into(),
603            event_id: generate_event_id(),
604            timestamp: "2026-04-05T08:00:00Z".into(),
605            sequence_no: 0,
606            trace_id: generate_trace_id(),
607            span_id: generate_span_id(),
608            parent_span_id: None,
609            agent_id: "agent://test".into(),
610            agent_instance_id: "ai_test_1".into(),
611            agent_name: "test-agent".into(),
612            agent_role: None,
613            host_id: "host_test".into(),
614            tool_runtime_id: None,
615            event_type,
616            artifact_ref: None,
617            meta: None,
618        }
619    }
620
621    #[test]
622    fn append_and_read_back() {
623        let dir =
624            std::env::temp_dir().join(format!("treeship-evtlog-test-{}", rand::random::<u32>()));
625        let log = EventLog::open(&dir).unwrap();
626
627        let mut e1 = make_event("ssn_001", EventType::SessionStarted);
628        let mut e2 = make_event(
629            "ssn_001",
630            EventType::AgentStarted {
631                parent_agent_instance_id: None,
632            },
633        );
634
635        log.append(&mut e1).unwrap();
636        log.append(&mut e2).unwrap();
637
638        assert_eq!(log.event_count(), 2);
639        assert_eq!(e1.sequence_no, 0);
640        assert_eq!(e2.sequence_no, 1);
641
642        let events = log.read_all().unwrap();
643        assert_eq!(events.len(), 2);
644        assert_eq!(events[0].sequence_no, 0);
645        assert_eq!(events[1].sequence_no, 1);
646
647        let _ = std::fs::remove_dir_all(&dir);
648    }
649
650    #[test]
651    fn read_all_skips_malformed_lines() {
652        // Regression: a single malformed line in events.jsonl used to
653        // make read_all() return Err, and the caller's
654        // .unwrap_or_default() would drop EVERY event in the log. Real
655        // bug: hooks emitting events with an unknown `type` field made
656        // side_effects.files_written come back empty even though every
657        // other event in the log was a perfectly valid agent.wrote_file
658        // event. Now we skip-and-log the bad line and keep the rest.
659        let dir = std::env::temp_dir().join(format!(
660            "treeship-evtlog-malformed-{}",
661            rand::random::<u32>()
662        ));
663        let log = EventLog::open(&dir).unwrap();
664
665        let mut good1 = make_event(
666            "ssn_001",
667            EventType::AgentWroteFile {
668                file_path: "src/before.rs".into(),
669                digest: None,
670                operation: None,
671                additions: None,
672                deletions: None,
673            },
674        );
675        let mut good2 = make_event(
676            "ssn_001",
677            EventType::AgentWroteFile {
678                file_path: "src/after.rs".into(),
679                digest: None,
680                operation: None,
681                additions: None,
682                deletions: None,
683            },
684        );
685        log.append(&mut good1).unwrap();
686        log.append(&mut good2).unwrap();
687
688        // Manually inject a malformed line between the two good ones by
689        // truncating the file and rewriting. The malformed line has an
690        // unknown event type ("custom.weird") which the closed EventType
691        // enum can't deserialize.
692        let path = log.path().to_path_buf();
693        let original = std::fs::read_to_string(&path).unwrap();
694        let mut lines: Vec<&str> = original.lines().collect();
695        lines.insert(1, r#"{"session_id":"ssn_001","event_id":"evt_bad","timestamp":"2026-04-26T00:00:00Z","sequence_no":1,"trace_id":"x","span_id":"y","agent_id":"a","agent_instance_id":"i","agent_name":"n","host_id":"h","type":"custom.weird","payload":42}"#);
696        std::fs::write(&path, lines.join("\n") + "\n").unwrap();
697
698        let events = log.read_all().unwrap();
699        assert_eq!(
700            events.len(),
701            2,
702            "expected the two valid events to come through; got {}",
703            events.len()
704        );
705        // Confirm the valid events are the file-write events and not
706        // some default fallback.
707        let written_paths: Vec<&str> = events
708            .iter()
709            .filter_map(|e| match &e.event_type {
710                EventType::AgentWroteFile { file_path, .. } => Some(file_path.as_str()),
711                _ => None,
712            })
713            .collect();
714        assert_eq!(written_paths, vec!["src/before.rs", "src/after.rs"]);
715
716        // Codex finding #8: the count must be exposed in-band so a
717        // sealed receipt can carry the incompleteness signal. Verify
718        // read_all_with_stats reports it.
719        let (events2, skipped) = log.read_all_with_stats().unwrap();
720        assert_eq!(events2.len(), 2);
721        assert_eq!(
722            skipped, 1,
723            "exactly one malformed line was injected; expected skipped == 1"
724        );
725
726        let _ = std::fs::remove_dir_all(&dir);
727    }
728
729    #[test]
730    fn read_all_with_stats_reports_zero_when_clean() {
731        // No malformed lines -> skipped == 0 -> the receipt's
732        // event_log_skipped field stays default (0) and gets omitted
733        // from canonical JSON. This preserves byte-identical receipts
734        // for the common case where the event log is clean.
735        let dir =
736            std::env::temp_dir().join(format!("treeship-evtlog-clean-{}", rand::random::<u32>()));
737        let log = EventLog::open(&dir).unwrap();
738
739        let mut e = make_event(
740            "ssn_001",
741            EventType::AgentWroteFile {
742                file_path: "x.rs".into(),
743                digest: None,
744                operation: None,
745                additions: None,
746                deletions: None,
747            },
748        );
749        log.append(&mut e).unwrap();
750
751        let (events, skipped) = log.read_all_with_stats().unwrap();
752        assert_eq!(events.len(), 1);
753        assert_eq!(skipped, 0);
754
755        let _ = std::fs::remove_dir_all(&dir);
756    }
757
758    #[test]
759    fn non_utf8_byte_is_skipped_not_aborted() {
760        // A single non-UTF-8 byte in the log (partial write / corruption)
761        // must skip THAT line and keep the rest, with a nonzero skip count —
762        // NOT abort the whole read and seal an empty receipt with skipped=0.
763        use std::io::Write;
764        let dir =
765            std::env::temp_dir().join(format!("treeship-evtlog-badbyte-{}", rand::random::<u32>()));
766        let log = EventLog::open(&dir).unwrap();
767        let mut e = make_event(
768            "ssn_001",
769            EventType::AgentWroteFile {
770                file_path: "ok.rs".into(),
771                digest: None,
772                operation: None,
773                additions: None,
774                deletions: None,
775            },
776        );
777        log.append(&mut e).unwrap();
778
779        // Append a raw non-UTF-8 line directly to the events file.
780        let events_path = dir.join("events.jsonl");
781        let mut f = std::fs::OpenOptions::new()
782            .append(true)
783            .open(&events_path)
784            .unwrap();
785        f.write_all(&[0xff, 0xfe, b'\n']).unwrap(); // invalid UTF-8 line
786        drop(f);
787
788        let (events, skipped) = log.read_all_with_stats().unwrap();
789        assert_eq!(events.len(), 1, "the one good event must survive");
790        assert_eq!(
791            skipped, 1,
792            "the bad line must be counted as skipped, not silently dropped"
793        );
794
795        let _ = std::fs::remove_dir_all(&dir);
796    }
797
798    #[test]
799    fn reopen_preserves_sequence() {
800        let dir =
801            std::env::temp_dir().join(format!("treeship-evtlog-reopen-{}", rand::random::<u32>()));
802
803        {
804            let log = EventLog::open(&dir).unwrap();
805            let mut e = make_event("ssn_001", EventType::SessionStarted);
806            log.append(&mut e).unwrap();
807        }
808
809        // Reopen
810        let log = EventLog::open(&dir).unwrap();
811        assert_eq!(log.event_count(), 1);
812
813        let mut e2 = make_event(
814            "ssn_001",
815            EventType::AgentStarted {
816                parent_agent_instance_id: None,
817            },
818        );
819        log.append(&mut e2).unwrap();
820        assert_eq!(e2.sequence_no, 1);
821
822        let _ = std::fs::remove_dir_all(&dir);
823    }
824
825    /// Regression test for #1 in the v0.9.3 Codex adversarial review.
826    ///
827    /// Multiple `EventLog` instances opened against the same directory must
828    /// not collide on `sequence_no`. This simulates what happens when each
829    /// `treeship session event` invocation (one per PostToolUse hook firing)
830    /// creates a fresh `EventLog` on a shared events.jsonl. Without the
831    /// flock-based re-derivation in `append_locked`, every instance sees
832    /// the same on-disk count at open time and assigns duplicate sequence
833    /// numbers.
834    #[cfg(not(target_family = "wasm"))]
835    #[test]
836    fn concurrent_appends_have_unique_sequence_numbers() {
837        use std::sync::Arc;
838        use std::thread;
839
840        let dir =
841            std::env::temp_dir().join(format!("treeship-evtlog-race-{}", rand::random::<u32>()));
842        std::fs::create_dir_all(&dir).unwrap();
843
844        const WRITERS: usize = 16;
845        let dir = Arc::new(dir);
846        let mut handles = Vec::with_capacity(WRITERS);
847
848        for _ in 0..WRITERS {
849            let dir = Arc::clone(&dir);
850            handles.push(thread::spawn(move || {
851                // Each thread opens its OWN EventLog -- mimics a separate
852                // process invocation. Without flock, all threads would see
853                // the same line count at open() time.
854                let log = EventLog::open(&dir).unwrap();
855                let mut e = make_event("ssn_race", EventType::SessionStarted);
856                log.append(&mut e).unwrap();
857                e.sequence_no
858            }));
859        }
860
861        let mut seqs: Vec<u64> = handles.into_iter().map(|h| h.join().unwrap()).collect();
862        seqs.sort();
863
864        // All sequence numbers must be unique and contiguous 0..WRITERS.
865        let expected: Vec<u64> = (0..WRITERS as u64).collect();
866        assert_eq!(seqs, expected, "sequence_no collisions under contention");
867
868        // Same invariant from the on-disk file's perspective.
869        let log = EventLog::open(&dir).unwrap();
870        let read = log.read_all().unwrap();
871        assert_eq!(read.len(), WRITERS);
872        let mut on_disk: Vec<u64> = read.iter().map(|e| e.sequence_no).collect();
873        on_disk.sort();
874        assert_eq!(on_disk, expected);
875
876        let _ = std::fs::remove_dir_all(&*dir);
877    }
878
879    /// Sidecar lock file must be created mode 0o600 (owner-only) on Unix.
880    /// Regression test for #5 in the second Codex adversarial review.
881    #[cfg(all(not(target_family = "wasm"), unix))]
882    #[test]
883    fn lock_file_has_owner_only_permissions() {
884        use std::os::unix::fs::PermissionsExt;
885
886        let dir =
887            std::env::temp_dir().join(format!("treeship-evtlog-perms-{}", rand::random::<u32>()));
888        let log = EventLog::open(&dir).unwrap();
889
890        let mut e = make_event("ssn_perms", EventType::SessionStarted);
891        log.append(&mut e).unwrap();
892
893        let lock_path = log.path().with_extension("jsonl.lock");
894        let meta = std::fs::metadata(&lock_path).expect("lock file must exist after first append");
895        let mode = meta.permissions().mode() & 0o777;
896        assert_eq!(
897            mode, 0o600,
898            "lock file mode is {:o}, expected 0o600 (owner-only)",
899            mode
900        );
901
902        let _ = std::fs::remove_dir_all(&dir);
903    }
904
905    /// A pre-existing lock file (e.g. from a v0.9.2 era crash) with looser
906    /// permissions must be tightened to 0o600 on next `EventLog::open`.
907    /// Regression test for the third Codex adversarial review.
908    #[cfg(all(not(target_family = "wasm"), unix))]
909    #[test]
910    fn existing_lock_file_is_re_tightened() {
911        use std::os::unix::fs::PermissionsExt;
912
913        let dir = std::env::temp_dir().join(format!(
914            "treeship-evtlog-retighten-{}",
915            rand::random::<u32>()
916        ));
917        std::fs::create_dir_all(&dir).unwrap();
918
919        // Pre-create a lock file with deliberately loose perms, simulating
920        // an upgrade from a CLI version that didn't set 0o600.
921        let lock_path = dir.join("events.jsonl.lock");
922        std::fs::write(&lock_path, b"").unwrap();
923        std::fs::set_permissions(&lock_path, std::fs::Permissions::from_mode(0o644)).unwrap();
924        let pre_mode = std::fs::metadata(&lock_path).unwrap().permissions().mode() & 0o777;
925        assert_eq!(
926            pre_mode, 0o644,
927            "test setup: pre-existing perms should be 0o644"
928        );
929
930        // First append after upgrade -- should re-tighten.
931        let log = EventLog::open(&dir).unwrap();
932        let mut e = make_event("ssn_retighten", EventType::SessionStarted);
933        log.append(&mut e).unwrap();
934
935        let post_mode = std::fs::metadata(&lock_path).unwrap().permissions().mode() & 0o777;
936        assert_eq!(
937            post_mode, 0o600,
938            "lock file should be re-tightened to 0o600 after open; got {:o}",
939            post_mode
940        );
941
942        let _ = std::fs::remove_dir_all(&dir);
943    }
944
945    /// Counter sidecar must exist after the first append and contain
946    /// (count=1, byte_size=size of events.jsonl). This is the happy path
947    /// that lets every subsequent append skip the O(N) rescan.
948    #[cfg(not(target_family = "wasm"))]
949    #[test]
950    fn counter_sidecar_written_after_append() {
951        let dir =
952            std::env::temp_dir().join(format!("treeship-evtlog-counter-{}", rand::random::<u32>()));
953        let log = EventLog::open(&dir).unwrap();
954
955        let mut e = make_event("ssn_counter", EventType::SessionStarted);
956        log.append(&mut e).unwrap();
957
958        let counter = log.path().with_extension("jsonl.count");
959        let bytes = std::fs::read(&counter).expect("counter sidecar must exist after append");
960        assert_eq!(bytes.len(), 16, "counter sidecar must be 16 bytes");
961
962        let count = u64::from_le_bytes(bytes[0..8].try_into().unwrap());
963        let recorded_size = u64::from_le_bytes(bytes[8..16].try_into().unwrap());
964        let actual_size = std::fs::metadata(log.path()).unwrap().len();
965        assert_eq!(count, 1, "counter must reflect the one appended event");
966        assert_eq!(
967            recorded_size, actual_size,
968            "counter byte_size ({}) must match events.jsonl size ({})",
969            recorded_size, actual_size
970        );
971
972        let _ = std::fs::remove_dir_all(&dir);
973    }
974
975    /// A missing counter sidecar (fresh install, deleted by user, etc.)
976    /// must not break sequence_no assignment. The next append falls back
977    /// to an O(N) recount and rewrites the counter.
978    #[cfg(not(target_family = "wasm"))]
979    #[test]
980    fn counter_sidecar_recovers_when_missing() {
981        let dir = std::env::temp_dir().join(format!(
982            "treeship-evtlog-missing-counter-{}",
983            rand::random::<u32>()
984        ));
985
986        // Append two events, then nuke the counter sidecar.
987        {
988            let log = EventLog::open(&dir).unwrap();
989            let mut e1 = make_event("ssn_x", EventType::SessionStarted);
990            let mut e2 = make_event(
991                "ssn_x",
992                EventType::AgentStarted {
993                    parent_agent_instance_id: None,
994                },
995            );
996            log.append(&mut e1).unwrap();
997            log.append(&mut e2).unwrap();
998        }
999        let counter = dir.join("events.jsonl.count");
1000        std::fs::remove_file(&counter).expect("counter must exist before deletion");
1001
1002        // Reopen + append. The third event must get sequence_no=2 even
1003        // though the counter sidecar is gone.
1004        let log = EventLog::open(&dir).unwrap();
1005        assert_eq!(
1006            log.event_count(),
1007            2,
1008            "open() must recount when counter is missing"
1009        );
1010
1011        let mut e3 = make_event(
1012            "ssn_x",
1013            EventType::SessionClosed {
1014                summary: None,
1015                duration_ms: None,
1016            },
1017        );
1018        log.append(&mut e3).unwrap();
1019        assert_eq!(e3.sequence_no, 2);
1020        assert!(counter.exists(), "counter must be rewritten after recount");
1021
1022        let _ = std::fs::remove_dir_all(&dir);
1023    }
1024
1025    /// A short-read or garbage counter sidecar (corrupted, partial write,
1026    /// truncated by external tool) must not be trusted. The size mismatch
1027    /// path covers the "wrong content" case for a 16-byte file too.
1028    #[cfg(not(target_family = "wasm"))]
1029    #[test]
1030    fn counter_sidecar_recovers_when_corrupt() {
1031        let dir = std::env::temp_dir().join(format!(
1032            "treeship-evtlog-corrupt-counter-{}",
1033            rand::random::<u32>()
1034        ));
1035
1036        {
1037            let log = EventLog::open(&dir).unwrap();
1038            let mut e = make_event("ssn_corrupt", EventType::SessionStarted);
1039            log.append(&mut e).unwrap();
1040        }
1041        // Truncate the counter to a non-16 length.
1042        let counter = dir.join("events.jsonl.count");
1043        std::fs::write(&counter, b"junk").unwrap();
1044
1045        let log = EventLog::open(&dir).unwrap();
1046        assert_eq!(
1047            log.event_count(),
1048            1,
1049            "short-read counter must be ignored, recount kicks in"
1050        );
1051
1052        let _ = std::fs::remove_dir_all(&dir);
1053    }
1054
1055    /// A counter that recorded the wrong byte_size (someone or something
1056    /// appended to events.jsonl behind our back) must not be trusted.
1057    /// This is the crash-recovery path: peer wrote events.jsonl but
1058    /// crashed before fsyncing the counter, so the recorded size is stale.
1059    #[cfg(not(target_family = "wasm"))]
1060    #[test]
1061    fn counter_sidecar_recovers_when_size_disagrees() {
1062        let dir = std::env::temp_dir().join(format!(
1063            "treeship-evtlog-stale-counter-{}",
1064            rand::random::<u32>()
1065        ));
1066
1067        {
1068            let log = EventLog::open(&dir).unwrap();
1069            let mut e = make_event("ssn_stale", EventType::SessionStarted);
1070            log.append(&mut e).unwrap();
1071        }
1072
1073        // Simulate a crash mid-append: append one extra raw line to
1074        // events.jsonl WITHOUT updating the counter. Now the counter
1075        // says (1, S) but events.jsonl is (S + |line|) bytes.
1076        let events_path = dir.join("events.jsonl");
1077        let mut extra = make_event(
1078            "ssn_stale",
1079            EventType::AgentStarted {
1080                parent_agent_instance_id: None,
1081            },
1082        );
1083        extra.sequence_no = 999; // intentionally wrong; will be overwritten on read
1084        let mut line = serde_json::to_vec(&extra).unwrap();
1085        line.push(b'\n');
1086        let mut f = std::fs::OpenOptions::new()
1087            .append(true)
1088            .open(&events_path)
1089            .unwrap();
1090        std::io::Write::write_all(&mut f, &line).unwrap();
1091        std::io::Write::flush(&mut f).unwrap();
1092
1093        // Re-open. The size mismatch must trigger a recount; we should see 2.
1094        let log = EventLog::open(&dir).unwrap();
1095        assert_eq!(
1096            log.event_count(),
1097            2,
1098            "size mismatch must force recount, ignoring stale counter"
1099        );
1100
1101        let _ = std::fs::remove_dir_all(&dir);
1102    }
1103
1104    /// The counter sidecar fix must not break the cross-process race
1105    /// safety established by the flock layer. This is the same shape as
1106    /// `concurrent_appends_have_unique_sequence_numbers` but exists to
1107    /// guard against a regression where the counter is read OUTSIDE the
1108    /// lock, which would let two writers both see count=N and assign N
1109    /// to two different events.
1110    #[cfg(not(target_family = "wasm"))]
1111    #[test]
1112    fn counter_sidecar_preserves_concurrent_uniqueness() {
1113        use std::sync::Arc;
1114        use std::thread;
1115
1116        let dir = std::env::temp_dir().join(format!(
1117            "treeship-evtlog-counter-race-{}",
1118            rand::random::<u32>()
1119        ));
1120        std::fs::create_dir_all(&dir).unwrap();
1121
1122        const WRITERS: usize = 16;
1123        let dir = Arc::new(dir);
1124        let mut handles = Vec::with_capacity(WRITERS);
1125
1126        for _ in 0..WRITERS {
1127            let dir = Arc::clone(&dir);
1128            handles.push(thread::spawn(move || {
1129                let log = EventLog::open(&dir).unwrap();
1130                let mut e = make_event("ssn_counter_race", EventType::SessionStarted);
1131                log.append(&mut e).unwrap();
1132                e.sequence_no
1133            }));
1134        }
1135
1136        let mut seqs: Vec<u64> = handles.into_iter().map(|h| h.join().unwrap()).collect();
1137        seqs.sort();
1138        let expected: Vec<u64> = (0..WRITERS as u64).collect();
1139        assert_eq!(
1140            seqs, expected,
1141            "counter must not bypass the flock race protection"
1142        );
1143
1144        // Counter should reflect the final state.
1145        let log = EventLog::open(&dir).unwrap();
1146        assert_eq!(log.event_count(), WRITERS as u64);
1147
1148        let _ = std::fs::remove_dir_all(&*dir);
1149    }
1150
1151    /// Counter sidecar must be created mode 0o600 (owner-only) on Unix --
1152    /// same scoping as the lock file; the existence of a counter is a
1153    /// session signal that doesn't need to leak to other users.
1154    #[cfg(all(not(target_family = "wasm"), unix))]
1155    #[test]
1156    fn counter_sidecar_has_owner_only_permissions() {
1157        use std::os::unix::fs::PermissionsExt;
1158
1159        let dir = std::env::temp_dir().join(format!(
1160            "treeship-evtlog-counter-perms-{}",
1161            rand::random::<u32>()
1162        ));
1163        let log = EventLog::open(&dir).unwrap();
1164
1165        let mut e = make_event("ssn_counter_perms", EventType::SessionStarted);
1166        log.append(&mut e).unwrap();
1167
1168        let counter = log.path().with_extension("jsonl.count");
1169        let mode = std::fs::metadata(&counter).unwrap().permissions().mode() & 0o777;
1170        assert_eq!(
1171            mode, 0o600,
1172            "counter sidecar mode is {:o}, expected 0o600 (owner-only)",
1173            mode
1174        );
1175
1176        let _ = std::fs::remove_dir_all(&dir);
1177    }
1178
1179    /// P0 regression (audit lane F): under heavy hook contention, multiple
1180    /// writers used to fall through and append without the flock when the
1181    /// 500ms poll exhausted, producing duplicate `sequence_no` values. The
1182    /// blocking `lock_exclusive` fix means every writer must hold the
1183    /// flock across both the counter read and the event write.
1184    ///
1185    /// This test spawns 8 threads, each calling `append` 25 times on its
1186    /// own `EventLog` (mimicking 8 separate hook processes each appending
1187    /// a burst of events). After the join, the on-disk log must contain
1188    /// exactly 8*25=200 events with `sequence_no` exactly the contiguous
1189    /// range 0..200, no duplicates and no gaps.
1190    #[cfg(not(target_family = "wasm"))]
1191    #[test]
1192    fn p0_no_duplicate_sequence_under_burst_contention() {
1193        use std::sync::Arc;
1194        use std::thread;
1195
1196        const THREADS: usize = 8;
1197        const PER_THREAD: usize = 25;
1198        const EXPECTED: usize = THREADS * PER_THREAD;
1199
1200        let dir = std::env::temp_dir().join(format!(
1201            "treeship-evtlog-p0-burst-{}",
1202            rand::random::<u32>()
1203        ));
1204        std::fs::create_dir_all(&dir).unwrap();
1205        let dir = Arc::new(dir);
1206
1207        let mut handles = Vec::with_capacity(THREADS);
1208        for t in 0..THREADS {
1209            let dir = Arc::clone(&dir);
1210            handles.push(thread::spawn(move || -> Vec<u64> {
1211                // Each thread opens its own EventLog -- this is the
1212                // per-process model the audit flagged: every PostToolUse
1213                // invocation is a fresh handle on the shared log.
1214                let log = EventLog::open(&dir).unwrap();
1215                let mut seen = Vec::with_capacity(PER_THREAD);
1216                for i in 0..PER_THREAD {
1217                    let mut e =
1218                        make_event(&format!("ssn_burst_{}_{}", t, i), EventType::SessionStarted);
1219                    log.append(&mut e).unwrap();
1220                    seen.push(e.sequence_no);
1221                }
1222                seen
1223            }));
1224        }
1225
1226        // Collect what each thread saw locally.
1227        let mut all_returned: Vec<u64> = handles
1228            .into_iter()
1229            .flat_map(|h| h.join().unwrap())
1230            .collect();
1231        all_returned.sort();
1232
1233        let expected: Vec<u64> = (0..EXPECTED as u64).collect();
1234        assert_eq!(
1235            all_returned, expected,
1236            "returned sequence_no values must be a contiguous range 0..{} \
1237             with no duplicates and no gaps",
1238            EXPECTED
1239        );
1240
1241        // Truth source: the on-disk log itself. Verify (a) count and
1242        // (b) sequence_no is exactly 0..EXPECTED on the persisted events.
1243        let log = EventLog::open(&dir).unwrap();
1244        let events = log.read_all().unwrap();
1245        assert_eq!(
1246            events.len(),
1247            EXPECTED,
1248            "on-disk event count must be exactly {} (got {})",
1249            EXPECTED,
1250            events.len()
1251        );
1252        let mut on_disk: Vec<u64> = events.iter().map(|e| e.sequence_no).collect();
1253        on_disk.sort();
1254        assert_eq!(
1255            on_disk, expected,
1256            "on-disk sequence_no must be a contiguous range with no duplicates and no gaps"
1257        );
1258
1259        // Counter sidecar must agree with both.
1260        assert_eq!(log.event_count(), EXPECTED as u64);
1261
1262        let _ = std::fs::remove_dir_all(&*dir);
1263    }
1264
1265    /// Companion stress test for the lock-file lifecycle. Each append
1266    /// opens the lock file, locks it, writes, unlocks, and drops the FD.
1267    /// Repeatedly creating + dropping `EventLog`s in a tight loop must
1268    /// not panic, must not leak FDs we can detect (no `EMFILE` after
1269    /// hundreds of iterations on a default ulimit), and must produce a
1270    /// log with contiguous sequence numbers.
1271    #[cfg(not(target_family = "wasm"))]
1272    #[test]
1273    fn lock_file_handles_drop_cleanly_under_churn() {
1274        let dir = std::env::temp_dir().join(format!(
1275            "treeship-evtlog-fd-churn-{}",
1276            rand::random::<u32>()
1277        ));
1278        std::fs::create_dir_all(&dir).unwrap();
1279
1280        // 500 sequential open + append + drop cycles. Far below the
1281        // default macOS/Linux soft limit (256/1024) for a sustained
1282        // leak, but plenty to catch one-per-iteration FD leaks.
1283        const ITERS: usize = 500;
1284        for i in 0..ITERS {
1285            let log = EventLog::open(&dir).unwrap();
1286            let mut e = make_event(&format!("ssn_churn_{}", i), EventType::SessionStarted);
1287            log.append(&mut e).unwrap();
1288            // log drops here -> lock_file FD already closed inside append.
1289        }
1290
1291        let log = EventLog::open(&dir).unwrap();
1292        let events = log.read_all().unwrap();
1293        assert_eq!(events.len(), ITERS);
1294        let mut seqs: Vec<u64> = events.iter().map(|e| e.sequence_no).collect();
1295        seqs.sort();
1296        let expected: Vec<u64> = (0..ITERS as u64).collect();
1297        assert_eq!(
1298            seqs, expected,
1299            "no FD leak should still produce contiguous seqs"
1300        );
1301
1302        let _ = std::fs::remove_dir_all(&dir);
1303    }
1304}
1305
1306#[cfg(test)]
1307mod open_race_tests {
1308    use super::tests::make_event;
1309    use super::*;
1310    use crate::session::EventType;
1311    use std::sync::Arc;
1312    use std::thread;
1313
1314    /// Reproduction harness for #275.
1315    ///
1316    /// The burst test opens one `EventLog` per thread, which is the real
1317    /// model: every hook invocation is a fresh handle on a shared log. That
1318    /// makes `open()` part of the contended path, and `open()` performs an
1319    /// unlocked read-modify-write of the counter sidecar
1320    /// (`read_counter_or_recount` rewrites it on a stale/missing counter)
1321    /// while another process may hold the append lock.
1322    ///
1323    /// More threads and more rounds than the original: the failure is timing
1324    /// dependent and did not reproduce locally at 8x25 across 20 runs.
1325    #[test]
1326    fn open_and_append_interleaved_keep_sequences_unique() {
1327        const THREADS: usize = 16;
1328        const PER_THREAD: usize = 12;
1329        const ROUNDS: usize = 12;
1330
1331        for round in 0..ROUNDS {
1332            let dir =
1333                std::env::temp_dir().join(format!("ts-open-race-{}-{}", std::process::id(), round));
1334            let _ = std::fs::remove_dir_all(&dir);
1335            std::fs::create_dir_all(&dir).unwrap();
1336            let dir = Arc::new(dir);
1337
1338            let handles: Vec<_> = (0..THREADS)
1339                .map(|t| {
1340                    let dir = Arc::clone(&dir);
1341                    thread::spawn(move || -> Vec<u64> {
1342                        let mut seen = Vec::with_capacity(PER_THREAD);
1343                        for i in 0..PER_THREAD {
1344                            // Re-open every iteration. This is the point: it
1345                            // puts open() inside the contention window
1346                            // instead of once before it.
1347                            let log = EventLog::open(&dir).unwrap();
1348                            let mut e =
1349                                make_event(&format!("ssn_race_{t}_{i}"), EventType::SessionStarted);
1350                            log.append(&mut e).unwrap();
1351                            seen.push(e.sequence_no);
1352                        }
1353                        seen
1354                    })
1355                })
1356                .collect();
1357
1358            let mut all: Vec<u64> = handles
1359                .into_iter()
1360                .flat_map(|h| h.join().unwrap())
1361                .collect();
1362            all.sort_unstable();
1363            let expected: Vec<u64> = (0..(THREADS * PER_THREAD) as u64).collect();
1364            assert_eq!(
1365                all, expected,
1366                "round {round}: sequence_no values must be a contiguous range \
1367                 with no duplicates and no gaps"
1368            );
1369            let _ = std::fs::remove_dir_all(&*dir);
1370        }
1371    }
1372}