Skip to main content

vtcode_memory/
event_log.rs

1//! Append-only per-session `ThreadEvent` log plus index and manifest.
2
3use std::collections::{BTreeMap, HashMap, VecDeque};
4use std::fs::File;
5use std::io::{BufRead, Read, Seek, SeekFrom, Write};
6use std::path::{Path, PathBuf};
7use std::sync::atomic::{AtomicBool, Ordering};
8use std::sync::{Arc, Mutex, OnceLock, Weak};
9
10use chrono::Utc;
11use serde::{Deserialize, Serialize};
12use vtcode_commons::VtCodePaths;
13use vtcode_exec_events::{EVENT_SCHEMA_VERSION, ThreadEvent, ThreadItemDetails, VersionedThreadEvent};
14
15use crate::error::SessionStoreError;
16
17/// Sidecar lock file flock-held while a session's event-log handles are open.
18pub(crate) const SESSION_LOCK_FILE: &str = "session.lock";
19use crate::manifest::{ManifestStore, PendingCapRewrite};
20use crate::session_dir;
21
22/// Default maximum number of events retained per session before the oldest
23/// completed turns are evicted.
24pub const DEFAULT_MAX_EVENTS: usize = 10_000;
25
26/// Maximum serialized event bytes retained before an append forces a write.
27/// Turn boundaries and reads still flush immediately.
28const MAX_WRITE_BUFFER_BYTES: usize = 64 * 1024;
29const MAX_EVICTION_GROUNDED_FACTS: usize = 32;
30const MAX_EVICTION_GROUNDED_FACT_BYTES: usize = 512;
31
32/// Callback used to persist a summary of events before they are evicted.
33///
34/// The callback runs after the event bytes have been flushed and decoded, but
35/// before the canonical log is rewritten. A failure leaves the original log
36/// and in-memory index untouched, so retention never silently discards
37/// history.
38pub type EvictionSummaryHook = Arc<dyn Fn(&[ThreadEvent]) -> Result<(), SessionStoreError> + Send + Sync>;
39
40/// Minimal envelope used while rebuilding the turn index.
41///
42/// The index only needs the event discriminator. Deserializing a complete
43/// [`VersionedThreadEvent`] here would allocate every nested tool argument,
44/// output, and thread item even though none of that payload is retained.
45#[derive(Debug, Deserialize)]
46struct VersionedEventKind<'a> {
47    #[serde(rename = "schema_version", borrow)]
48    _schema_version: &'a str,
49    #[serde(borrow)]
50    event: EventKind<'a>,
51}
52
53#[derive(Debug, Deserialize)]
54struct EventKind<'a> {
55    #[serde(rename = "type", borrow)]
56    kind: &'a str,
57}
58
59/// Zero-clone serialization envelope for `ThreadEvent`.
60///
61/// Produces JSON byte-identical to `VersionedThreadEvent` but borrows the
62/// event by reference instead of cloning it. `append` is called for every
63/// runtime event, and `ThreadEvent` can carry large tool outputs / thread
64/// items — cloning just to feed `serde_json::to_string` was pure waste.
65#[derive(Serialize)]
66struct BorrowedVersionedEvent<'a> {
67    schema_version: &'a str,
68    event: &'a ThreadEvent,
69}
70
71/// Turn-lifecycle discriminator extracted from either a `ThreadEvent` (at
72/// append time) or a raw `&str` kind (during scan).  This is the single
73/// representation that both code paths feed into
74/// [`LogState::apply_lifecycle_event`], eliminating a duplicated state machine.
75#[derive(Debug, Clone, Copy, PartialEq, Eq)]
76enum LifecycleKind {
77    ThreadStarted,
78    ThreadCompleted,
79    TurnStarted,
80    TurnCompleted,
81    TurnFailed,
82    Other,
83}
84
85impl LifecycleKind {
86    /// Discriminate from a runtime `ThreadEvent` at append time.
87    #[inline]
88    fn from_event(event: &ThreadEvent) -> Self {
89        match event {
90            ThreadEvent::ThreadStarted(_) => Self::ThreadStarted,
91            ThreadEvent::ThreadCompleted(_) => Self::ThreadCompleted,
92            ThreadEvent::TurnStarted(_) => Self::TurnStarted,
93            ThreadEvent::TurnCompleted(_) => Self::TurnCompleted,
94            ThreadEvent::TurnFailed(_) => Self::TurnFailed,
95            _ => Self::Other,
96        }
97    }
98
99    /// Discriminate from a raw event-type string at scan time.
100    #[inline]
101    fn from_kind(kind: &str) -> Self {
102        match kind {
103            "thread.started" => Self::ThreadStarted,
104            "thread.completed" => Self::ThreadCompleted,
105            "turn.started" => Self::TurnStarted,
106            "turn.completed" => Self::TurnCompleted,
107            "turn.failed" => Self::TurnFailed,
108            _ => Self::Other,
109        }
110    }
111}
112
113/// In-memory state protected by a mutex (cheap; appends are infrequent relative
114/// to model inference).
115struct LogState {
116    manifest: SessionManifest,
117    index: TurnIndex,
118    /// Whether we are currently inside a turn (between TurnStarted and
119    /// TurnCompleted/TurnFailed). Used to update the last index entry's
120    /// offsets as intermediate events arrive.
121    in_turn: bool,
122    /// Running byte offset of the next append. Avoids a `stat` syscall per
123    /// event (the previous implementation re-statted the file twice on every
124    /// `append`); initialized from the file length on `open`.
125    next_offset: u64,
126    /// Buffered pending writes to batch syscalls. Events are appended here
127    /// and flushed to disk at turn boundaries or before read operations.
128    write_buf: Vec<u8>,
129}
130
131#[derive(Debug, Clone, Copy)]
132struct CapEvictionPlan {
133    truncate_offset: u64,
134    evicted_event_count: u64,
135    evicted_turn_count: usize,
136}
137
138impl LogState {
139    fn new(session_id: &str) -> Self {
140        Self {
141            manifest: SessionManifest::new(session_id),
142            index: TurnIndex::default(),
143            in_turn: false,
144            next_offset: 0,
145            write_buf: Vec::with_capacity(65536),
146        }
147    }
148
149    /// Serialize `event` directly into the reusable write buffer with rollback
150    /// on failure.
151    ///
152    /// This encapsulates the invariant that `write_buf` never contains a
153    /// partial JSON document: if `serde_json::to_writer` fails mid-write the
154    /// buffer is truncated back to its pre-serialization boundary.  Returns
155    /// the `(start, end)` byte offsets of the serialized event so the caller
156    /// can feed them to [`Self::apply_lifecycle_event`].
157    fn serialize_event(&mut self, event: &ThreadEvent) -> Result<(u64, u64), SessionStoreError> {
158        let start = self.next_offset;
159        let buf_len_before = self.write_buf.len();
160        if let Err(err) = serde_json::to_writer(
161            &mut self.write_buf,
162            &BorrowedVersionedEvent { schema_version: EVENT_SCHEMA_VERSION, event },
163        ) {
164            self.write_buf.truncate(buf_len_before);
165            return Err(err.into());
166        }
167        self.write_buf.push(b'\n');
168        let written = self.write_buf.len() - buf_len_before;
169        let end = start + written as u64;
170        self.next_offset = end;
171        Ok((start, end))
172    }
173
174    /// Update the in-memory turn index and manifest counters for a single
175    /// event.
176    ///
177    /// This is the single implementation of the turn-lifecycle state machine;
178    /// both the append path (via [`LifecycleKind::from_event`]) and the scan
179    /// path (via [`LifecycleKind::from_kind`]) route through here, eliminating
180    /// a previously duplicated match block.
181    ///
182    /// Returns `true` when the event closes a turn boundary
183    /// (`TurnCompleted` / `TurnFailed`) so the caller can persist metadata
184    /// at the appropriate time (append persists immediately; scan persists
185    /// once after the full scan).
186    fn apply_lifecycle_event(&mut self, kind: LifecycleKind, start: u64, end: u64) -> bool {
187        let is_boundary = match kind {
188            LifecycleKind::ThreadStarted => {
189                self.manifest.status = "active".to_string();
190                false
191            }
192            LifecycleKind::ThreadCompleted => {
193                self.manifest.status = "completed".to_string();
194                true
195            }
196            LifecycleKind::TurnStarted => {
197                self.manifest.status = "active".to_string();
198                self.in_turn = true;
199                let n = self.manifest.turn_count + 1;
200                self.index.entries.push_back(TurnIndexEntry {
201                    turn_number: n,
202                    start_offset: start,
203                    end_offset: end,
204                    event_count: 1,
205                    ts: now_rfc3339(),
206                });
207                false
208            }
209            LifecycleKind::TurnCompleted | LifecycleKind::TurnFailed => {
210                if self.in_turn {
211                    if let Some(entry) = self.index.entries.back_mut() {
212                        entry.end_offset = end;
213                        entry.event_count += 1;
214                        // `turn_number` remains monotonic when older completed
215                        // turns have been evicted. The manifest is the source
216                        // for the next ordinal, so never replace it with the
217                        // retained index length.
218                        self.manifest.turn_count = self.manifest.turn_count.max(entry.turn_number);
219                    }
220                    self.in_turn = false;
221                }
222                true
223            }
224            LifecycleKind::Other => {
225                if self.in_turn
226                    && let Some(entry) = self.index.entries.back_mut()
227                {
228                    entry.end_offset = end;
229                    entry.event_count += 1;
230                }
231                false
232            }
233        };
234        // Persist the open-turn marker alongside the manifest. The marker is
235        // deliberately optional for backwards compatibility: an older
236        // manifest without it forces a scan on reopen so the state can be
237        // reconstructed from the canonical event log.
238        self.manifest.in_turn = Some(self.in_turn);
239        is_boundary
240    }
241
242    /// Plan a cap-enforcement eviction: pop the oldest completed turns from
243    /// the index until `event_count` is within `max_events`.
244    ///
245    /// Returns the byte offset at which the file should be rewritten and the
246    /// counts needed to apply the eviction after successful I/O. Returns
247    /// `None` when no eviction is needed.
248    fn plan_cap_eviction(&self, max_events: usize) -> Option<CapEvictionPlan> {
249        if max_events == 0 || self.manifest.event_count <= max_events as u64 {
250            return None;
251        }
252        let mut evicted_event_count = 0u64;
253        let mut truncate_offset = 0u64;
254        let mut evicted_turn_count = 0;
255        for (entry_index, oldest) in self.index.entries.iter().enumerate() {
256            // Never evict the active turn. Its entry has a provisional end
257            // offset and will be closed by a later completion/failure event;
258            // removing it here would make that event unindexed and lose the
259            // in-flight turn from reconstruction. If the completed history
260            // alone cannot bring the log under the cap, retain the active
261            // turn until it reaches a terminal boundary.
262            if self.in_turn && entry_index + 1 == self.index.entries.len() {
263                break;
264            }
265            if self.manifest.event_count.saturating_sub(evicted_event_count) <= max_events as u64 {
266                break;
267            }
268            truncate_offset = oldest.end_offset;
269            evicted_event_count += oldest.event_count;
270            evicted_turn_count += 1;
271        }
272        if truncate_offset == 0 || evicted_turn_count == 0 {
273            None
274        } else {
275            Some(CapEvictionPlan {
276                truncate_offset,
277                evicted_event_count,
278                evicted_turn_count,
279            })
280        }
281    }
282
283    fn apply_cap_eviction(&mut self, plan: CapEvictionPlan, next_offset: u64) {
284        for _ in 0..plan.evicted_turn_count {
285            let _ = self.index.entries.pop_front();
286        }
287        for entry in &mut self.index.entries {
288            entry.start_offset = entry.start_offset.saturating_sub(plan.truncate_offset);
289            entry.end_offset = entry.end_offset.saturating_sub(plan.truncate_offset);
290        }
291        self.next_offset = next_offset;
292        self.manifest.event_count = self.manifest.event_count.saturating_sub(plan.evicted_event_count);
293        self.manifest.retained_turn_base = Some(self.retained_turn_base_after_eviction(0));
294    }
295
296    fn retained_turn_base_after_eviction(&self, evicted_turn_count: usize) -> u64 {
297        self.index
298            .entries
299            .get(evicted_turn_count)
300            .map(|entry| entry.turn_number)
301            .or_else(|| self.manifest.turn_count.checked_add(1))
302            .unwrap_or(1)
303            .max(1)
304    }
305}
306
307/// State and file handle shared by every owner of one session in this process.
308///
309/// A keyed operation lock alone is insufficient: two independently opened
310/// handles could still carry stale turn counters and overwrite each other's
311/// metadata after taking the lock. Sharing the mutable state and append file
312/// makes the lock a true session boundary while preserving the value-type
313/// `SessionEventLog` API.
314struct SessionShared {
315    file: Mutex<Option<File>>,
316    state: Mutex<LogState>,
317    eviction_lock: Mutex<()>,
318    initialized: AtomicBool,
319    /// Exclusive flock on the session's `session.lock`, held for as long as
320    /// any handle to this session exists. Retention reads it as a liveness
321    /// signal so an open-but-idle session is never marked or evicted.
322    liveness_lock: Option<File>,
323}
324
325/// Return the process-wide shared state for one session's canonical event file.
326///
327/// The weak registry avoids retaining closed sessions forever while still
328/// making repeated `open` calls converge on one file handle and turn state.
329fn shared_session(events_path: &Path, session_id: &str) -> Result<Arc<SessionShared>, SessionStoreError> {
330    static SESSION_SHARED: OnceLock<Mutex<HashMap<PathBuf, Weak<SessionShared>>>> = OnceLock::new();
331
332    let key = events_path
333        .parent()
334        .and_then(|parent| vtcode_commons::paths::canonicalize(parent).ok())
335        .and_then(|parent| events_path.file_name().map(|name| parent.join(name)))
336        .unwrap_or_else(|| events_path.to_path_buf());
337    let registry = SESSION_SHARED.get_or_init(|| Mutex::new(HashMap::new()));
338    let mut shared_by_path = match registry.lock() {
339        Ok(locks) => locks,
340        Err(poisoned) => poisoned.into_inner(),
341    };
342    // Do not retain dead weak entries for every session ever opened by a
343    // long-running process. The registry is only an in-process coordination
344    // aid, so removing entries with no live owners is safe.
345    shared_by_path.retain(|_, shared| shared.strong_count() > 0);
346    if let Some(shared) = shared_by_path.get(&key).and_then(Weak::upgrade) {
347        return Ok(shared);
348    }
349
350    let file = VtCodePaths::open_private_append_file(events_path)
351        .map_err(|error| SessionStoreError::io(events_path.to_path_buf(), std::io::Error::other(error)))?;
352    let shared = Arc::new(SessionShared {
353        file: Mutex::new(Some(file)),
354        state: Mutex::new(LogState::new(session_id)),
355        eviction_lock: Mutex::new(()),
356        initialized: AtomicBool::new(false),
357        liveness_lock: acquire_liveness_lock(events_path),
358    });
359    shared_by_path.insert(key, Arc::downgrade(&shared));
360    Ok(shared)
361}
362
363/// Best-effort exclusive flock on the session's `session.lock`.
364///
365/// The lock is the liveness signal retention reads (`session_dir_is_live`): a
366/// held lock means a live process still has the session open. The crate has
367/// no logging surface and this is an advisory signal, not persistence, so a
368/// failure to create/open/lock degrades to unlocked (previous retention
369/// behavior) instead of failing the session open or swallowing a persistence
370/// error. The fd is released automatically when the shared handle drops.
371fn acquire_liveness_lock(events_path: &Path) -> Option<File> {
372    let dir = events_path.parent()?;
373    let lock_path = dir.join(SESSION_LOCK_FILE);
374    VtCodePaths::write_private_file_atomic_if_absent(&lock_path, b"").ok()?;
375    let file = std::fs::OpenOptions::new().read(true).write(true).open(&lock_path).ok()?;
376    match file.try_lock() {
377        Ok(()) => Some(file),
378        Err(_) => None,
379    }
380}
381
382/// Canonical append-only event log for a single session.
383///
384/// All session history is reconstructable from this log. Live conversation
385/// state is never read back into context from here; the log is only consumed
386/// for revert, compaction, analytics, and long-term-learning queries.
387pub struct SessionEventLog {
388    events_path: PathBuf,
389    manifest_store: ManifestStore,
390    shared: Arc<SessionShared>,
391    max_events: usize,
392    eviction_summary_hook: EvictionSummaryHook,
393}
394
395impl SessionEventLog {
396    /// Open the log for `session_id`, creating the session directory tree and
397    /// rebuilding the index from `events.jsonl` if it already exists.
398    pub(crate) fn open(workspace: &Path, session_id: &str, max_events: usize) -> Result<Self, SessionStoreError> {
399        let dir = session_dir(workspace, session_id);
400        let hook = default_eviction_summary_hook(dir.join(crate::DERIVED_DIR), session_id.to_string());
401        Self::open_with_eviction_summary(workspace, session_id, max_events, hook)
402    }
403
404    /// Open a log with an explicit eviction-summary callback.
405    ///
406    /// This is useful for hosts that keep derived memory in another store and
407    /// for deterministic failure-path tests. The callback must persist its
408    /// summary before returning `Ok(())`.
409    pub fn open_with_eviction_summary(
410        workspace: &Path,
411        session_id: &str,
412        max_events: usize,
413        eviction_summary_hook: EvictionSummaryHook,
414    ) -> Result<Self, SessionStoreError> {
415        let dir = session_dir(workspace, session_id);
416        crate::ensure_private_directory(&crate::sessions_root(workspace))?;
417        crate::ensure_private_directory(&dir)?;
418        crate::ensure_private_directory(&dir.join(crate::DERIVED_DIR))?;
419        crate::ensure_private_directory(&dir.join("index"))?;
420        let events_path = dir.join("events.jsonl");
421        let manifest_store = ManifestStore::new(dir.clone());
422        let pending_rewrite = manifest_store.load_pending_cap_rewrite()?;
423        let shared = shared_session(&events_path, session_id)?;
424        let log = Self {
425            events_path: events_path.clone(),
426            manifest_store,
427            shared,
428            max_events,
429            eviction_summary_hook,
430        };
431        let _eviction_guard = log.shared.eviction_lock.lock().map_err(poison)?;
432        // Try the fast path: read the persisted manifest + index and skip
433        // the O(n) scan when they are present and consistent.
434        if !log.shared.initialized.load(Ordering::Acquire) {
435            let manifest_opt = log.manifest_store.load_manifest()?;
436            let index_opt = log.manifest_store.load_turn_index()?;
437            let file_len = log.event_file_metadata_len()?;
438            let pending_rewrite_matches_file = pending_rewrite.as_ref().is_some_and(|pending| {
439                pending.new_file_len == file_len && pending.new_file_len < pending.previous_file_len
440            });
441            match (&manifest_opt, &index_opt) {
442                (Some(manifest), Some(index))
443                    if !pending_rewrite_matches_file
444                        && manifest.in_turn.is_some()
445                        && manifest.persisted_file_len == Some(file_len)
446                        && index.is_valid_for_file(file_len)
447                        && index.is_consistent_with_manifest(manifest) =>
448                {
449                    let mut st = log.shared.state.lock().map_err(poison)?;
450                    st.in_turn = manifest.in_turn.unwrap_or(false);
451                    st.manifest = manifest.clone();
452                    st.index = index.clone();
453                    st.next_offset = file_len;
454                }
455                _ => {
456                    let scan_turn_base = infer_scan_turn_base(
457                        manifest_opt.as_ref(),
458                        index_opt.as_ref(),
459                        pending_rewrite.as_ref().filter(|_| pending_rewrite_matches_file),
460                        file_len,
461                    );
462                    {
463                        let mut st = log.shared.state.lock().map_err(poison)?;
464                        if let Some(previous) = manifest_opt.as_ref() {
465                            st.manifest = previous.clone();
466                        }
467                        // The canonical event file is authoritative after any
468                        // stale/corrupt metadata. Preserve the ordinal of the
469                        // first retained turn while rebuilding all counters.
470                        st.manifest.turn_count = scan_turn_base.saturating_sub(1);
471                        st.manifest.event_count = 0;
472                        st.manifest.status = "active".to_string();
473                        st.manifest.in_turn = Some(false);
474                        st.manifest.retained_turn_base = Some(scan_turn_base);
475                        st.index = TurnIndex::default();
476                        st.in_turn = false;
477                    }
478                    log.scan()?;
479                    let mut st = log.shared.state.lock().map_err(poison)?;
480                    st.next_offset = file_len;
481                    log.persist_meta_locked(&mut st)?;
482                }
483            }
484            if pending_rewrite.is_some() {
485                // A marker whose file length did not match either side of the
486                // rewrite is stale, while a matching marker has now been
487                // incorporated into the repaired metadata. In both cases it
488                // is safe to remove it after the open path has persisted the
489                // authoritative state.
490                log.manifest_store.clear_pending_cap_rewrite()?;
491            }
492            log.shared.initialized.store(true, Ordering::Release);
493        }
494        drop(_eviction_guard);
495        Ok(log)
496    }
497
498    /// Append an event to the log and update the in-memory index/manifest.
499    pub fn append(&self, event: &ThreadEvent) -> Result<(), SessionStoreError> {
500        let _eviction_guard = self.shared.eviction_lock.lock().map_err(poison)?;
501        let mut st = self.shared.state.lock().map_err(poison)?;
502
503        // Serialize into the write buffer with rollback on failure — the
504        // invariant that `write_buf` never contains partial JSON is
505        // encapsulated in `serialize_event`.
506        let (start, end) = st.serialize_event(event)?;
507
508        st.manifest.event_count += 1;
509        st.manifest.updated_at = now_rfc3339();
510
511        // Route through the single turn-lifecycle state machine.  When the
512        // event closes a turn, persist metadata immediately so a reopen
513        // after a mid-turn crash sees a consistent index.
514        let is_turn_boundary = st.apply_lifecycle_event(LifecycleKind::from_event(event), start, end);
515        if is_turn_boundary {
516            self.persist_meta_locked(&mut st)?;
517        }
518
519        if st.write_buf.len() >= MAX_WRITE_BUFFER_BYTES {
520            // Persist metadata with the bounded byte flush so a reopen after
521            // a mid-turn crash does not trust an index that predates these
522            // already-written events.
523            self.persist_meta_locked(&mut st)?;
524        }
525        drop(st);
526        self.enforce_event_cap()
527    }
528
529    /// Enforce the per-session event cap by evicting the oldest completed
530    /// turns when the log exceeds [`Self::max_events`]. Returns `Ok(())` even
531    /// when no truncation is needed or the cap is disabled (`max_events == 0`).
532    fn enforce_event_cap(&self) -> Result<(), SessionStoreError> {
533        let mut st = self.shared.state.lock().map_err(poison)?;
534
535        // `plan_cap_eviction` encapsulates the index arithmetic and returns
536        // `None` when the cap is disabled or not yet exceeded.
537        let Some(plan) = st.plan_cap_eviction(self.max_events) else {
538            return Ok(());
539        };
540
541        // Keep ordinary appends in memory until a turn boundary or an
542        // explicit read. Cap enforcement is the one append-time path that
543        // needs the complete on-disk file before rewriting it.
544        self.flush_write_buf_locked(&mut st)?;
545
546        let (evicted, mut remaining, previous_file_len) = {
547            let mut file_slot = self.shared.file.lock().map_err(poison)?;
548            let file = file_slot.as_mut().ok_or_else(|| self.event_file_unavailable())?;
549            let file_len = file
550                .metadata()
551                .map_err(|error| SessionStoreError::io(&self.events_path, error))?
552                .len();
553            if plan.truncate_offset > file_len {
554                return Err(SessionStoreError::io(
555                    &self.events_path,
556                    std::io::Error::new(std::io::ErrorKind::InvalidData, "cap offset exceeds event log length"),
557                ));
558            }
559            file.seek(SeekFrom::Start(0))
560                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
561            let mut evicted = vec![
562                0u8;
563                usize::try_from(plan.truncate_offset).map_err(|error| {
564                    SessionStoreError::io(
565                        &self.events_path,
566                        std::io::Error::new(std::io::ErrorKind::InvalidData, error),
567                    )
568                })?
569            ];
570            file.read_exact(&mut evicted)
571                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
572            file.seek(SeekFrom::Start(plan.truncate_offset))
573                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
574            let mut remaining = Vec::new();
575            file.read_to_end(&mut remaining)
576                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
577            (evicted, remaining, file_len)
578        };
579        let retained_turn_base = st.retained_turn_base_after_eviction(plan.evicted_turn_count);
580        drop(st);
581
582        // The turn index counts only records inside indexed turns. The bytes
583        // removed by a cap rewrite may also contain valid session-level
584        // records (for example `thread.started`) before the first turn, so
585        // reconcile the manifest against the actual persisted prefix rather
586        // than the turn-only estimate from `plan_cap_eviction`.
587        let mut plan = plan;
588        plan.evicted_event_count = count_persisted_event_records(&evicted);
589        let evicted_events = decode_events(&evicted);
590        (self.eviction_summary_hook)(&evicted_events)?;
591
592        // Matrix snapshots are complete session-level checkpoints. Preserve the
593        // latest evicted checkpoint only when no newer snapshot survives.
594        let mut matrix_checkpoints = BTreeMap::new();
595        for event in &evicted_events {
596            if let ThreadEvent::MatrixUpdated(snapshot) = event {
597                matrix_checkpoints.insert(snapshot.spec.id.clone(), event);
598            }
599        }
600        for event in decode_events(&remaining) {
601            if let ThreadEvent::MatrixUpdated(snapshot) = event {
602                matrix_checkpoints.remove(&snapshot.spec.id);
603            }
604        }
605        let mut checkpoint_prefix = Vec::new();
606        for event in matrix_checkpoints.values() {
607            serde_json::to_writer(&mut checkpoint_prefix, &VersionedThreadEvent::new((*event).clone()))
608                .map_err(|error| SessionStoreError::io(&self.events_path, std::io::Error::other(error)))?;
609            checkpoint_prefix.push(b'\n');
610        }
611        let checkpoint_bytes = checkpoint_prefix.len() as u64;
612        let checkpoint_count = matrix_checkpoints.len() as u64;
613        if !checkpoint_prefix.is_empty() {
614            checkpoint_prefix.extend_from_slice(&remaining);
615            remaining = checkpoint_prefix;
616        }
617
618        let new_file_len = u64::try_from(remaining.len()).map_err(|error| {
619            SessionStoreError::io(&self.events_path, std::io::Error::new(std::io::ErrorKind::InvalidData, error))
620        })?;
621        self.manifest_store.write_pending_cap_rewrite(&PendingCapRewrite {
622            previous_file_len,
623            new_file_len,
624            retained_turn_base,
625        })?;
626        let next_offset = self.replace_event_file_contents(&remaining)?;
627        let mut st = self.shared.state.lock().map_err(poison)?;
628        st.apply_cap_eviction(plan, next_offset);
629        for entry in &mut st.index.entries {
630            entry.start_offset += checkpoint_bytes;
631            entry.end_offset += checkpoint_bytes;
632        }
633        st.manifest.event_count += checkpoint_count;
634        // The rewrite changed byte offsets and retained counts; persist the
635        // derived metadata before exposing the append as successful.
636        self.persist_meta_locked(&mut st)?;
637        self.manifest_store.clear_pending_cap_rewrite()?;
638        Ok(())
639    }
640
641    /// Reconstruct every event belonging to `turn`.
642    pub(crate) fn reconstruct_turn(&self, turn: u64) -> Result<Vec<ThreadEvent>, SessionStoreError> {
643        // Keep the index snapshot and byte-range read together with cap
644        // rewriting. Otherwise an eviction can replace the file between these
645        // steps and leave the snapshot offsets pointing into unrelated events.
646        let _eviction_guard = self.shared.eviction_lock.lock().map_err(poison)?;
647        let entry = {
648            let st = self.shared.state.lock().map_err(poison)?;
649            st.index
650                .entries
651                .iter()
652                .find(|e| e.turn_number == turn)
653                .cloned()
654                .ok_or(SessionStoreError::TurnNotFound { session: st.manifest.session_id.clone(), turn })?
655        };
656        {
657            let mut st = self.shared.state.lock().map_err(poison)?;
658            self.flush_write_buf_locked(&mut st)?;
659        }
660        let buf = {
661            let mut file_slot = self.shared.file.lock().map_err(poison)?;
662            let file = file_slot.as_mut().ok_or_else(|| self.event_file_unavailable())?;
663            file.seek(SeekFrom::Start(entry.start_offset))
664                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
665            let len = usize::try_from(entry.end_offset.checked_sub(entry.start_offset).ok_or_else(|| {
666                SessionStoreError::io(
667                    &self.events_path,
668                    std::io::Error::new(std::io::ErrorKind::InvalidData, "turn index offsets are out of order"),
669                )
670            })?)
671            .map_err(|error| {
672                SessionStoreError::io(&self.events_path, std::io::Error::new(std::io::ErrorKind::InvalidData, error))
673            })?;
674            let mut buf = vec![0u8; len];
675            file.read_exact(&mut buf)
676                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
677            buf
678        };
679        let text = String::from_utf8_lossy(&buf);
680        let mut events = Vec::new();
681        for line in text.lines() {
682            let line = line.trim();
683            if line.is_empty() {
684                continue;
685            }
686            // The index scan only validates the event envelope (plus the
687            // lifecycle shape) so it can rebuild cheaply. A line accepted by
688            // the scan can therefore still fail full decoding here; skip it
689            // instead of failing the whole reconstruction (revert, compaction,
690            // and analytics must not break on a single malformed record).
691            let v: VersionedThreadEvent = match serde_json::from_str(line) {
692                Ok(v) => v,
693                Err(_) => continue,
694            };
695            events.push(v.into_event());
696        }
697        Ok(events)
698    }
699
700    /// Number of turns recorded.
701    #[must_use]
702    pub(crate) fn turn_count(&self) -> u64 {
703        self.shared.state.lock().map_err(poison).map_or(0, |s| s.manifest.turn_count)
704    }
705
706    /// Number of events recorded.
707    #[must_use]
708    pub fn event_count(&self) -> u64 {
709        self.shared.state.lock().map_err(poison).map_or(0, |s| s.manifest.event_count)
710    }
711
712    /// Flush pending event bytes and metadata to the session store.
713    /// Visit one consistent retained range, including records outside turns.
714    /// Reads are streamed while cap rewrites and appends are excluded.
715    pub fn visit_snapshot<F>(&self, mut visitor: F) -> Result<SessionManifest, SessionStoreError>
716    where
717        F: FnMut(u64, &[u8]),
718    {
719        use std::io::{BufRead, BufReader};
720        let _eviction_guard = self.shared.eviction_lock.lock().map_err(poison)?;
721        let mut st = self.shared.state.lock().map_err(poison)?;
722        self.persist_meta_locked(&mut st)?;
723        let mut file_slot = self.shared.file.lock().map_err(poison)?;
724        let file = file_slot.as_mut().ok_or_else(|| self.event_file_unavailable())?;
725        file.seek(SeekFrom::Start(0))
726            .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
727        let mut reader = BufReader::new(file.take(st.next_offset));
728        let mut offset = 0;
729        let mut line = Vec::new();
730        loop {
731            line.clear();
732            let length = reader
733                .read_until(b'\n', &mut line)
734                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
735            if length == 0 {
736                break;
737            }
738            visitor(offset, &line);
739            offset += length as u64;
740        }
741        Ok(st.manifest.clone())
742    }
743
744    /// Flush pending event bytes and metadata to the session store.
745    pub fn flush(&self) -> Result<(), SessionStoreError> {
746        let _eviction_guard = self.shared.eviction_lock.lock().map_err(poison)?;
747        let mut st = self.shared.state.lock().map_err(poison)?;
748        self.persist_meta_locked(&mut st)
749    }
750
751    /// Snapshot of the session manifest.
752    #[must_use]
753    pub fn manifest(&self) -> SessionManifest {
754        self.shared
755            .state
756            .lock()
757            .map_err(poison)
758            .map(|s| s.manifest.clone())
759            .unwrap_or_else(|_| SessionManifest::new(""))
760    }
761
762    /// Snapshot of the turn index.
763    #[must_use]
764    pub fn turn_index(&self) -> TurnIndex {
765        self.shared
766            .state
767            .lock()
768            .map_err(poison)
769            .map(|s| s.index.clone())
770            .unwrap_or_default()
771    }
772
773    /// Flush metadata for callers that explicitly close a log handle.
774    ///
775    /// Terminal status is intentionally controlled only by a persisted
776    /// `thread.completed` event. This method does not synthesize lifecycle
777    /// state for callers that merely release a store handle.
778    pub(crate) fn complete(&self) -> Result<(), SessionStoreError> {
779        let _eviction_guard = self.shared.eviction_lock.lock().map_err(poison)?;
780        let mut st = self.shared.state.lock().map_err(poison)?;
781        st.manifest.updated_at = now_rfc3339();
782        self.persist_meta_locked(&mut st)
783    }
784
785    /// Rebuild index + manifest by scanning `events.jsonl` (authoritative).
786    ///
787    /// Reads the file line-by-line via `BufReader` to avoid loading the entire
788    /// log into memory. Long-lived sessions can otherwise produce multi-megabyte
789    /// logs that spike memory on every reopen.
790    fn scan(&self) -> Result<(), SessionStoreError> {
791        let mut st = self.shared.state.lock().map_err(poison)?;
792        let file = self
793            .shared
794            .file
795            .lock()
796            .map_err(poison)?
797            .as_ref()
798            .ok_or_else(|| self.event_file_unavailable())?
799            .try_clone()
800            .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
801        let mut reader = std::io::BufReader::new(file);
802        reader
803            .seek(SeekFrom::Start(0))
804            .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
805        let mut buf = Vec::new();
806        let mut pos = 0u64;
807        let mut first_ts: Option<String> = None;
808        loop {
809            buf.clear();
810            let n = reader
811                .read_until(b'\n', &mut buf)
812                .map_err(|e| SessionStoreError::io(&self.events_path, e))?;
813            if n == 0 {
814                break;
815            }
816            let line_end = pos + n as u64;
817            let trimmed = std::str::from_utf8(&buf).unwrap_or("").trim();
818            if !trimmed.is_empty()
819                && let Ok(v) = serde_json::from_str::<VersionedEventKind<'_>>(trimmed)
820            {
821                let kind = v.event.kind;
822                if requires_full_lifecycle_validation(kind) && !valid_lifecycle_payload(kind, trimmed) {
823                    pos = line_end;
824                    continue;
825                }
826                st.manifest.event_count += 1;
827                // `thread.started` is not part of the turn lifecycle — it
828                // only seeds `created_at` on the first occurrence.
829                if kind == "thread.started" && first_ts.is_none() {
830                    first_ts = Some(now_rfc3339());
831                }
832                // Route turn-lifecycle events through the same state machine
833                // as `append`, eliminating a previously duplicated match block.
834                st.apply_lifecycle_event(LifecycleKind::from_kind(kind), pos, line_end);
835            }
836            pos = line_end;
837        }
838        // Keep the open-turn state reconstructed from the canonical event log.
839        // This lets a reopened session continue a turn that was flushed before
840        // its completion event was written.
841        st.manifest.in_turn = Some(st.in_turn);
842        if let Some(ts) = first_ts
843            && st.manifest.created_at.is_empty()
844        {
845            st.manifest.created_at = ts;
846        }
847        Ok(())
848    }
849
850    fn persist_meta_locked(&self, st: &mut LogState) -> Result<(), SessionStoreError> {
851        self.flush_write_buf_locked(st)?;
852        // The manifest is only eligible for the fast reopen path when it
853        // describes the complete on-disk event file. Drop intentionally
854        // flushes bytes without metadata, so a length mismatch safely forces
855        // the authoritative scan on the next open.
856        st.manifest.persisted_file_len = Some(st.next_offset);
857        // Publish the derived index first. If a process stops between these
858        // two atomic renames, the older manifest still carries a stale file
859        // length and forces a scan instead of allowing the new manifest to
860        // pair with an older, apparently valid index.
861        self.manifest_store.write_turn_index(&st.index)?;
862        self.manifest_store.write_manifest(&st.manifest)?;
863        Ok(())
864    }
865
866    /// Flush the in-memory write buffer to the underlying file.
867    fn flush_write_buf_locked(&self, st: &mut LogState) -> Result<(), SessionStoreError> {
868        if st.write_buf.is_empty() {
869            return Ok(());
870        }
871        let mut file_slot = self.shared.file.lock().map_err(poison)?;
872        let file = file_slot.as_mut().ok_or_else(|| self.event_file_unavailable())?;
873        let previous_len = file.metadata().map_err(|e| SessionStoreError::io(&self.events_path, e))?.len();
874        if let Err(error) = file.write_all(&st.write_buf) {
875            if file.set_len(previous_len).is_err() {
876                st.write_buf.clear();
877            }
878            return Err(SessionStoreError::io(&self.events_path, error));
879        }
880        if let Err(error) = file.sync_data() {
881            st.write_buf.clear();
882            return Err(SessionStoreError::io(&self.events_path, error));
883        }
884        st.write_buf.clear();
885        Ok(())
886    }
887
888    fn event_file_metadata_len(&self) -> Result<u64, SessionStoreError> {
889        let file_slot = self.shared.file.lock().map_err(poison)?;
890        file_slot
891            .as_ref()
892            .ok_or_else(|| self.event_file_unavailable())?
893            .metadata()
894            .map(|metadata| metadata.len())
895            .map_err(|error| SessionStoreError::io(&self.events_path, error))
896    }
897
898    fn open_event_file(&self) -> Result<File, SessionStoreError> {
899        VtCodePaths::open_private_append_file(&self.events_path)
900            .map_err(|error| SessionStoreError::io(&self.events_path, std::io::Error::other(error)))
901    }
902
903    fn replace_event_file_contents(&self, contents: &[u8]) -> Result<u64, SessionStoreError> {
904        let old_file = {
905            let mut file_slot = self.shared.file.lock().map_err(poison)?;
906            file_slot.take().ok_or_else(|| self.event_file_unavailable())?
907        };
908        drop(old_file);
909
910        if let Err(error) = VtCodePaths::write_private_file_atomic(&self.events_path, contents)
911            .map_err(|error| SessionStoreError::io(&self.events_path, std::io::Error::other(error)))
912        {
913            let restored = self.open_event_file();
914            if let Ok(file) = restored {
915                let mut file_slot = self.shared.file.lock().map_err(poison)?;
916                *file_slot = Some(file);
917                return Err(error);
918            }
919            return Err(SessionStoreError::io(
920                &self.events_path,
921                std::io::Error::other(format!("{error}; failed to restore event log handle")),
922            ));
923        }
924
925        let replacement = self.open_event_file()?;
926        let next_offset = replacement
927            .metadata()
928            .map_err(|error| SessionStoreError::io(&self.events_path, error))?
929            .len();
930        let mut file_slot = self.shared.file.lock().map_err(poison)?;
931        *file_slot = Some(replacement);
932        Ok(next_offset)
933    }
934
935    fn event_file_unavailable(&self) -> SessionStoreError {
936        SessionStoreError::io(&self.events_path, std::io::Error::other("event log file is unavailable"))
937    }
938}
939
940/// Recover the ordinal of the first retained turn when metadata is stale.
941///
942/// A cap rewrite atomically replaces the event file before it publishes the
943/// shortened index and manifest. If the process crashes in that interval,
944/// the old index still records the pre-rewrite offsets. The difference between
945/// its persisted length and the current file length is exactly the removed
946/// prefix, so the first old index entry at that boundary supplies the retained
947/// turn base. Other stale metadata falls back to the durable base field.
948fn infer_scan_turn_base(
949    manifest: Option<&SessionManifest>,
950    index: Option<&TurnIndex>,
951    pending: Option<&PendingCapRewrite>,
952    file_len: u64,
953) -> u64 {
954    let marker_base = pending
955        .filter(|pending| pending.new_file_len == file_len && pending.new_file_len < pending.previous_file_len)
956        .map(|pending| pending.retained_turn_base);
957    let rewritten_base = manifest.and_then(|manifest| {
958        let previous_len = manifest.persisted_file_len?;
959        if previous_len <= file_len {
960            return None;
961        }
962        let removed_prefix = previous_len - file_len;
963        let index = index.filter(|index| index.is_valid_for_file(previous_len))?;
964        let first_retained = index
965            .entries
966            .iter()
967            .find(|entry| entry.start_offset >= removed_prefix && entry.end_offset <= previous_len)
968            .map(|entry| entry.turn_number);
969        first_retained.or_else(|| {
970            // If no indexed turn starts in the shortened file, the rewrite
971            // evicted every previously indexed turn. Preserve the next
972            // ordinal from the stale manifest so a subsequent append cannot
973            // reuse an already-observed turn number.
974            (removed_prefix >= previous_len).then(|| manifest.turn_count.saturating_add(1))
975        })
976    });
977
978    let legacy_index_base = index
979        .filter(|index| index.is_valid_for_file(file_len))
980        .and_then(|index| index.entries.front().map(|entry| entry.turn_number));
981
982    marker_base
983        .or(rewritten_base)
984        .or_else(|| manifest.and_then(|manifest| manifest.retained_turn_base))
985        .or(legacy_index_base)
986        .unwrap_or(1)
987        .max(1)
988}
989
990fn requires_full_lifecycle_validation(kind: &str) -> bool {
991    matches!(kind, "thread.started" | "thread.completed" | "turn.started" | "turn.completed" | "turn.failed")
992}
993
994fn valid_lifecycle_payload(kind: &str, line: &str) -> bool {
995    if serde_json::from_str::<VersionedThreadEvent>(line).is_err() {
996        return false;
997    }
998    if kind != "turn.completed" {
999        return true;
1000    }
1001    let Ok(value) = serde_json::from_str::<serde_json::Value>(line) else {
1002        return false;
1003    };
1004    value
1005        .get("event")
1006        .and_then(|event| event.get("usage"))
1007        .is_some_and(serde_json::Value::is_object)
1008}
1009
1010fn decode_events(bytes: &[u8]) -> Vec<ThreadEvent> {
1011    bytes
1012        .split(|byte| *byte == b'\n')
1013        .filter_map(|line| {
1014            let line = std::str::from_utf8(line).ok()?.trim();
1015            if line.is_empty() {
1016                return None;
1017            }
1018            serde_json::from_str::<VersionedThreadEvent>(line)
1019                .ok()
1020                .map(VersionedThreadEvent::into_event)
1021        })
1022        .collect()
1023}
1024
1025/// Count records that the authoritative scan would include in the manifest.
1026/// This deliberately parses the lightweight event envelope instead of using
1027/// the turn index: a cap rewrite can remove session-level records that never
1028/// belong to an indexed turn.
1029fn count_persisted_event_records(bytes: &[u8]) -> u64 {
1030    bytes
1031        .split(|byte| *byte == b'\n')
1032        .filter_map(|line| std::str::from_utf8(line).ok())
1033        .map(str::trim)
1034        .filter(|line| !line.is_empty())
1035        .filter(|line| {
1036            let Ok(value) = serde_json::from_str::<VersionedEventKind<'_>>(line) else {
1037                return false;
1038            };
1039            let kind = value.event.kind;
1040            !requires_full_lifecycle_validation(kind) || valid_lifecycle_payload(kind, line)
1041        })
1042        .count() as u64
1043}
1044
1045#[derive(Debug, Serialize)]
1046struct EvictionSummary {
1047    session_id: String,
1048    evicted_event_count: usize,
1049    event_types: BTreeMap<String, u64>,
1050    grounded_facts: Vec<String>,
1051    created_at: String,
1052}
1053
1054/// Extract a bounded, deterministic set of facts from canonical event
1055/// payloads. This is intentionally structural rather than model-generated:
1056/// eviction must remain synchronous, reproducible, and safe when the model is
1057/// unavailable. Only completed item snapshots and terminal thread errors are
1058/// considered, so streaming deltas and raw tool output cannot flood the
1059/// derived summary.
1060fn extract_grounded_facts(events: &[ThreadEvent]) -> Vec<String> {
1061    let mut facts = Vec::new();
1062    for event in events {
1063        let candidates: Vec<String> = match event {
1064            ThreadEvent::ItemCompleted(completed) => item_facts(&completed.item.details),
1065            ThreadEvent::TurnFailed(failed) => vec![failed.message.clone()],
1066            ThreadEvent::TurnBlocked(blocked) => vec![blocked.message.clone()],
1067            ThreadEvent::Error(error) => vec![error.message.clone()],
1068            _ => Vec::new(),
1069        };
1070        for candidate in candidates {
1071            let fact = normalize_eviction_fact(&candidate);
1072            if fact.is_empty() || facts.iter().any(|existing| existing == &fact) {
1073                continue;
1074            }
1075            facts.push(fact);
1076            if facts.len() == MAX_EVICTION_GROUNDED_FACTS {
1077                return facts;
1078            }
1079        }
1080    }
1081    facts
1082}
1083
1084fn item_facts(details: &ThreadItemDetails) -> Vec<String> {
1085    match details {
1086        ThreadItemDetails::AgentMessage(item) => vec![item.text.clone()],
1087        ThreadItemDetails::Plan(item) => vec![item.text.clone()],
1088        ThreadItemDetails::FileChange(item) => item
1089            .changes
1090            .iter()
1091            .map(|change| {
1092                let kind = match change.kind {
1093                    vtcode_exec_events::PatchChangeKind::Add => "add",
1094                    vtcode_exec_events::PatchChangeKind::Delete => "delete",
1095                    vtcode_exec_events::PatchChangeKind::Update => "update",
1096                };
1097                format!("file {kind}: {}", change.path)
1098            })
1099            .collect(),
1100        ThreadItemDetails::Harness(item) => item.message.clone().into_iter().collect(),
1101        _ => Vec::new(),
1102    }
1103}
1104
1105fn normalize_eviction_fact(value: &str) -> String {
1106    let normalized = value.split_whitespace().collect::<Vec<_>>().join(" ");
1107    if normalized.len() <= MAX_EVICTION_GROUNDED_FACT_BYTES {
1108        return normalized;
1109    }
1110    let mut end = MAX_EVICTION_GROUNDED_FACT_BYTES - '…'.len_utf8();
1111    while !normalized.is_char_boundary(end) {
1112        end -= 1;
1113    }
1114    format!("{}…", &normalized[..end])
1115}
1116
1117fn default_eviction_summary_hook(derived_dir: PathBuf, session_id: String) -> EvictionSummaryHook {
1118    Arc::new(move |events| {
1119        let mut event_types = BTreeMap::new();
1120        for event in events {
1121            let kind = serde_json::to_value(event)
1122                .ok()
1123                .and_then(|value| value.get("type").and_then(|value| value.as_str()).map(str::to_owned))
1124                .unwrap_or_else(|| "unknown".to_owned());
1125            *event_types.entry(kind).or_insert(0) += 1;
1126        }
1127        let summary = EvictionSummary {
1128            session_id: session_id.clone(),
1129            evicted_event_count: events.len(),
1130            event_types,
1131            grounded_facts: extract_grounded_facts(events),
1132            created_at: now_rfc3339(),
1133        };
1134        let path = derived_dir.join(format!("eviction-summary-{}.json", uuid::Uuid::new_v4().simple()));
1135        let bytes = serde_json::to_vec(&summary)?;
1136        VtCodePaths::write_private_file_atomic(&path, &bytes)
1137            .map_err(|error| SessionStoreError::io(path, std::io::Error::other(error)))
1138    })
1139}
1140
1141impl Drop for SessionEventLog {
1142    fn drop(&mut self) {
1143        if let Ok(_eviction_guard) = self.shared.eviction_lock.lock()
1144            && let Ok(mut st) = self.shared.state.lock()
1145        {
1146            // The fallible `flush` method is the authoritative shutdown path;
1147            // Drop only provides a best-effort byte flush for callers that do
1148            // not explicitly close the log. Rewriting metadata here could
1149            // overwrite a manifest update made by another owner after the
1150            // last append.
1151            let _ = self.flush_write_buf_locked(&mut st);
1152        }
1153    }
1154}
1155
1156/// Locate the next newline at or after `from`, returning a past-the-end index.
1157fn poison<T>(_e: std::sync::PoisonError<T>) -> SessionStoreError {
1158    SessionStoreError::Io {
1159        path: PathBuf::new(),
1160        source: std::io::Error::other("session store lock poisoned"),
1161    }
1162}
1163
1164fn now_rfc3339() -> String {
1165    Utc::now().to_rfc3339()
1166}
1167
1168/// Session-level metadata persisted to `manifest.json`.
1169#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1170pub struct SessionManifest {
1171    /// Stable session identifier (directory name).
1172    pub session_id: String,
1173    /// Layout schema version (`SESSION_STORE_SCHEMA_VERSION`).
1174    schema_version: u32,
1175    /// RFC3339 creation timestamp.
1176    pub created_at: String,
1177    /// RFC3339 last-update timestamp.
1178    pub updated_at: String,
1179    /// Number of completed turns.
1180    pub turn_count: u64,
1181    /// Total number of events recorded.
1182    pub event_count: u64,
1183    /// Lifecycle status (`active` | `completed`).
1184    pub status: String,
1185    /// Whether the canonical log currently ends inside an open turn.
1186    ///
1187    /// This is optional on read so manifests written before open-turn
1188    /// persistence was introduced trigger a safe event-log scan instead of
1189    /// silently losing lifecycle state.
1190    #[serde(default)]
1191    in_turn: Option<bool>,
1192    /// Byte length covered by the persisted manifest and turn index.
1193    ///
1194    /// This is optional for compatibility with manifests written before the
1195    /// fast-path freshness guard existed; those manifests are rebuilt from the
1196    /// canonical event log on reopen.
1197    #[serde(default)]
1198    persisted_file_len: Option<u64>,
1199    /// Ordinal of the first turn retained in the canonical log.
1200    ///
1201    /// Cap eviction removes completed turns but must keep later turn numbers
1202    /// monotonic. The field lets an authoritative scan restore those ordinals
1203    /// even when the derived index is stale or missing.
1204    #[serde(default)]
1205    retained_turn_base: Option<u64>,
1206}
1207
1208impl SessionManifest {
1209    /// Number of turns preceding the retained canonical range.
1210    #[must_use]
1211    pub fn evicted_turn_count(&self) -> u64 {
1212        self.retained_turn_base.unwrap_or(1).saturating_sub(1)
1213    }
1214    /// Create a fresh manifest for a session.
1215    #[must_use]
1216    pub(crate) fn new(session_id: &str) -> Self {
1217        let ts = now_rfc3339();
1218        Self {
1219            session_id: session_id.to_string(),
1220            schema_version: crate::SESSION_STORE_SCHEMA_VERSION,
1221            created_at: ts.clone(),
1222            updated_at: ts,
1223            turn_count: 0,
1224            event_count: 0,
1225            status: "active".to_string(),
1226            in_turn: Some(false),
1227            persisted_file_len: Some(0),
1228            retained_turn_base: Some(1),
1229        }
1230    }
1231}
1232
1233/// Byte-offset index of a single turn within `events.jsonl`.
1234#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1235pub struct TurnIndexEntry {
1236    /// Turn ordinal (1-based).
1237    turn_number: u64,
1238    /// Byte offset of the turn's first event.
1239    start_offset: u64,
1240    /// Byte offset just past the turn's last event.
1241    end_offset: u64,
1242    /// Number of events in the turn.
1243    event_count: u64,
1244    /// RFC3339 timestamp of turn start.
1245    ts: String,
1246}
1247
1248/// Ordered index of all turns in a session.
1249#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
1250pub struct TurnIndex {
1251    /// Turn entries in ordinal order.
1252    entries: VecDeque<TurnIndexEntry>,
1253}
1254
1255impl TurnIndex {
1256    /// Number of indexed turns.
1257    #[must_use]
1258    pub fn len(&self) -> usize {
1259        self.entries.len()
1260    }
1261
1262    /// Whether the index is empty.
1263    #[must_use]
1264    pub fn is_empty(&self) -> bool {
1265        self.entries.is_empty()
1266    }
1267
1268    fn is_valid_for_file(&self, file_len: u64) -> bool {
1269        let mut previous_end = 0u64;
1270        self.entries.iter().all(|entry| {
1271            let valid = entry.event_count > 0
1272                && entry.start_offset >= previous_end
1273                && entry.start_offset <= entry.end_offset
1274                && entry.end_offset <= file_len;
1275            if valid {
1276                previous_end = entry.end_offset;
1277            }
1278            valid
1279        })
1280    }
1281
1282    fn is_consistent_with_manifest(&self, manifest: &SessionManifest) -> bool {
1283        let expected_last_turn = if manifest.in_turn == Some(true) {
1284            manifest.turn_count.saturating_add(1)
1285        } else {
1286            manifest.turn_count
1287        };
1288        let entries_are_contiguous = self
1289            .entries
1290            .iter()
1291            .map(|entry| entry.turn_number)
1292            .try_fold(None::<u64>, |previous, turn_number| {
1293                if previous.is_some_and(|previous| turn_number != previous.saturating_add(1)) {
1294                    return Err(());
1295                }
1296                Ok(Some(turn_number))
1297            })
1298            .is_ok();
1299        if !entries_are_contiguous {
1300            return false;
1301        }
1302
1303        let expected_first_turn = manifest.retained_turn_base.unwrap_or(1).max(1);
1304        match (self.entries.front(), self.entries.back()) {
1305            (Some(first), Some(last)) => {
1306                first.turn_number == expected_first_turn && last.turn_number == expected_last_turn
1307            }
1308            (None, None) => {
1309                expected_last_turn == 0
1310                    || manifest
1311                        .retained_turn_base
1312                        .is_some_and(|retained_turn_base| retained_turn_base > manifest.turn_count)
1313            }
1314            _ => false,
1315        }
1316    }
1317}
1318
1319#[cfg(test)]
1320mod borrowed_envelope_tests {
1321    use super::{BorrowedVersionedEvent, EVENT_SCHEMA_VERSION};
1322    use vtcode_exec_events::{
1323        ThreadEvent, ThreadStartedEvent, TurnCompletedEvent, TurnStartedEvent, Usage, VersionedThreadEvent,
1324    };
1325
1326    /// The borrowed envelope must produce JSON byte-identical to
1327    /// `VersionedThreadEvent::new(event.clone())`. This guards against drift if
1328    /// either the envelope or the canonical wrapper is modified.
1329    #[test]
1330    fn borrowed_envelope_matches_versioned_envelope() {
1331        for event in [
1332            ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread".to_string() }),
1333            ThreadEvent::TurnStarted(TurnStartedEvent::default()),
1334            ThreadEvent::TurnCompleted(TurnCompletedEvent {
1335                completed_at: None,
1336                usage: Usage::default(),
1337                in_progress_exec_sessions: Vec::new(),
1338            }),
1339        ] {
1340            let canonical =
1341                serde_json::to_string(&VersionedThreadEvent::new(event.clone())).expect("canonical serialize");
1342            let borrowed = serde_json::to_string(&BorrowedVersionedEvent {
1343                schema_version: EVENT_SCHEMA_VERSION,
1344                event: &event,
1345            })
1346            .expect("borrowed serialize");
1347            assert_eq!(canonical, borrowed, "JSON differs for {event:?}");
1348        }
1349    }
1350}
1351
1352#[cfg(test)]
1353mod lifecycle_state_machine_tests {
1354    use super::{LifecycleKind, LogState};
1355    use vtcode_exec_events::{
1356        ThreadCompletedEvent, ThreadCompletionSubtype, ThreadEvent, ThreadStartedEvent, TurnCompletedEvent,
1357        TurnFailedEvent, TurnStartedEvent, Usage,
1358    };
1359
1360    fn fresh_state() -> LogState {
1361        LogState::new("test-session")
1362    }
1363
1364    #[test]
1365    fn lifecycle_kind_from_event_covers_all_variants() {
1366        assert_eq!(
1367            LifecycleKind::from_event(&ThreadEvent::TurnStarted(TurnStartedEvent::default())),
1368            LifecycleKind::TurnStarted
1369        );
1370        assert_eq!(
1371            LifecycleKind::from_event(&ThreadEvent::TurnCompleted(TurnCompletedEvent {
1372                completed_at: None,
1373                usage: Usage::default(),
1374                in_progress_exec_sessions: Vec::new()
1375            })),
1376            LifecycleKind::TurnCompleted
1377        );
1378        assert_eq!(
1379            LifecycleKind::from_event(&ThreadEvent::TurnFailed(TurnFailedEvent {
1380                completed_at: None,
1381                message: "err".to_string(),
1382                usage: None,
1383            })),
1384            LifecycleKind::TurnFailed
1385        );
1386        assert_eq!(
1387            LifecycleKind::from_event(&ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "x".to_string() })),
1388            LifecycleKind::ThreadStarted
1389        );
1390        assert_eq!(
1391            LifecycleKind::from_event(&ThreadEvent::ThreadCompleted(Box::new(ThreadCompletedEvent {
1392                completed_at: None,
1393                thread_id: "x".to_string(),
1394                session_id: "x".to_string(),
1395                subtype: ThreadCompletionSubtype::Success,
1396                outcome_code: "completed".to_string(),
1397                result: None,
1398                stop_reason: None,
1399                usage: Usage::default(),
1400                total_cost_usd: None,
1401                num_turns: 1,
1402            }))),
1403            LifecycleKind::ThreadCompleted
1404        );
1405        assert_eq!(LifecycleKind::from_kind("thread.started"), LifecycleKind::ThreadStarted);
1406        assert_eq!(LifecycleKind::from_kind("thread.completed"), LifecycleKind::ThreadCompleted);
1407    }
1408
1409    #[test]
1410    fn lifecycle_kind_from_str_matches_event_discriminator() {
1411        assert_eq!(LifecycleKind::from_kind("turn.started"), LifecycleKind::TurnStarted);
1412        assert_eq!(LifecycleKind::from_kind("turn.completed"), LifecycleKind::TurnCompleted);
1413        assert_eq!(LifecycleKind::from_kind("turn.failed"), LifecycleKind::TurnFailed);
1414        assert_eq!(LifecycleKind::from_kind("tool.called"), LifecycleKind::Other);
1415        assert_eq!(LifecycleKind::from_kind("thread.started"), LifecycleKind::ThreadStarted);
1416        assert_eq!(LifecycleKind::from_kind("thread.completed"), LifecycleKind::ThreadCompleted);
1417    }
1418
1419    #[test]
1420    fn turn_started_pushes_index_entry_and_sets_in_turn() {
1421        let mut st = fresh_state();
1422        st.manifest.status = "completed".to_string();
1423        let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
1424        assert!(!is_boundary, "TurnStarted is not a turn boundary");
1425        assert!(st.in_turn);
1426        assert_eq!(st.manifest.status, "active");
1427        assert_eq!(st.index.entries.len(), 1);
1428        let entry = &st.index.entries[0];
1429        assert_eq!(entry.turn_number, 1);
1430        assert_eq!(entry.start_offset, 0);
1431        assert_eq!(entry.end_offset, 100);
1432        assert_eq!(entry.event_count, 1);
1433    }
1434
1435    #[test]
1436    fn intermediate_events_extend_current_turn() {
1437        let mut st = fresh_state();
1438        st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
1439        // Simulate two intermediate events.
1440        let is_b1 = st.apply_lifecycle_event(LifecycleKind::Other, 100, 200);
1441        let is_b2 = st.apply_lifecycle_event(LifecycleKind::Other, 200, 300);
1442        assert!(!is_b1 && !is_b2);
1443        assert!(st.in_turn);
1444        assert_eq!(st.index.entries.len(), 1);
1445        let entry = &st.index.entries[0];
1446        assert_eq!(entry.end_offset, 300);
1447        assert_eq!(entry.event_count, 3);
1448    }
1449
1450    #[test]
1451    fn turn_completed_closes_turn_and_returns_boundary() {
1452        let mut st = fresh_state();
1453        st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
1454        st.apply_lifecycle_event(LifecycleKind::Other, 100, 200);
1455        let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnCompleted, 200, 300);
1456        assert!(is_boundary);
1457        assert!(!st.in_turn);
1458        assert_eq!(st.manifest.turn_count, 1);
1459        assert_eq!(st.manifest.status, "active");
1460        st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 300, 400);
1461        assert_eq!(st.manifest.status, "completed");
1462        let entry = &st.index.entries[0];
1463        assert_eq!(entry.end_offset, 300);
1464        assert_eq!(entry.event_count, 3);
1465    }
1466
1467    #[test]
1468    fn turn_failed_closes_turn_without_terminal_thread_status() {
1469        let mut st = fresh_state();
1470        st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
1471        let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnFailed, 100, 200);
1472        assert!(is_boundary);
1473        assert!(!st.in_turn);
1474        assert_eq!(st.manifest.turn_count, 1);
1475        assert_eq!(st.manifest.status, "active");
1476        st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 200, 300);
1477        assert_eq!(st.manifest.status, "completed");
1478    }
1479
1480    #[test]
1481    fn turn_completed_without_turn_started_is_idempotent() {
1482        let mut st = fresh_state();
1483        // Receiving TurnCompleted without a preceding TurnStarted should not
1484        // panic or corrupt the index; terminal status remains active until the
1485        // thread lifecycle itself completes.
1486        let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnCompleted, 0, 100);
1487        assert!(is_boundary);
1488        assert!(!st.in_turn);
1489        assert_eq!(st.manifest.turn_count, 0, "no turn was started");
1490        assert_eq!(st.manifest.status, "active");
1491        st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 100, 200);
1492        assert_eq!(st.manifest.status, "completed");
1493        assert!(st.index.entries.is_empty());
1494    }
1495
1496    #[test]
1497    fn multiple_turns_get_incrementing_ordinals() {
1498        let mut st = fresh_state();
1499        for n in 1..=3 {
1500            st.apply_lifecycle_event(LifecycleKind::TurnStarted, n * 100, n * 100 + 50);
1501            st.apply_lifecycle_event(LifecycleKind::TurnCompleted, n * 100 + 50, n * 100 + 100);
1502        }
1503        assert_eq!(st.index.entries.len(), 3);
1504        for (i, entry) in st.index.entries.iter().enumerate() {
1505            assert_eq!(entry.turn_number, (i + 1) as u64);
1506        }
1507        assert_eq!(st.manifest.turn_count, 3);
1508    }
1509}
1510
1511#[cfg(test)]
1512mod cap_eviction_tests {
1513    use super::{LogState, TurnIndexEntry};
1514
1515    /// Build a `LogState` with `turns` fake turns, each having `events_per_turn`
1516    /// events, starting at byte offset 0.
1517    fn state_with_turns(turns: usize, events_per_turn: u64) -> LogState {
1518        let mut st = LogState::new("cap-test");
1519        st.manifest.event_count = (turns as u64) * events_per_turn;
1520        let mut offset = 0u64;
1521        for n in 1..=turns {
1522            st.index.entries.push_back(TurnIndexEntry {
1523                turn_number: n as u64,
1524                start_offset: offset,
1525                end_offset: offset + events_per_turn * 10,
1526                event_count: events_per_turn,
1527                ts: "2026-01-01T00:00:00Z".to_string(),
1528            });
1529            offset += events_per_turn * 10;
1530        }
1531        st
1532    }
1533
1534    #[test]
1535    fn no_eviction_when_under_cap() {
1536        let st = state_with_turns(3, 2); // 6 events
1537        assert!(st.plan_cap_eviction(10).is_none());
1538        assert_eq!(st.index.entries.len(), 3, "no turns should be evicted");
1539    }
1540
1541    #[test]
1542    fn no_eviction_when_cap_disabled() {
1543        let st = state_with_turns(5, 2); // 10 events
1544        assert!(st.plan_cap_eviction(0).is_none());
1545        assert_eq!(st.index.entries.len(), 5);
1546    }
1547
1548    #[test]
1549    fn evicts_oldest_turns_to_meet_cap() {
1550        // 5 turns × 2 events = 10 events; cap = 6 → need to evict 2 turns (4 events).
1551        let st = state_with_turns(5, 2);
1552        let plan = st.plan_cap_eviction(6).expect("eviction planned");
1553        assert_eq!(plan.evicted_event_count, 4, "should evict 4 events (2 turns)");
1554        assert_eq!(st.index.entries.len(), 5, "planning must not mutate state");
1555        // Truncate offset is the end of the last evicted turn.
1556        assert_eq!(plan.truncate_offset, 40); // 2 turns × 20 bytes each
1557
1558        // Applying the plan leaves turns 3, 4, 5.
1559        let mut st = st;
1560        st.apply_cap_eviction(plan, 60);
1561        assert_eq!(st.index.entries.len(), 3, "should keep 3 turns");
1562        assert_eq!(st.index.entries[0].turn_number, 3);
1563        assert_eq!(st.index.entries[2].turn_number, 5);
1564    }
1565
1566    #[test]
1567    fn evicts_all_turns_when_cap_smaller_than_one_turn() {
1568        // 3 turns × 5 events = 15 events; cap = 3 → evict turns until ≤ 3 remain.
1569        // Each turn has 5 events, so evicting 2 turns leaves 5 (>3), evicting
1570        // 3 turns leaves 0.
1571        let st = state_with_turns(3, 5);
1572        let plan = st.plan_cap_eviction(3).expect("eviction planned");
1573        assert_eq!(plan.evicted_event_count, 15, "all events evicted");
1574        let mut st = st;
1575        st.apply_cap_eviction(plan, 0);
1576        assert_eq!(st.index.entries.len(), 0);
1577    }
1578}