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,
Title,
}
impl SessionDiagnosticOperation {
fn as_str(self) -> &'static str {
match self {
Self::Replay => "replay",
Self::Metadata => "metadata",
Self::Title => "title",
}
}
}
pub(crate) fn report_session_diagnostic(
operation: SessionDiagnosticOperation,
_path: &Path,
error: impl std::fmt::Display,
) {
let error = redact_sensitive_text(&error.to_string())
.chars()
.take(500)
.collect::<String>();
crate::output::emit_terminal_warning(format!(
"warning: session operation={} category=session_jsonl failed: {error}",
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 JsonlMarker {
fn from_metadata(metadata: &fs::Metadata) -> Self {
Self {
len: metadata.len(),
modified_ns: metadata.modified().ok().map(system_time_ns),
}
}
}
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::new(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");
};
Ok(JsonlMarker::from_metadata(&file.metadata()?))
}
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 metadata_for_append(
session: &Session,
current_metadata: &fs::Metadata,
) -> anyhow::Result<SessionMetadataRecord> {
let marker = JsonlMarker::from_metadata(current_metadata);
if let Some(record) = session
.metadata_cache()
.lock()
.map_err(|_| anyhow::anyhow!("session metadata cache was poisoned"))?
.clone()
&& record.complete
&& record.jsonl_len == marker.len
&& record.jsonl_modified_ns == marker.modified_ns
{
return Ok(record);
}
let record = rebuild_session_metadata_record_from_jsonl(session)?;
cache_session_metadata(session, Some(record.clone()))?;
Ok(record)
}
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))
}
fn rebuild_session_metadata_record_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;
Ok(record)
}
fn cache_session_metadata(
session: &Session,
record: Option<SessionMetadataRecord>,
) -> anyhow::Result<()> {
*session
.metadata_cache()
.lock()
.map_err(|_| anyhow::anyhow!("session metadata cache was poisoned"))? = record;
Ok(())
}
pub(crate) fn invalidate_session_metadata_cache(session: &Session) -> anyhow::Result<()> {
cache_session_metadata(session, None)
}
pub(crate) fn rebuild_session_metadata_from_jsonl(
session: &Session,
) -> anyhow::Result<SessionMetadataRecord> {
let record = rebuild_session_metadata_record_from_jsonl(session)?;
cache_session_metadata(session, Some(record.clone()))?;
if let Err(error) = write_session_metadata(session, &record)
&& !is_session_file_not_found(&error)
{
report_session_diagnostic(
SessionDiagnosticOperation::Metadata,
&metadata_path_for_session(session),
&error,
);
}
Ok(record)
}
fn is_session_file_not_found(error: &anyhow::Error) -> bool {
error.chain().any(|cause| {
cause
.downcast_ref::<std::io::Error>()
.is_some_and(|io_error| io_error.kind() == std::io::ErrorKind::NotFound)
})
}
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;
cache_session_metadata(session, Some(record.clone()))?;
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)? {
cache_session_metadata(session, Some(record.clone()))?;
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],
sidecar_exists: bool,
) -> anyhow::Result<()> {
let mut record = previous.unwrap_or_else(|| {
SessionMetadataRecord::empty(
session,
JsonlMarker {
len: 0,
modified_ns: None,
},
false,
)
});
let was_complete = record.complete;
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;
record.complete = true;
let requires_write = !sidecar_exists
|| events.iter().any(|event| {
matches!(
event.kind(),
Some(
SessionEventKind::SessionTitle
| SessionEventKind::UserInput
| SessionEventKind::AssistantOutput
| SessionEventKind::TurnStatus
)
)
})
|| (!was_complete && record.complete);
cache_session_metadata(session, Some(record.clone()))?;
if requires_write {
write_session_metadata(session, &record)?;
}
Ok(())
}
pub(crate) fn latest_valid_compaction_checkpoint_for_replay(
session_id: &str,
events: &[SessionEvent],
) -> (Option<CompactionCheckpoint>, Vec<String>) {
latest_valid_compaction_checkpoint_with_policy(session_id, events, true)
}
fn latest_valid_compaction_checkpoint_with_policy(
session_id: &str,
events: &[SessionEvent],
skip_foreign_session_events: bool,
) -> (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;
}
if skip_foreign_session_events && event.session_id != session_id {
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)
}
pub(super) 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)
}