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}