use super::event::{SessionEvent, SessionEventKind};
use super::manager::{Session, SessionInternalDiagnostic, SessionListReport, SessionManager};
use super::read::{MAX_METADATA_VISIT_BYTES, MAX_METADATA_VISIT_LINES, validate_session_id};
use super::store::{open_existing_named, validate_existing_file, validate_path_file};
use super::write::{
COMPACTION_SCHEMA_VERSION, sanitize_compaction_summary, sanitize_session_title,
};
use crate::{output::redact_sensitive_text, persistence::atomic_write_with_permissions};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::{
fs,
io::Read,
path::{Path, PathBuf},
time::{SystemTime, UNIX_EPOCH},
};
const MAX_SESSION_METADATA_BYTES: u64 = 64 * 1024;
const MAX_SESSION_INTERNAL_DIAGNOSTICS: usize = 64;
pub(crate) fn push_session_diagnostic(
diagnostics: &mut Vec<SessionInternalDiagnostic>,
session_id: Option<String>,
message: String,
) {
if diagnostics.len() < MAX_SESSION_INTERNAL_DIAGNOSTICS {
let redacted = redact_sensitive_text(&message);
let mut chars = redacted.chars();
let mut bounded = chars.by_ref().take(500).collect::<String>();
if chars.next().is_some() {
bounded.push('…');
}
diagnostics.push(SessionInternalDiagnostic {
session_id,
message: bounded,
});
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) enum SessionDiagnosticOperation {
Replay,
Metadata,
Listing,
Title,
RecentContext,
}
impl SessionDiagnosticOperation {
fn as_str(self) -> &'static str {
match self {
Self::Replay => "replay",
Self::Metadata => "metadata",
Self::Listing => "listing",
Self::Title => "title",
Self::RecentContext => "recent_context",
}
}
}
pub(crate) fn report_session_diagnostic(
operation: SessionDiagnosticOperation,
_path: &Path,
_error: impl std::fmt::Display,
) {
eprintln!(
"warning: session operation={} category=session_jsonl failed",
operation.as_str()
);
}
const SESSION_METADATA_SCHEMA_VERSION: u64 = 3;
const SESSION_METADATA_EXTENSION: &str = "metadata.json";
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub(in crate::sessions) struct SessionMetadataRecord {
schema_version: u64,
session_id: String,
latest_activity_timestamp: Option<DateTime<Utc>>,
latest_title: Option<String>,
session_path: Option<PathBuf>,
#[serde(default)]
first_user_input_text: Option<String>,
user_input_count: usize,
jsonl_len: u64,
jsonl_modified_ns: Option<u128>,
complete: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SessionMetadataSummary {
pub(crate) session: Session,
pub(crate) latest_activity_time: SystemTime,
pub(crate) latest_title: Option<String>,
pub(crate) session_path: Option<PathBuf>,
}
impl SessionMetadataRecord {
fn empty(session: &Session, marker: JsonlMarker, complete: bool) -> Self {
Self {
schema_version: SESSION_METADATA_SCHEMA_VERSION,
session_id: session.id.clone(),
latest_activity_timestamp: None,
latest_title: None,
session_path: None,
first_user_input_text: None,
user_input_count: 0,
jsonl_len: marker.len,
jsonl_modified_ns: marker.modified_ns,
complete,
}
}
fn apply_event(&mut self, event: &SessionEvent) {
if self.session_path.is_none() {
self.session_path = event.session_path.clone();
}
self.latest_activity_timestamp = Some(
self.latest_activity_timestamp
.map_or(event.timestamp, |current| current.max(event.timestamp)),
);
match event.kind() {
Some(SessionEventKind::SessionTitle) => {
if let Some(title) = event.payload.get("title").and_then(Value::as_str)
&& let Some(title) = sanitize_session_title(title)
{
self.latest_title = Some(title);
}
self.first_user_input_text = None;
}
Some(SessionEventKind::UserInput) => {
self.user_input_count = self.user_input_count.saturating_add(1);
if self.first_user_input_text.is_none()
&& let Some(text) = event.payload.get("text").and_then(Value::as_str)
&& !text.trim().is_empty()
{
self.first_user_input_text = Some(text.to_string());
}
}
Some(SessionEventKind::Compaction) => {
if let Some(aggregate) = event.payload.get("aggregate") {
if let Some(title) = aggregate.get("latest_title").and_then(Value::as_str) {
self.latest_title = sanitize_session_title(title);
}
if let Some(count) = aggregate.get("user_input_count").and_then(Value::as_u64)
&& let Ok(count) = usize::try_from(count)
{
self.user_input_count = count;
}
self.first_user_input_text = aggregate
.get("first_user_input_text")
.and_then(Value::as_str)
.map(|text| redact_sensitive_text(text).chars().take(4_000).collect());
if let Some(path) = aggregate.get("session_path").and_then(Value::as_str) {
self.session_path = Some(PathBuf::from(path));
}
if let Some(timestamp) = aggregate.get("latest_activity_timestamp")
&& let Ok(timestamp) =
serde_json::from_value::<chrono::DateTime<Utc>>(timestamp.clone())
{
self.latest_activity_timestamp = Some(
self.latest_activity_timestamp
.map_or(timestamp, |current| current.max(timestamp)),
);
}
}
}
_ => {}
}
}
pub(crate) fn checkpoint_aggregate(&self) -> Value {
serde_json::json!({
"latest_title": self.latest_title,
"user_input_count": self.user_input_count,
"first_user_input_text": self.first_user_input_text.as_deref().map(|text| redact_sensitive_text(text).chars().take(4_000).collect::<String>()),
"session_path": self.session_path.as_ref().map(|path| path.to_string_lossy().to_string()),
"latest_activity_timestamp": self.latest_activity_timestamp,
})
}
pub(crate) fn activity_time(&self, session: &Session) -> SystemTime {
self.latest_activity_timestamp
.map(SystemTime::from)
.unwrap_or_else(|| jsonl_modified_time(session).unwrap_or(SystemTime::UNIX_EPOCH))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct JsonlMarker {
len: u64,
modified_ns: Option<u128>,
}
impl SessionManager {
pub(crate) fn list_metadata_summaries(&self) -> anyhow::Result<Vec<SessionMetadataSummary>> {
Ok(self.list_metadata_report()?.summaries)
}
pub(crate) fn list_metadata_report(&self) -> anyhow::Result<SessionListReport> {
if !self.root.exists() {
return Ok(SessionListReport {
summaries: Vec::new(),
diagnostics: Vec::new(),
});
}
super::store::validate_session_root(&self.root)?;
let mut summaries = Vec::new();
let mut diagnostics = Vec::new();
for entry in fs::read_dir(&self.root)? {
let entry = match entry {
Ok(entry) => entry,
Err(error) => {
push_session_diagnostic(
&mut diagnostics,
None,
format!("failed to read session directory entry: {error}"),
);
continue;
}
};
let path = entry.path();
let Some((id, path)) = session_from_entry(path) else {
continue;
};
let id = match validate_session_id(id) {
Ok(id) => id,
Err(error) => {
push_session_diagnostic(
&mut diagnostics,
None,
format!("ignored invalid session file {}: {error}", path.display()),
);
continue;
}
};
if validate_path_file(&self.root, &id, &path).is_err() {
push_session_diagnostic(
&mut diagnostics,
Some(id),
"ignored unsafe session file".to_string(),
);
continue;
}
let session = Session { id, path };
match session_metadata_summary(session) {
Ok(summary) => summaries.push(summary),
Err(error) => push_session_diagnostic(
&mut diagnostics,
None,
format!("failed to load session metadata: {error}"),
),
}
}
summaries.sort_by(|left, right| {
left.latest_activity_time
.cmp(&right.latest_activity_time)
.then_with(|| left.session.id.cmp(&right.session.id))
});
Ok(SessionListReport {
summaries,
diagnostics,
})
}
}
pub(crate) fn session_from_entry(path: PathBuf) -> Option<(String, PathBuf)> {
if path
.extension()
.is_none_or(|extension| extension != "jsonl")
{
return None;
}
let id = path.file_stem()?.to_string_lossy().to_string();
Some((id, path))
}
pub(crate) fn metadata_path_for_session(session: &Session) -> PathBuf {
session.path.with_extension(SESSION_METADATA_EXTENSION)
}
fn jsonl_marker(session: &Session) -> anyhow::Result<JsonlMarker> {
let root = session
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
let Some(file) = super::store::open_existing_primary(root, &session.id)? else {
anyhow::bail!("session JSONL is missing");
};
let metadata = file.metadata()?;
Ok(JsonlMarker {
len: metadata.len(),
modified_ns: metadata.modified().ok().map(system_time_ns),
})
}
fn jsonl_modified_time(session: &Session) -> Option<SystemTime> {
let root = session.path.parent()?;
super::store::open_existing_primary(root, &session.id)
.ok()
.flatten()
.and_then(|file| file.metadata().ok())
.and_then(|metadata| metadata.modified().ok())
}
fn system_time_ns(time: SystemTime) -> u128 {
time.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_nanos()
}
fn read_session_metadata_with_policy(
session: &Session,
require_complete: bool,
) -> anyhow::Result<Option<SessionMetadataRecord>> {
let path = metadata_path_for_session(session);
let root = session
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
let name = path
.file_name()
.ok_or_else(|| anyhow::anyhow!("metadata path has no filename"))?
.to_string_lossy();
if !session.path.exists() {
return Ok(None);
}
let Some(file) = open_existing_named(root, &name)? else {
return Ok(None);
};
let mut bytes = Vec::new();
file.take(MAX_SESSION_METADATA_BYTES.saturating_add(1))
.read_to_end(&mut bytes)?;
if bytes.len() as u64 > MAX_SESSION_METADATA_BYTES {
return Ok(None);
}
let Ok(record) = serde_json::from_slice::<SessionMetadataRecord>(&bytes) else {
return Ok(None);
};
if record.schema_version != SESSION_METADATA_SCHEMA_VERSION
|| record.session_id != session.id
|| (require_complete && !record.complete)
{
return Ok(None);
}
let marker = jsonl_marker(session)?;
if record.jsonl_len != marker.len || record.jsonl_modified_ns != marker.modified_ns {
return Ok(None);
}
Ok(Some(record))
}
pub(crate) fn read_complete_session_metadata(
session: &Session,
) -> anyhow::Result<Option<SessionMetadataRecord>> {
read_session_metadata(session)
}
fn read_session_metadata(session: &Session) -> anyhow::Result<Option<SessionMetadataRecord>> {
read_session_metadata_with_policy(session, true)
}
pub(crate) fn read_session_metadata_for_append(
session: &Session,
) -> anyhow::Result<Option<SessionMetadataRecord>> {
read_session_metadata_with_policy(session, false)
}
fn write_session_metadata(session: &Session, record: &SessionMetadataRecord) -> anyhow::Result<()> {
let path = metadata_path_for_session(session);
validate_existing_file(&path)?;
let bytes = serde_json::to_vec_pretty(record)?;
atomic_write_with_permissions(&path, &bytes, Some(super::store::SESSION_FILE_MODE))
}
pub(crate) fn rebuild_session_metadata_from_jsonl(
session: &Session,
) -> anyhow::Result<SessionMetadataRecord> {
validate_session_id(session.id.clone())?;
let initial_marker = jsonl_marker(session)?;
let mut record = SessionMetadataRecord::empty(session, initial_marker, true);
let (_, _) = session.visit_events_tolerant_bounded(
MAX_METADATA_VISIT_LINES,
MAX_METADATA_VISIT_BYTES,
|event| {
if event.session_id == session.id {
record.apply_event(&event);
}
},
)?;
let final_marker = jsonl_marker(session)?;
record.jsonl_len = final_marker.len;
record.jsonl_modified_ns = final_marker.modified_ns;
record.complete = initial_marker == final_marker;
if let Err(error) = write_session_metadata(session, &record) {
report_session_diagnostic(
SessionDiagnosticOperation::Metadata,
&metadata_path_for_session(session),
&error,
);
}
Ok(record)
}
pub(crate) fn write_rotated_session_metadata(
session: &Session,
mut record: SessionMetadataRecord,
checkpoint: &SessionEvent,
) -> anyhow::Result<()> {
record.apply_event(checkpoint);
let final_marker = jsonl_marker(session)?;
record.jsonl_len = final_marker.len;
record.jsonl_modified_ns = final_marker.modified_ns;
record.complete = true;
write_session_metadata(session, &record)
}
pub(crate) fn session_metadata_for_listing(
session: &Session,
) -> anyhow::Result<SessionMetadataRecord> {
if let Some(record) = read_session_metadata(session)? {
return Ok(record);
}
rebuild_session_metadata_from_jsonl(session)
}
fn session_metadata_summary(session: Session) -> anyhow::Result<SessionMetadataSummary> {
let record = session_metadata_for_listing(&session)?;
let latest_activity_time = record.activity_time(&session);
Ok(SessionMetadataSummary {
session,
latest_activity_time,
latest_title: record.latest_title,
session_path: record.session_path,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct CompactionCheckpoint {
pub(crate) summary: String,
pub(crate) provider: String,
pub(crate) model: String,
pub(crate) cutoff_event_count: usize,
pub(crate) event_index: usize,
}
pub(crate) fn update_session_metadata_after_append_batch(
session: &Session,
previous: Option<SessionMetadataRecord>,
events: &[SessionEvent],
) -> anyhow::Result<()> {
let mut record = match previous {
Some(record) => record,
None => SessionMetadataRecord::empty(session, jsonl_marker(session)?, false),
};
for event in events {
record.apply_event(event);
}
let marker = jsonl_marker(session)?;
record.jsonl_len = marker.len;
record.jsonl_modified_ns = marker.modified_ns;
write_session_metadata(session, &record)
}
pub(crate) fn latest_valid_compaction_checkpoint(
session_id: &str,
events: &[SessionEvent],
) -> (Option<CompactionCheckpoint>, Vec<String>) {
let mut diagnostics = Vec::new();
for (event_index, event) in events.iter().enumerate().rev() {
if event.kind() != Some(SessionEventKind::Compaction) {
continue;
}
match validate_compaction_checkpoint(session_id, event, event_index, events.len()) {
Ok(checkpoint) => return (Some(checkpoint), diagnostics),
Err(message) => diagnostics.push(message),
}
}
(None, diagnostics)
}
fn validate_compaction_checkpoint(
session_id: &str,
event: &SessionEvent,
event_index: usize,
total_event_count: usize,
) -> Result<CompactionCheckpoint, String> {
if event.session_id != session_id {
return Err(format!(
"ignored malformed compaction checkpoint at event_index={event_index}: session_id_mismatch"
));
}
let schema_version = event
.payload
.get("schema_version")
.and_then(Value::as_u64)
.ok_or_else(|| {
format!(
"ignored malformed compaction checkpoint at event_index={event_index}: missing_schema_version"
)
})?;
if schema_version != COMPACTION_SCHEMA_VERSION {
return Err(format!(
"ignored malformed compaction checkpoint at event_index={event_index}: unsupported_schema_version"
));
}
let summary = event
.payload
.get("summary")
.and_then(Value::as_str)
.and_then(sanitize_compaction_summary)
.ok_or_else(|| {
format!(
"ignored malformed compaction checkpoint at event_index={event_index}: empty_summary"
)
})?;
let provider = required_non_blank_payload_string(event, "provider", event_index)?;
let model = required_non_blank_payload_string(event, "model", event_index)?;
let cutoff_event_count = event
.payload
.get("cutoff_event_count")
.and_then(Value::as_u64)
.and_then(|value| usize::try_from(value).ok())
.ok_or_else(|| {
format!(
"ignored malformed compaction checkpoint at event_index={event_index}: missing_cutoff_event_count"
)
})?;
if cutoff_event_count > event_index || cutoff_event_count > total_event_count {
return Err(format!(
"ignored malformed compaction checkpoint at event_index={event_index}: cutoff_event_count_out_of_bounds"
));
}
Ok(CompactionCheckpoint {
summary,
provider,
model,
cutoff_event_count,
event_index,
})
}
fn required_non_blank_payload_string(
event: &SessionEvent,
key: &str,
event_index: usize,
) -> Result<String, String> {
event.payload.get(key).and_then(Value::as_str).map(str::trim)
.filter(|value| !value.is_empty())
.map(ToString::to_string)
.ok_or_else(|| format!("ignored malformed compaction checkpoint at event_index={event_index}: missing_{key}"))
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub(crate) struct SessionTitleMetadata {
pub(crate) latest_title: Option<String>,
pub(crate) user_input_count: usize,
pub(crate) first_user_input_text: Option<String>,
}
pub(crate) fn session_title_metadata(session: &Session) -> anyhow::Result<SessionTitleMetadata> {
if let Some(record) = read_session_metadata(session)? {
return Ok(SessionTitleMetadata {
latest_title: record.latest_title,
user_input_count: record.user_input_count,
first_user_input_text: record.first_user_input_text,
});
}
let mut metadata = SessionTitleMetadata::default();
session
.visit_events_tolerant_bounded(
MAX_METADATA_VISIT_LINES,
MAX_METADATA_VISIT_BYTES,
|event| match event.kind() {
Some(SessionEventKind::SessionTitle) => {
if let Some(title) = event
.payload
.get("title")
.and_then(Value::as_str)
.and_then(sanitize_session_title)
{
metadata.latest_title = Some(title);
}
}
Some(SessionEventKind::UserInput) => {
metadata.user_input_count = metadata.user_input_count.saturating_add(1);
if metadata.first_user_input_text.is_none()
&& let Some(text) = event.payload.get("text").and_then(Value::as_str)
&& !text.trim().is_empty()
{
metadata.first_user_input_text = Some(text.to_string());
}
}
_ => {}
},
)
.inspect_err(|error| {
report_session_diagnostic(SessionDiagnosticOperation::Title, session.path(), error);
})?;
Ok(metadata)
}
pub fn latest_session_title(session: &Session) -> anyhow::Result<Option<String>> {
Ok(session_title_metadata(session)?.latest_title)
}
pub fn session_user_input_count(session: &Session) -> anyhow::Result<usize> {
Ok(session_title_metadata(session)?.user_input_count)
}
#[cfg(test)]
mod tests {
use super::super::read::latest_valid_event_timestamp_streaming;
use super::super::write::{
SESSION_TITLE_EVENT, record_session_compaction, record_session_event, record_session_title,
};
use super::*;
use chrono::{TimeZone, Timelike};
use serde_json::json;
use tempfile::TempDir;
#[test]
fn compaction_checkpoint_appends_round_trip_and_redacts_summary() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
record_session_compaction(
&session,
temp.path(),
"summary with sk-testSecret123456",
"provider-a",
"model-a",
0,
)
.unwrap();
let events = session.read_events().unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, "compaction");
assert_eq!(
events[0].payload["schema_version"],
COMPACTION_SCHEMA_VERSION
);
assert_eq!(events[0].payload["provider"], "provider-a");
assert_eq!(events[0].payload["model"], "model-a");
assert_eq!(events[0].payload["cutoff_event_count"], 0);
assert!(
!events[0].payload["summary"]
.to_string()
.contains("sk-testSecret")
);
let (checkpoint, diagnostics) = latest_valid_compaction_checkpoint(session.id(), &events);
assert!(diagnostics.is_empty(), "{diagnostics:?}");
let checkpoint = checkpoint.unwrap();
assert_eq!(checkpoint.provider, "provider-a");
assert_eq!(checkpoint.model, "model-a");
assert_eq!(checkpoint.cutoff_event_count, 0);
assert_eq!(checkpoint.event_index, 0);
}
#[test]
fn compaction_checkpoint_rejects_empty_summary() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let error = record_session_compaction(&session, temp.path(), " \n\t ", "p", "m", 0)
.unwrap_err()
.to_string();
assert!(error.contains("empty"), "{error}");
assert!(session.read_events().unwrap().is_empty());
}
#[test]
fn latest_valid_compaction_falls_back_after_malformed_or_mismatched_latest() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let valid = SessionEvent::new_kind(
SessionEventKind::Compaction,
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"schema_version": COMPACTION_SCHEMA_VERSION,
"summary":"valid summary",
"provider":"p",
"model":"m",
"cutoff_event_count": 0,
}),
);
let malformed = SessionEvent::new_kind(
SessionEventKind::Compaction,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"schema_version": COMPACTION_SCHEMA_VERSION, "summary":" "}),
);
let mismatched = SessionEvent::new_kind(
SessionEventKind::Compaction,
"other".to_string(),
temp.path().to_path_buf(),
json!({
"schema_version": COMPACTION_SCHEMA_VERSION,
"summary":"wrong session",
"provider":"p",
"model":"m",
"cutoff_event_count": 2,
}),
);
let (checkpoint, diagnostics) = latest_valid_compaction_checkpoint(
session.id(),
&[valid.clone(), malformed, mismatched],
);
assert_eq!(checkpoint.unwrap().summary, "valid summary");
assert_eq!(diagnostics.len(), 2);
assert!(
diagnostics
.iter()
.any(|message| message.contains("session_id_mismatch"))
);
assert!(
diagnostics
.iter()
.any(|message| message.contains("empty_summary"))
);
assert!(!diagnostics.join("\n").contains("wrong session"));
}
#[test]
fn session_titles_sanitize_round_trip_and_latest_wins() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
assert_eq!(session_title_metadata(&session).unwrap().latest_title, None);
assert_eq!(
sanitize_session_title(" \"Fix parser\"\nnow "),
Some("Fix parser now".to_string())
);
assert_eq!(sanitize_session_title("\n\t\r"), None);
assert_eq!(
sanitize_session_title(&"é".repeat(60))
.unwrap()
.chars()
.count(),
50
);
record_session_title(
&session,
temp.path(),
"First title",
"provider-a",
"model-a",
)
.unwrap();
record_session_title(
&session,
temp.path(),
"'Second title'",
"provider-b",
"model-b",
)
.unwrap();
let events = session.read_events().unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].event_type, SESSION_TITLE_EVENT);
assert_eq!(events[0].payload["title"], "First title");
assert_eq!(events[0].payload["provider"], "provider-a");
assert_eq!(events[0].payload["model"], "model-a");
assert!(events[0].payload.get("access_token").is_none());
assert_eq!(
session_title_metadata(&session)
.unwrap()
.latest_title
.as_deref(),
Some("Second title")
);
}
#[test]
fn session_user_input_count_supports_first_message_trigger() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
assert_eq!(session_user_input_count(&session).unwrap(), 0);
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":"one"}),
)
.unwrap();
assert_eq!(session_user_input_count(&session).unwrap(), 1);
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":"two"}),
)
.unwrap();
assert_eq!(session_user_input_count(&session).unwrap(), 2);
}
#[test]
fn session_title_metadata_uses_first_persisted_non_empty_text() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
record_session_event(
Some(&session),
temp.path(),
"diagnostic",
json!({"text":"ignore"}),
)
.unwrap();
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":" "}),
)
.unwrap();
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":"first durable prompt"}),
)
.unwrap();
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":"second prompt"}),
)
.unwrap();
assert_eq!(
session_title_metadata(&session)
.unwrap()
.first_user_input_text
.as_deref(),
Some("first durable prompt")
);
}
#[test]
fn session_title_metadata_rejects_oversized_jsonl() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("oversized-title").unwrap();
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(session.path())
.unwrap()
.set_len((MAX_METADATA_VISIT_BYTES as u64) + 1)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
let error = session_title_metadata(&session).unwrap_err().to_string();
assert!(error.contains("tolerant read limit exceeded"), "{error}");
assert!(error.contains("bytes"), "{error}");
}
#[test]
fn metadata_rebuild_rejects_oversized_jsonl() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("oversized-metadata").unwrap();
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::OpenOptions::new()
.create(true)
.write(true)
.truncate(true)
.open(session.path())
.unwrap()
.set_len((MAX_METADATA_VISIT_BYTES as u64) + 1)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
let error = session_metadata_for_listing(&session)
.unwrap_err()
.to_string();
assert!(error.contains("tolerant read limit exceeded"), "{error}");
assert!(error.contains("bytes"), "{error}");
}
#[test]
fn session_title_and_user_count_tolerate_malformed_lines() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let first_title = SessionEvent::new_kind(
SessionEventKind::SessionTitle,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"title":"First title", "provider":"p", "model":"m"}),
);
let user_input = SessionEvent::new_kind(
SessionEventKind::UserInput,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"one"}),
);
let latest_title = SessionEvent::new_kind(
SessionEventKind::SessionTitle,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"title":"Latest title", "provider":"p", "model":"m"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\nnot json\n{}\n{}\n",
serde_json::to_string(&first_title).unwrap(),
serde_json::to_string(&user_input).unwrap(),
serde_json::to_string(&latest_title).unwrap()
),
)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
assert!(session.read_events().is_err());
let metadata = session_title_metadata(&session).unwrap();
assert_eq!(metadata.latest_title.as_deref(), Some("Latest title"));
assert_eq!(metadata.user_input_count, 1);
assert_eq!(metadata.first_user_input_text.as_deref(), Some("one"));
assert_eq!(
session_title_metadata(&session)
.unwrap()
.latest_title
.as_deref(),
Some("Latest title")
);
assert_eq!(session_user_input_count(&session).unwrap(), 1);
}
#[test]
fn continue_uses_most_recent_session_by_latest_event_timestamp() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
fs::create_dir_all(temp.path().join("sessions")).unwrap();
let z_old = manager.open("zzzz").unwrap();
let a_new = manager.open("aaaa").unwrap();
let z_old_event = event_at(&z_old, temp.path(), 2024, 1, 3);
fs::write(
z_old.path(),
format!("{}\n", serde_json::to_string(&z_old_event).unwrap()),
)
.unwrap();
let a_latest_outside_tail = event_at(&a_new, temp.path(), 2024, 1, 4);
let a_tail_older = event_at(&a_new, temp.path(), 2024, 1, 2);
fs::write(
a_new.path(),
format!(
"{}\n{}\n{}\n",
serde_json::to_string(&a_latest_outside_tail).unwrap(),
"not-json".repeat(300 * 1024),
serde_json::to_string(&a_tail_older).unwrap()
),
)
.unwrap();
crate::sessions::store::secure_test_session_root(temp.path().join("sessions").as_path());
let sessions = manager.list().unwrap();
assert_eq!(sessions.last().unwrap().id(), "aaaa");
assert_eq!(manager.most_recent().unwrap().unwrap().id(), "aaaa");
}
#[test]
fn session_list_activity_streams_latest_valid_timestamp_without_event_vec() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let older = manager.open("older").unwrap();
let newer = manager.open("newer").unwrap();
let older_event = event_at(&older, temp.path(), 2024, 1, 1);
let newer_event = event_at(&newer, temp.path(), 2024, 1, 2);
fs::create_dir_all(newer.path().parent().unwrap()).unwrap();
fs::write(
older.path(),
format!("{}\n", serde_json::to_string(&older_event).unwrap()),
)
.unwrap();
fs::write(
newer.path(),
format!(
"not json\n{}\n",
serde_json::to_string(&newer_event).unwrap()
),
)
.unwrap();
crate::sessions::store::secure_test_session_root(newer.path().parent().unwrap());
let sessions = manager.list().unwrap();
assert_eq!(sessions.last().unwrap().id(), "newer");
assert_eq!(
latest_valid_event_timestamp_streaming(&newer),
Some(SystemTime::from(newer_event.timestamp))
);
}
#[test]
fn session_activity_streaming_uses_latest_valid_event_across_full_file() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("tail-session").unwrap();
let old = event_at(&session, temp.path(), 2024, 1, 1);
let recent = event_at(&session, temp.path(), 2024, 1, 2);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!(
"{}\n{}\nnot json\n",
serde_json::to_string(&old).unwrap(),
serde_json::to_string(&recent).unwrap()
),
)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
assert_eq!(
latest_valid_event_timestamp_streaming(&session),
Some(SystemTime::from(recent.timestamp))
);
}
#[test]
fn record_session_event_propagates_append_failure() {
let temp = TempDir::new().unwrap();
let invalid_session =
Session::unchecked_for_test("../unsafe".to_string(), temp.path().join("unsafe.jsonl"));
let error = record_session_event(
Some(&invalid_session),
temp.path(),
"diagnostic",
json!({"message":"persist me"}),
)
.unwrap_err();
assert!(error.to_string().contains("must not contain '..'"));
assert!(!invalid_session.path().exists());
}
#[test]
fn session_metadata_sidecar_excludes_raw_transcript_and_secrets() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"raw prompt secret sk-testSecret123456", "api_key":"abc123"}),
))
.unwrap();
session
.append(&SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"result":{"content":"tool output bearer secret-token"}}),
))
.unwrap();
record_session_title(
&session,
temp.path(),
"Safe Indexed Title",
"provider-a",
"model-a",
)
.unwrap();
let sidecar = fs::read_to_string(metadata_path_for_session(&session)).unwrap();
assert!(sidecar.contains("Safe Indexed Title"));
assert!(sidecar.contains("session_path"));
for forbidden in [
"raw prompt secret",
"sk-testSecret123456",
"api_key",
"abc123",
"tool output",
"bearer secret-token",
"provider-a",
"model-a",
"access_token",
"account_id",
] {
assert!(!sidecar.contains(forbidden), "sidecar leaked {forbidden}");
}
}
#[test]
fn metadata_indexes_first_session_path_and_keeps_legacy_missing_path_none() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let scoped = temp.path().join("scoped");
fs::create_dir(&scoped).unwrap();
let session = manager.open("scoped-session").unwrap();
session
.append(&SessionEvent::new(
"diagnostic",
session.id().to_string(),
scoped.clone(),
json!({}),
))
.unwrap();
let record: SessionMetadataRecord =
serde_json::from_slice(&fs::read(metadata_path_for_session(&session)).unwrap())
.unwrap();
assert_eq!(record.schema_version, SESSION_METADATA_SCHEMA_VERSION);
assert_eq!(record.session_path.as_deref(), Some(scoped.as_path()));
let summary = manager
.list_metadata_summaries()
.unwrap()
.into_iter()
.find(|summary| summary.session.id() == session.id())
.unwrap();
assert_eq!(summary.session_path.as_deref(), Some(scoped.as_path()));
let legacy = manager.open("legacy-session").unwrap();
let mut legacy_event = SessionEvent::new(
"diagnostic",
legacy.id().to_string(),
temp.path().to_path_buf(),
json!({}),
);
legacy_event.session_path = None;
fs::create_dir_all(legacy.path().parent().unwrap()).unwrap();
fs::write(
legacy.path(),
format!("{}\n", serde_json::to_string(&legacy_event).unwrap()),
)
.unwrap();
crate::sessions::store::secure_test_session_root(legacy.path().parent().unwrap());
let legacy_record = session_metadata_for_listing(&legacy).unwrap();
assert_eq!(legacy_record.session_path, None);
}
#[test]
fn most_recent_uses_valid_sidecar_without_parsing_jsonl() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("indexed-session").unwrap();
let event = event_at(&session, temp.path(), 2024, 1, 5);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!("{}\n", serde_json::to_string(&event).unwrap()),
)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
let mut record = rebuild_session_metadata_from_jsonl(&session).unwrap();
fs::write(
session.path(),
"not valid jsonl but marker matches sidecar\n",
)
.unwrap();
let marker = jsonl_marker(&session).unwrap();
record.jsonl_len = marker.len;
record.jsonl_modified_ns = marker.modified_ns;
write_session_metadata(&session, &record).unwrap();
let latest = manager.most_recent().unwrap().unwrap();
assert_eq!(latest.id(), session.id());
}
#[test]
fn corrupt_sidecar_falls_back_to_jsonl_and_rebuilds() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("fallback-session").unwrap();
let event = event_at(&session, temp.path(), 2024, 2, 1);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!("{}\n", serde_json::to_string(&event).unwrap()),
)
.unwrap();
fs::write(metadata_path_for_session(&session), "{not json").unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
let sessions = manager.list().unwrap();
let rebuilt: SessionMetadataRecord =
serde_json::from_slice(&fs::read(metadata_path_for_session(&session)).unwrap())
.unwrap();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].id(), session.id());
assert_eq!(rebuilt.session_id, session.id());
assert_eq!(rebuilt.latest_activity_timestamp, Some(event.timestamp));
}
#[test]
fn list_metadata_report_surfaces_invalid_session_entries_without_failing() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
fs::create_dir_all(temp.path().join("sessions")).unwrap();
crate::sessions::store::secure_test_session_root(temp.path().join("sessions").as_path());
let valid = manager.open("valid-session").unwrap();
append_event_at(&valid, temp.path(), 2024, 2, 1);
fs::write(temp.path().join("sessions/bad.name.jsonl"), "{}").unwrap();
let report = manager.list_metadata_report().unwrap();
assert_eq!(report.summaries.len(), 1);
assert_eq!(report.summaries[0].session.id(), "valid-session");
assert_eq!(report.diagnostics.len(), 1);
assert!(report.diagnostics[0].message.contains("invalid session"));
}
#[test]
fn append_without_metadata_writes_incomplete_sidecar_without_jsonl_rebuild() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("missing-sidecar").unwrap();
let old_title = SessionEvent::new_kind(
SessionEventKind::SessionTitle,
session.id().to_string(),
temp.path().to_path_buf(),
json!({"title":"Old Title", "provider":"p", "model":"m"}),
);
fs::create_dir_all(session.path().parent().unwrap()).unwrap();
fs::write(
session.path(),
format!("{}\n", serde_json::to_string(&old_title).unwrap()),
)
.unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
session
.append(&SessionEvent::new(
"diagnostic",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"message":"new append"}),
))
.unwrap();
let append_record: SessionMetadataRecord =
serde_json::from_slice(&fs::read(metadata_path_for_session(&session)).unwrap())
.unwrap();
assert!(!append_record.complete);
assert_eq!(append_record.latest_title, None);
let summaries = manager.list_metadata_summaries().unwrap();
let rebuilt: SessionMetadataRecord =
serde_json::from_slice(&fs::read(metadata_path_for_session(&session)).unwrap())
.unwrap();
assert_eq!(summaries[0].latest_title.as_deref(), Some("Old Title"));
assert!(rebuilt.complete);
}
#[test]
fn metadata_rebuild_does_not_publish_stale_complete_under_append() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
for iteration in 0..32 {
let session = manager.open(format!("rebuild-race-{iteration}")).unwrap();
record_session_event(
Some(&session),
temp.path(),
"user_input",
json!({"text":"seed"}),
)
.unwrap();
fs::remove_file(metadata_path_for_session(&session)).unwrap();
let append_session = session.clone();
let cwd = temp.path().to_path_buf();
let handle = std::thread::spawn(move || {
std::thread::sleep(std::time::Duration::from_micros(50 * iteration));
append_session
.append(&SessionEvent::new(
"user_input",
append_session.id().to_string(),
cwd,
json!({"text":"concurrent"}),
))
.unwrap();
});
let record = session_metadata_for_listing(&session).unwrap();
handle.join().unwrap();
let final_marker = jsonl_marker(&session).unwrap();
let actual_user_inputs = session
.read_events_tolerant()
.unwrap()
.events
.iter()
.filter(|event| event.kind() == Some(SessionEventKind::UserInput))
.count();
if record.complete
&& record.jsonl_len == final_marker.len
&& record.jsonl_modified_ns == final_marker.modified_ns
{
assert_eq!(record.user_input_count, actual_user_inputs);
}
}
}
#[test]
fn metadata_sidecar_write_failure_warns_but_append_succeeds() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("metadata-fails").unwrap();
fs::create_dir_all(metadata_path_for_session(&session)).unwrap();
crate::sessions::store::secure_test_session_root(session.path().parent().unwrap());
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"durable"}),
))
.unwrap();
let events = session.read_events().unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].payload["text"], "durable");
assert!(metadata_path_for_session(&session).is_dir());
}
#[test]
fn concurrent_session_appends_update_metadata_without_losing_latest() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let mut handles = Vec::new();
for index in 0..16 {
let session = session.clone();
let cwd = temp.path().to_path_buf();
handles.push(std::thread::spawn(move || {
let mut event = SessionEvent::new(
"diagnostic",
session.id().to_string(),
cwd,
json!({"index": index}),
);
event.timestamp = Utc.with_ymd_and_hms(2024, 3, 1, 0, 0, index).unwrap();
session.append(&event).unwrap();
}));
}
for handle in handles {
handle.join().unwrap();
}
let record: SessionMetadataRecord =
serde_json::from_slice(&fs::read(metadata_path_for_session(&session)).unwrap())
.unwrap();
let tolerant = session.read_events_tolerant().unwrap();
assert!(tolerant.diagnostics.is_empty());
assert_eq!(tolerant.events.len(), 16);
assert_eq!(record.latest_activity_timestamp.unwrap().second(), 15);
}
fn append_event_at(session: &Session, cwd: &Path, year: i32, month: u32, day: u32) {
session
.append(&event_at(session, cwd, year, month, day))
.unwrap();
}
fn event_at(session: &Session, cwd: &Path, year: i32, month: u32, day: u32) -> SessionEvent {
let mut event = SessionEvent::new(
"diagnostic",
session.id().to_string(),
cwd.to_path_buf(),
json!({}),
);
event.timestamp = Utc.with_ymd_and_hms(year, month, day, 0, 0, 0).unwrap();
event
}
}