use std::collections::VecDeque;
use std::fs::{File, OpenOptions};
use std::io::{BufRead, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Mutex;
use chrono::Utc;
use serde::{Deserialize, Serialize};
use vtcode_exec_events::{EVENT_SCHEMA_VERSION, ThreadEvent, VersionedThreadEvent};
use crate::error::SessionStoreError;
use crate::manifest::ManifestStore;
use crate::session_dir;
pub const DEFAULT_MAX_EVENTS: usize = 10_000;
const MAX_WRITE_BUFFER_BYTES: usize = 64 * 1024;
#[derive(Debug, Deserialize)]
struct VersionedEventKind<'a> {
#[serde(rename = "schema_version", borrow)]
_schema_version: &'a str,
#[serde(borrow)]
event: EventKind<'a>,
}
#[derive(Debug, Deserialize)]
struct EventKind<'a> {
#[serde(rename = "type", borrow)]
kind: &'a str,
}
#[derive(Serialize)]
struct BorrowedVersionedEvent<'a> {
schema_version: &'a str,
event: &'a ThreadEvent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum LifecycleKind {
ThreadStarted,
ThreadCompleted,
TurnStarted,
TurnCompleted,
TurnFailed,
Other,
}
impl LifecycleKind {
#[inline]
fn from_event(event: &ThreadEvent) -> Self {
match event {
ThreadEvent::ThreadStarted(_) => Self::ThreadStarted,
ThreadEvent::ThreadCompleted(_) => Self::ThreadCompleted,
ThreadEvent::TurnStarted(_) => Self::TurnStarted,
ThreadEvent::TurnCompleted(_) => Self::TurnCompleted,
ThreadEvent::TurnFailed(_) => Self::TurnFailed,
_ => Self::Other,
}
}
#[inline]
fn from_kind(kind: &str) -> Self {
match kind {
"thread.started" => Self::ThreadStarted,
"thread.completed" => Self::ThreadCompleted,
"turn.started" => Self::TurnStarted,
"turn.completed" => Self::TurnCompleted,
"turn.failed" => Self::TurnFailed,
_ => Self::Other,
}
}
}
struct LogState {
manifest: SessionManifest,
index: TurnIndex,
in_turn: bool,
next_offset: u64,
write_buf: Vec<u8>,
}
impl LogState {
fn new(session_id: &str) -> Self {
Self {
manifest: SessionManifest::new(session_id),
index: TurnIndex::default(),
in_turn: false,
next_offset: 0,
write_buf: Vec::with_capacity(65536),
}
}
fn serialize_event(&mut self, event: &ThreadEvent) -> Result<(u64, u64), SessionStoreError> {
let start = self.next_offset;
let buf_len_before = self.write_buf.len();
if let Err(err) = serde_json::to_writer(
&mut self.write_buf,
&BorrowedVersionedEvent { schema_version: EVENT_SCHEMA_VERSION, event },
) {
self.write_buf.truncate(buf_len_before);
return Err(err.into());
}
self.write_buf.push(b'\n');
let written = self.write_buf.len() - buf_len_before;
let end = start + written as u64;
self.next_offset = end;
Ok((start, end))
}
fn apply_lifecycle_event(&mut self, kind: LifecycleKind, start: u64, end: u64) -> bool {
match kind {
LifecycleKind::ThreadStarted => {
self.manifest.status = "active".to_string();
false
}
LifecycleKind::ThreadCompleted => {
self.manifest.status = "completed".to_string();
true
}
LifecycleKind::TurnStarted => {
self.manifest.status = "active".to_string();
self.in_turn = true;
let n = self.manifest.turn_count + 1;
self.index.entries.push_back(TurnIndexEntry {
turn_number: n,
start_offset: start,
end_offset: end,
event_count: 1,
ts: now_rfc3339(),
});
false
}
LifecycleKind::TurnCompleted | LifecycleKind::TurnFailed => {
if self.in_turn {
if let Some(entry) = self.index.entries.back_mut() {
entry.end_offset = end;
entry.event_count += 1;
}
self.in_turn = false;
self.manifest.turn_count = self.index.entries.len() as u64;
}
true
}
LifecycleKind::Other => {
if self.in_turn
&& let Some(entry) = self.index.entries.back_mut()
{
entry.end_offset = end;
entry.event_count += 1;
}
false
}
}
}
fn plan_cap_eviction(&mut self, max_events: usize) -> Option<(u64, u64)> {
if max_events == 0 || self.manifest.event_count <= max_events as u64 {
return None;
}
let mut evicted_event_count = 0u64;
let mut truncate_offset = 0u64;
while self.manifest.event_count - evicted_event_count > max_events as u64
&& let Some(oldest) = self.index.entries.front()
{
truncate_offset = oldest.end_offset;
evicted_event_count += oldest.event_count;
self.index.entries.pop_front();
}
if truncate_offset == 0 {
None
} else {
Some((truncate_offset, evicted_event_count))
}
}
}
pub struct SessionEventLog {
events_path: PathBuf,
manifest_store: ManifestStore,
file: Mutex<File>,
state: Mutex<LogState>,
max_events: usize,
}
impl SessionEventLog {
pub(crate) fn open(workspace: &Path, session_id: &str, max_events: usize) -> Result<Self, SessionStoreError> {
let dir = session_dir(workspace, session_id);
std::fs::create_dir_all(dir.join(crate::DERIVED_DIR))
.map_err(|e| SessionStoreError::CreateDir { path: dir.clone(), source: e })?;
std::fs::create_dir_all(dir.join("index"))
.map_err(|e| SessionStoreError::CreateDir { path: dir.clone(), source: e })?;
let events_path = dir.join("events.jsonl");
let file = OpenOptions::new()
.create(true)
.read(true)
.append(true)
.open(&events_path)
.map_err(|e| SessionStoreError::io(events_path.clone(), e))?;
let manifest_store = ManifestStore::new(dir.clone());
let log = Self {
events_path: events_path.clone(),
manifest_store,
file: Mutex::new(file),
state: Mutex::new(LogState::new(session_id)),
max_events,
};
let manifest_opt = log.manifest_store.load_manifest()?;
let index_opt = log.manifest_store.load_turn_index()?;
let file_len = std::fs::metadata(&events_path)
.map_err(|e| SessionStoreError::io(events_path.clone(), e))?
.len();
match (manifest_opt, index_opt) {
(Some(manifest), Some(index)) => {
let mut st = log.state.lock().map_err(poison)?;
st.manifest = manifest;
st.index = index;
st.next_offset = file_len;
}
_ => {
log.scan()?;
let mut st = log.state.lock().map_err(poison)?;
st.next_offset = file_len;
}
}
Ok(log)
}
pub fn append(&self, event: &ThreadEvent) -> Result<(), SessionStoreError> {
let mut st = self.state.lock().map_err(poison)?;
let (start, end) = st.serialize_event(event)?;
st.manifest.event_count += 1;
st.manifest.updated_at = now_rfc3339();
let is_turn_boundary = st.apply_lifecycle_event(LifecycleKind::from_event(event), start, end);
if is_turn_boundary {
self.persist_meta_locked(&mut st)?;
}
if st.write_buf.len() >= MAX_WRITE_BUFFER_BYTES {
self.persist_meta_locked(&mut st)?;
}
drop(st);
self.enforce_event_cap()
}
fn enforce_event_cap(&self) -> Result<(), SessionStoreError> {
let mut st = self.state.lock().map_err(poison)?;
let Some((truncate_offset, evicted_event_count)) = st.plan_cap_eviction(self.max_events) else {
return Ok(());
};
self.flush_write_buf_locked(&mut st)?;
{
let mut file = self.file.lock().map_err(poison)?;
file.seek(SeekFrom::Start(truncate_offset))
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
let mut remaining = Vec::new();
file.read_to_end(&mut remaining)
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
file.set_len(0).map_err(|e| SessionStoreError::io(&self.events_path, e))?;
file.seek(SeekFrom::Start(0))
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
file.write_all(&remaining)
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
file.flush().map_err(|e| SessionStoreError::io(&self.events_path, e))?;
}
for entry in &mut st.index.entries {
entry.start_offset -= truncate_offset;
entry.end_offset -= truncate_offset;
}
st.next_offset -= truncate_offset;
st.manifest.event_count = st.manifest.event_count.saturating_sub(evicted_event_count);
self.persist_meta_locked(&mut st)?;
Ok(())
}
pub(crate) fn reconstruct_turn(&self, turn: u64) -> Result<Vec<ThreadEvent>, SessionStoreError> {
let entry = {
let st = self.state.lock().map_err(poison)?;
st.index
.entries
.iter()
.find(|e| e.turn_number == turn)
.cloned()
.ok_or(SessionStoreError::TurnNotFound { session: st.manifest.session_id.clone(), turn })?
};
{
let mut st = self.state.lock().map_err(poison)?;
self.flush_write_buf_locked(&mut st)?;
}
let buf = {
let mut file = self.file.lock().map_err(poison)?;
file.seek(SeekFrom::Start(entry.start_offset))
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
let len = (entry.end_offset - entry.start_offset) as usize;
let mut buf = vec![0u8; len];
file.read_exact(&mut buf)
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
buf
};
let text = String::from_utf8_lossy(&buf);
let mut events = Vec::new();
for line in text.lines() {
let line = line.trim();
if line.is_empty() {
continue;
}
let v: VersionedThreadEvent = match serde_json::from_str(line) {
Ok(v) => v,
Err(_) => continue,
};
events.push(v.into_event());
}
Ok(events)
}
#[must_use]
pub(crate) fn turn_count(&self) -> u64 {
self.state.lock().map_err(poison).map_or(0, |s| s.manifest.turn_count)
}
#[must_use]
pub fn event_count(&self) -> u64 {
self.state.lock().map_err(poison).map_or(0, |s| s.manifest.event_count)
}
pub fn flush(&self) -> Result<(), SessionStoreError> {
let mut st = self.state.lock().map_err(poison)?;
self.persist_meta_locked(&mut st)
}
#[must_use]
pub fn manifest(&self) -> SessionManifest {
self.state
.lock()
.map_err(poison)
.map(|s| s.manifest.clone())
.unwrap_or_else(|_| SessionManifest::new(""))
}
#[must_use]
pub fn turn_index(&self) -> TurnIndex {
self.state.lock().map_err(poison).map(|s| s.index.clone()).unwrap_or_default()
}
pub(crate) fn complete(&self) -> Result<(), SessionStoreError> {
let mut st = self.state.lock().map_err(poison)?;
st.manifest.updated_at = now_rfc3339();
self.persist_meta_locked(&mut st)
}
fn scan(&self) -> Result<(), SessionStoreError> {
let mut st = self.state.lock().map_err(poison)?;
if !self.events_path.exists() {
return Ok(());
}
let file = File::open(&self.events_path).map_err(|e| SessionStoreError::io(&self.events_path, e))?;
let mut reader = std::io::BufReader::new(file);
let mut buf = Vec::new();
let mut pos = 0u64;
let mut first_ts: Option<String> = None;
loop {
buf.clear();
let n = reader
.read_until(b'\n', &mut buf)
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
if n == 0 {
break;
}
let line_end = pos + n as u64;
let trimmed = std::str::from_utf8(&buf).unwrap_or("").trim();
if !trimmed.is_empty()
&& let Ok(v) = serde_json::from_str::<VersionedEventKind<'_>>(trimmed)
{
let kind = v.event.kind;
if requires_full_lifecycle_validation(kind)
&& serde_json::from_str::<VersionedThreadEvent>(trimmed).is_err()
{
pos = line_end;
continue;
}
st.manifest.event_count += 1;
if kind == "thread.started" && first_ts.is_none() {
first_ts = Some(now_rfc3339());
}
st.apply_lifecycle_event(LifecycleKind::from_kind(kind), pos, line_end);
}
pos = line_end;
}
st.in_turn = false;
if let Some(ts) = first_ts
&& st.manifest.created_at.is_empty()
{
st.manifest.created_at = ts;
}
Ok(())
}
fn persist_meta_locked(&self, st: &mut LogState) -> Result<(), SessionStoreError> {
self.flush_write_buf_locked(st)?;
self.manifest_store.write_manifest(&st.manifest)?;
self.manifest_store.write_turn_index(&st.index)?;
Ok(())
}
fn flush_write_buf_locked(&self, st: &mut LogState) -> Result<(), SessionStoreError> {
if st.write_buf.is_empty() {
return Ok(());
}
let mut file = self.file.lock().map_err(poison)?;
file.write_all(&st.write_buf)
.map_err(|e| SessionStoreError::io(&self.events_path, e))?;
st.write_buf.clear();
Ok(())
}
}
fn requires_full_lifecycle_validation(kind: &str) -> bool {
matches!(kind, "thread.started" | "thread.completed" | "turn.started" | "turn.completed" | "turn.failed")
}
impl Drop for SessionEventLog {
fn drop(&mut self) {
if let Ok(mut st) = self.state.lock() {
let _ = self.flush_write_buf_locked(&mut st);
}
}
}
fn poison<T>(_e: std::sync::PoisonError<T>) -> SessionStoreError {
SessionStoreError::Io {
path: PathBuf::new(),
source: std::io::Error::other("session store lock poisoned"),
}
}
fn now_rfc3339() -> String {
Utc::now().to_rfc3339()
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct SessionManifest {
pub session_id: String,
schema_version: u32,
pub created_at: String,
pub updated_at: String,
pub turn_count: u64,
pub event_count: u64,
pub status: String,
}
impl SessionManifest {
#[must_use]
pub(crate) fn new(session_id: &str) -> Self {
let ts = now_rfc3339();
Self {
session_id: session_id.to_string(),
schema_version: crate::SESSION_STORE_SCHEMA_VERSION,
created_at: ts.clone(),
updated_at: ts,
turn_count: 0,
event_count: 0,
status: "active".to_string(),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct TurnIndexEntry {
turn_number: u64,
start_offset: u64,
end_offset: u64,
event_count: u64,
ts: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
pub struct TurnIndex {
entries: VecDeque<TurnIndexEntry>,
}
impl TurnIndex {
#[must_use]
pub fn len(&self) -> usize {
self.entries.len()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.entries.is_empty()
}
}
#[cfg(test)]
mod borrowed_envelope_tests {
use super::{BorrowedVersionedEvent, EVENT_SCHEMA_VERSION};
use vtcode_exec_events::{
ThreadEvent, ThreadStartedEvent, TurnCompletedEvent, TurnStartedEvent, Usage, VersionedThreadEvent,
};
#[test]
fn borrowed_envelope_matches_versioned_envelope() {
for event in [
ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread".to_string() }),
ThreadEvent::TurnStarted(TurnStartedEvent::default()),
ThreadEvent::TurnCompleted(TurnCompletedEvent { usage: Usage::default() }),
] {
let canonical =
serde_json::to_string(&VersionedThreadEvent::new(event.clone())).expect("canonical serialize");
let borrowed = serde_json::to_string(&BorrowedVersionedEvent {
schema_version: EVENT_SCHEMA_VERSION,
event: &event,
})
.expect("borrowed serialize");
assert_eq!(canonical, borrowed, "JSON differs for {event:?}");
}
}
}
#[cfg(test)]
mod lifecycle_state_machine_tests {
use super::{LifecycleKind, LogState};
use vtcode_exec_events::{
ThreadCompletedEvent, ThreadCompletionSubtype, ThreadEvent, ThreadStartedEvent, TurnCompletedEvent,
TurnFailedEvent, TurnStartedEvent, Usage,
};
fn fresh_state() -> LogState {
LogState::new("test-session")
}
#[test]
fn lifecycle_kind_from_event_covers_all_variants() {
assert_eq!(
LifecycleKind::from_event(&ThreadEvent::TurnStarted(TurnStartedEvent::default())),
LifecycleKind::TurnStarted
);
assert_eq!(
LifecycleKind::from_event(&ThreadEvent::TurnCompleted(TurnCompletedEvent { usage: Usage::default() })),
LifecycleKind::TurnCompleted
);
assert_eq!(
LifecycleKind::from_event(&ThreadEvent::TurnFailed(TurnFailedEvent {
message: "err".to_string(),
usage: None,
})),
LifecycleKind::TurnFailed
);
assert_eq!(
LifecycleKind::from_event(&ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "x".to_string() })),
LifecycleKind::ThreadStarted
);
assert_eq!(
LifecycleKind::from_event(&ThreadEvent::ThreadCompleted(ThreadCompletedEvent {
thread_id: "x".to_string(),
session_id: "x".to_string(),
subtype: ThreadCompletionSubtype::Success,
outcome_code: "completed".to_string(),
result: None,
stop_reason: None,
usage: Usage::default(),
total_cost_usd: None,
num_turns: 1,
})),
LifecycleKind::ThreadCompleted
);
assert_eq!(LifecycleKind::from_kind("thread.started"), LifecycleKind::ThreadStarted);
assert_eq!(LifecycleKind::from_kind("thread.completed"), LifecycleKind::ThreadCompleted);
}
#[test]
fn lifecycle_kind_from_str_matches_event_discriminator() {
assert_eq!(LifecycleKind::from_kind("turn.started"), LifecycleKind::TurnStarted);
assert_eq!(LifecycleKind::from_kind("turn.completed"), LifecycleKind::TurnCompleted);
assert_eq!(LifecycleKind::from_kind("turn.failed"), LifecycleKind::TurnFailed);
assert_eq!(LifecycleKind::from_kind("tool.called"), LifecycleKind::Other);
assert_eq!(LifecycleKind::from_kind("thread.started"), LifecycleKind::ThreadStarted);
assert_eq!(LifecycleKind::from_kind("thread.completed"), LifecycleKind::ThreadCompleted);
}
#[test]
fn turn_started_pushes_index_entry_and_sets_in_turn() {
let mut st = fresh_state();
st.manifest.status = "completed".to_string();
let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
assert!(!is_boundary, "TurnStarted is not a turn boundary");
assert!(st.in_turn);
assert_eq!(st.manifest.status, "active");
assert_eq!(st.index.entries.len(), 1);
let entry = &st.index.entries[0];
assert_eq!(entry.turn_number, 1);
assert_eq!(entry.start_offset, 0);
assert_eq!(entry.end_offset, 100);
assert_eq!(entry.event_count, 1);
}
#[test]
fn intermediate_events_extend_current_turn() {
let mut st = fresh_state();
st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
let is_b1 = st.apply_lifecycle_event(LifecycleKind::Other, 100, 200);
let is_b2 = st.apply_lifecycle_event(LifecycleKind::Other, 200, 300);
assert!(!is_b1 && !is_b2);
assert!(st.in_turn);
assert_eq!(st.index.entries.len(), 1);
let entry = &st.index.entries[0];
assert_eq!(entry.end_offset, 300);
assert_eq!(entry.event_count, 3);
}
#[test]
fn turn_completed_closes_turn_and_returns_boundary() {
let mut st = fresh_state();
st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
st.apply_lifecycle_event(LifecycleKind::Other, 100, 200);
let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnCompleted, 200, 300);
assert!(is_boundary);
assert!(!st.in_turn);
assert_eq!(st.manifest.turn_count, 1);
assert_eq!(st.manifest.status, "active");
st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 300, 400);
assert_eq!(st.manifest.status, "completed");
let entry = &st.index.entries[0];
assert_eq!(entry.end_offset, 300);
assert_eq!(entry.event_count, 3);
}
#[test]
fn turn_failed_closes_turn_without_terminal_thread_status() {
let mut st = fresh_state();
st.apply_lifecycle_event(LifecycleKind::TurnStarted, 0, 100);
let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnFailed, 100, 200);
assert!(is_boundary);
assert!(!st.in_turn);
assert_eq!(st.manifest.turn_count, 1);
assert_eq!(st.manifest.status, "active");
st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 200, 300);
assert_eq!(st.manifest.status, "completed");
}
#[test]
fn turn_completed_without_turn_started_is_idempotent() {
let mut st = fresh_state();
let is_boundary = st.apply_lifecycle_event(LifecycleKind::TurnCompleted, 0, 100);
assert!(is_boundary);
assert!(!st.in_turn);
assert_eq!(st.manifest.turn_count, 0, "no turn was started");
assert_eq!(st.manifest.status, "active");
st.apply_lifecycle_event(LifecycleKind::ThreadCompleted, 100, 200);
assert_eq!(st.manifest.status, "completed");
assert!(st.index.entries.is_empty());
}
#[test]
fn multiple_turns_get_incrementing_ordinals() {
let mut st = fresh_state();
for n in 1..=3 {
st.apply_lifecycle_event(LifecycleKind::TurnStarted, n * 100, n * 100 + 50);
st.apply_lifecycle_event(LifecycleKind::TurnCompleted, n * 100 + 50, n * 100 + 100);
}
assert_eq!(st.index.entries.len(), 3);
for (i, entry) in st.index.entries.iter().enumerate() {
assert_eq!(entry.turn_number, (i + 1) as u64);
}
assert_eq!(st.manifest.turn_count, 3);
}
}
#[cfg(test)]
mod cap_eviction_tests {
use super::{LogState, TurnIndexEntry};
fn state_with_turns(turns: usize, events_per_turn: u64) -> LogState {
let mut st = LogState::new("cap-test");
st.manifest.event_count = (turns as u64) * events_per_turn;
let mut offset = 0u64;
for n in 1..=turns {
st.index.entries.push_back(TurnIndexEntry {
turn_number: n as u64,
start_offset: offset,
end_offset: offset + events_per_turn * 10,
event_count: events_per_turn,
ts: "2026-01-01T00:00:00Z".to_string(),
});
offset += events_per_turn * 10;
}
st
}
#[test]
fn no_eviction_when_under_cap() {
let mut st = state_with_turns(3, 2); assert!(st.plan_cap_eviction(10).is_none());
assert_eq!(st.index.entries.len(), 3, "no turns should be evicted");
}
#[test]
fn no_eviction_when_cap_disabled() {
let mut st = state_with_turns(5, 2); assert!(st.plan_cap_eviction(0).is_none());
assert_eq!(st.index.entries.len(), 5);
}
#[test]
fn evicts_oldest_turns_to_meet_cap() {
let mut st = state_with_turns(5, 2);
let (truncate_offset, evicted) = st.plan_cap_eviction(6).expect("eviction planned");
assert_eq!(evicted, 4, "should evict 4 events (2 turns)");
assert_eq!(st.index.entries.len(), 3, "should keep 3 turns");
assert_eq!(truncate_offset, 40); assert_eq!(st.index.entries[0].turn_number, 3);
assert_eq!(st.index.entries[2].turn_number, 5);
}
#[test]
fn evicts_all_turns_when_cap_smaller_than_one_turn() {
let mut st = state_with_turns(3, 5);
let (_truncate_offset, evicted) = st.plan_cap_eviction(3).expect("eviction planned");
assert_eq!(evicted, 15, "all events evicted");
assert_eq!(st.index.entries.len(), 0);
}
}