use super::event::{SessionEvent, SessionEventKind};
use super::manager::Session;
use super::metadata::{
metadata_path_for_session, read_session_metadata_for_append, report_session_diagnostic,
update_session_metadata_after_append_batch,
};
use super::read::validate_session_id;
use crate::{
output::{is_credential_like_key, redact_sensitive_text},
persistence::CrossProcessFileLock,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SessionAppendOutcome {
Durable,
DurableJsonlMetadataUpdateFailed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SessionAppendError {
NotWritten(String),
RolledBack(String),
Uncertain(String),
}
impl std::fmt::Display for SessionAppendError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::NotWritten(message) => write!(formatter, "session append not written: {message}"),
Self::RolledBack(message) => write!(formatter, "session append rolled back: {message}"),
Self::Uncertain(message) => {
write!(
formatter,
"session append outcome uncertain: on-disk state uncertain: {message}"
)
}
}
}
}
impl std::error::Error for SessionAppendError {}
use anyhow::Context;
use serde_json::Value;
use sha2::{Digest, Sha256};
use std::{
collections::HashMap,
fs,
io::{Read, Seek, Write},
path::{Path, PathBuf},
sync::{Arc, Mutex, OnceLock, Weak},
};
pub const SESSION_TITLE_EVENT: &str = SessionEventKind::SessionTitle.as_str();
pub fn record_session_event(
session: Option<&Session>,
cwd: &Path,
event_type: impl AsRef<str>,
payload: Value,
) -> anyhow::Result<SessionAppendOutcome> {
let Some(session) = session else {
return Ok(SessionAppendOutcome::Durable);
};
session.append_with_outcome(&SessionEvent::new(
event_type.as_ref(),
session.id().to_string(),
cwd.to_path_buf(),
payload,
))
}
pub fn try_record_session_event(
session: Option<&Session>,
cwd: &Path,
event_type: impl AsRef<str>,
payload: Value,
) -> anyhow::Result<SessionAppendOutcome> {
record_session_event(session, cwd, event_type, payload)
}
pub fn sanitize_session_title(raw: &str) -> Option<String> {
let trimmed = raw
.trim()
.trim_matches(|ch| matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’'));
let collapsed = trimmed
.chars()
.map(|ch| {
if matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’') {
return ' ';
}
if ch.is_control() || ch.is_whitespace() {
' '
} else {
ch
}
})
.collect::<String>();
let collapsed = collapsed.split_whitespace().collect::<Vec<_>>().join(" ");
let title = collapsed
.trim()
.trim_matches(|ch| matches!(ch, '"' | '\'' | '“' | '”' | '‘' | '’'))
.trim();
if title.is_empty() {
return None;
}
Some(
title
.chars()
.take(crate::sessions::titles::SESSION_TITLE_MAX_CHARS)
.collect(),
)
}
pub fn record_session_title(
session: &Session,
cwd: &Path,
title: &str,
provider: &str,
model: &str,
) -> anyhow::Result<()> {
let Some(title) = sanitize_session_title(title) else {
anyhow::bail!("session title is empty after sanitization");
};
session.append(&SessionEvent::new_kind(
SessionEventKind::SessionTitle,
session.id().to_string(),
cwd.to_path_buf(),
serde_json::json!({"title": title, "provider": provider, "model": model}),
))
}
pub(crate) const COMPACTION_SCHEMA_VERSION: u64 = 1;
pub(crate) fn sanitize_compaction_summary(raw: &str) -> Option<String> {
let cleaned = raw
.trim()
.chars()
.map(|ch| {
if ch.is_control() && !matches!(ch, '\n' | '\t') {
' '
} else {
ch
}
})
.collect::<String>();
let redacted = redact_sensitive_text(cleaned.trim());
if redacted.trim().is_empty() {
None
} else {
Some(redacted.trim().to_string())
}
}
#[cfg(test)]
pub(crate) fn record_session_compaction(
session: &Session,
cwd: &Path,
summary: &str,
provider: &str,
model: &str,
cutoff_event_count: usize,
) -> anyhow::Result<()> {
let bytes = session
.path
.exists()
.then(|| fs::read(&session.path))
.transpose()?
.unwrap_or_default();
let offset = bytes
.split_inclusive(|byte| *byte == b'\n')
.take(cutoff_event_count)
.map(|line| line.len())
.sum::<usize>();
record_session_compaction_at_byte_offset_with_snapshot(
session,
cwd,
summary,
provider,
model,
offset as u64,
capture_session_snapshot(&session.path)?,
)
}
pub(crate) fn record_session_compaction_at_byte_offset_with_snapshot(
session: &Session,
cwd: &Path,
summary: &str,
provider: &str,
model: &str,
cutoff_byte_offset: u64,
snapshot: Option<SessionSnapshot>,
) -> anyhow::Result<()> {
let Some(summary) = sanitize_compaction_summary(summary) else {
anyhow::bail!("compaction summary is empty after sanitization");
};
let provider = provider.trim();
let model = model.trim();
if provider.is_empty() || model.is_empty() {
anyhow::bail!("compaction provider and model must be non-empty");
}
rotate_session_history(
session,
cwd,
&summary,
provider,
model,
cutoff_byte_offset,
snapshot,
)
}
fn rotate_session_history(
session: &Session,
cwd: &Path,
summary: &str,
provider: &str,
model: &str,
cutoff_byte_offset: u64,
snapshot: Option<SessionSnapshot>,
) -> anyhow::Result<()> {
validate_session_id(session.id.clone())?;
let parent = session
.path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent directory"))?;
ensure_directory(parent)?;
let append_lock = session_append_lock(&session.path)?;
let _append_guard = append_lock
.lock()
.map_err(|_| anyhow::anyhow!("session append lock was poisoned"))?;
let _file_guard = CrossProcessFileLock::acquire(&session.path)?;
validate_session_snapshot(&session.path, snapshot.as_ref())?;
let previous_metadata = if session.path.exists() {
match super::metadata::read_complete_session_metadata(session)? {
Some(record) => Some(record),
None => Some(super::metadata::rebuild_session_metadata_from_jsonl(
session,
)?),
}
} else {
None
};
let old_len = open_secure_session_path(&session.path)?
.as_ref()
.map(|file| file.metadata().map(|metadata| metadata.len()))
.transpose()?
.unwrap_or(0);
let suffix_start = cutoff_byte_offset.min(old_len);
let checkpoint = SessionEvent::new_kind(
SessionEventKind::Compaction,
session.id.to_string(),
cwd.to_path_buf(),
serde_json::json!({"schema_version": COMPACTION_SCHEMA_VERSION, "summary": summary, "provider": provider, "model": model, "cutoff_event_count": 0, "aggregate": previous_metadata.as_ref().map(|record| record.checkpoint_aggregate()).unwrap_or_else(|| serde_json::json!({}))}),
);
let history_parent = parent.join(".history");
ensure_directory_or_create(&history_parent)?;
let history_root = history_parent.join(&session.id);
ensure_directory_or_create(&history_root)?;
remove_stale_rotation_temps(parent, &session.id)?;
let generation = fs::read_dir(&history_root)?
.filter_map(Result::ok)
.filter_map(|entry| {
entry
.file_name()
.to_str()
.and_then(|name| name.strip_suffix(".jsonl"))
.and_then(|name| name.parse::<u64>().ok())
})
.max()
.unwrap_or(0)
.saturating_add(1);
remove_stale_rotation_archive_temps(&history_root, generation)?;
let archive = history_root.join(format!("{generation}.jsonl"));
let archive_temp = history_root.join(format!(".{generation}.jsonl.compact-tmp"));
let archive_result = (|| -> anyhow::Result<()> {
let mut archive_file = fs::OpenOptions::new();
archive_file.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
archive_file.mode(super::store::SESSION_FILE_MODE);
}
let mut archive_file = archive_file.open(&archive_temp)?;
if let Some(mut source) = open_secure_session_path(&session.path)? {
std::io::copy(&mut source, &mut archive_file)?;
}
archive_file.sync_all()?;
fs::rename(&archive_temp, &archive)?;
sync_session_parent_dir(&history_root)?;
Ok(())
})();
if let Err(error) = archive_result {
remove_validated_rotation_temp(&archive_temp);
return Err(error);
}
let checkpoint_line = serialize_session_event_line(&checkpoint)?;
let temporary = parent.join(format!(
".{}.jsonl.compact-{}-{}",
session.id,
std::process::id(),
generation
));
let mut active_file = fs::OpenOptions::new();
active_file.write(true).create_new(true);
#[cfg(unix)]
{
use std::os::unix::fs::OpenOptionsExt;
active_file.mode(super::store::SESSION_FILE_MODE);
}
let mut active_file = active_file.open(&temporary)?;
active_file.write_all(&checkpoint_line)?;
if let Some(mut source) = open_secure_session_path(&session.path)? {
source.seek(std::io::SeekFrom::Start(suffix_start))?;
std::io::copy(&mut source, &mut active_file)?;
}
active_file.sync_all()?;
if let Err(error) = fs::rename(&temporary, &session.path) {
remove_validated_rotation_temp(&temporary);
return Err(error.into());
}
sync_session_parent_dir(parent).context("active session replacement committed but directory sync failed; active JSONL remains recoverable")?;
if let Some(previous) = previous_metadata
&& let Err(error) =
super::metadata::write_rotated_session_metadata(session, previous, &checkpoint)
{
report_session_diagnostic(
super::metadata::SessionDiagnosticOperation::Metadata,
&metadata_path_for_session(session),
error,
);
}
Ok(())
}
fn open_secure_session_path(path: &Path) -> anyhow::Result<Option<fs::File>> {
let root = path
.parent()
.ok_or_else(|| anyhow::anyhow!("session file has no parent"))?;
let name = path
.file_name()
.and_then(|name| name.to_str())
.ok_or_else(|| anyhow::anyhow!("session file has no filename"))?;
let id = name
.strip_suffix(".jsonl")
.ok_or_else(|| anyhow::anyhow!("session file has unexpected name"))?;
super::store::open_existing_primary(root, id)
}
#[derive(Debug)]
pub(crate) struct SessionSnapshot {
len: u64,
prefix: Vec<u8>,
digest: [u8; 32],
identity: Option<(u64, u64)>,
}
pub(crate) fn capture_session_snapshot(path: &Path) -> anyhow::Result<Option<SessionSnapshot>> {
let Some(mut file) = open_secure_session_path(path)? else {
return Ok(None);
};
let metadata = file.metadata()?;
let mut prefix = vec![0; usize::try_from(metadata.len().min(4096)).unwrap_or(4096)];
file.read_exact(&mut prefix)?;
let mut digest = Sha256::new();
digest.update(&prefix);
let mut remaining = metadata.len().saturating_sub(prefix.len() as u64);
let mut buffer = [0u8; 8192];
while remaining > 0 {
let read_len = usize::try_from(remaining.min(buffer.len() as u64)).unwrap_or(buffer.len());
file.read_exact(&mut buffer[..read_len])?;
digest.update(&buffer[..read_len]);
remaining -= read_len as u64;
}
Ok(Some(SessionSnapshot {
len: metadata.len(),
prefix,
digest: digest.finalize().into(),
identity: file_identity(&metadata),
}))
}
fn validate_session_snapshot(
path: &Path,
snapshot: Option<&SessionSnapshot>,
) -> anyhow::Result<()> {
let current = capture_session_snapshot(path)?;
match (snapshot, current.as_ref()) {
(None, None) => Ok(()),
(Some(expected), Some(actual))
if actual.len >= expected.len
&& actual.prefix.starts_with(&expected.prefix)
&& (expected.identity == actual.identity
|| (expected.identity.is_none()
&& actual.len == expected.len
&& actual.digest == expected.digest)) =>
{
Ok(())
}
_ => anyhow::bail!(
"session changed during compaction; rotation aborted safely (snapshot/current generation mismatch)"
),
}
}
fn file_identity(metadata: &fs::Metadata) -> Option<(u64, u64)> {
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
Some((metadata.dev(), metadata.ino()))
}
#[cfg(not(unix))]
{
let _ = metadata;
None
}
}
fn remove_validated_rotation_temp(path: &Path) {
if let Ok(metadata) = fs::symlink_metadata(path)
&& metadata.file_type().is_file()
&& !metadata.file_type().is_symlink()
{
let _ = fs::remove_file(path);
}
}
fn ensure_directory(path: &Path) -> anyhow::Result<()> {
let metadata = fs::symlink_metadata(path)?;
if !metadata.file_type().is_dir() || metadata.file_type().is_symlink() {
anyhow::bail!("session history path is not a directory");
}
Ok(())
}
fn ensure_directory_or_create(path: &Path) -> anyhow::Result<()> {
if !path.exists() {
fs::create_dir(path)?;
if let Some(parent) = path.parent() {
sync_session_parent_dir(parent)?;
}
}
ensure_directory(path)
}
fn remove_stale_rotation_temps(parent: &Path, id: &str) -> anyhow::Result<()> {
let prefix = format!(".{id}.jsonl.compact-");
for entry in fs::read_dir(parent)? {
let entry = entry?;
if !entry.file_name().to_string_lossy().starts_with(&prefix) {
continue;
}
let metadata = fs::symlink_metadata(entry.path())?;
if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
continue;
}
fs::remove_file(entry.path())?;
}
Ok(())
}
fn remove_stale_rotation_archive_temps(history_root: &Path, generation: u64) -> anyhow::Result<()> {
let prefix = format!(".{generation}.jsonl.compact-tmp");
for entry in fs::read_dir(history_root)? {
let entry = entry?;
if !entry.file_name().to_string_lossy().starts_with(&prefix) {
continue;
}
let metadata = fs::symlink_metadata(entry.path())?;
if metadata.file_type().is_symlink() || !metadata.file_type().is_file() {
continue;
}
fs::remove_file(entry.path())?;
}
Ok(())
}
const MAX_REDACTION_DEPTH: usize = 128;
fn redact_json_value(value: Value) -> Value {
redact_json_value_at_depth(value, 0)
}
fn redact_json_value_at_depth(value: Value, depth: usize) -> Value {
match value {
Value::String(text) => Value::String(redact_sensitive_text(&text)),
Value::Array(items) => {
if depth >= MAX_REDACTION_DEPTH {
return Value::String("<redacted: max depth>".to_string());
}
Value::Array(
items
.into_iter()
.map(|value| redact_json_value_at_depth(value, depth + 1))
.collect(),
)
}
Value::Object(map) => {
if depth >= MAX_REDACTION_DEPTH {
return Value::String("<redacted: max depth>".to_string());
}
Value::Object(
map.into_iter()
.map(|(key, value)| {
let redacted = if is_credential_like_key(&key) {
match value {
Value::String(_) => Value::String("<redacted>".to_string()),
other => redact_json_value_at_depth(other, depth + 1),
}
} else {
redact_json_value_at_depth(value, depth + 1)
};
(key, redacted)
})
.collect(),
)
}
other => other,
}
}
trait DurableSessionWriter {
fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()>;
fn flush_durable(&mut self) -> std::io::Result<()>;
fn sync_data_durable(&mut self) -> std::io::Result<()>;
fn set_len_durable(&mut self, len: u64) -> std::io::Result<()>;
}
impl DurableSessionWriter for fs::File {
fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
self.write_all(bytes)
}
fn flush_durable(&mut self) -> std::io::Result<()> {
self.flush()
}
fn sync_data_durable(&mut self) -> std::io::Result<()> {
self.sync_data()
}
fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
self.set_len(len)
}
}
fn serialize_session_event_line(event: &SessionEvent) -> anyhow::Result<Vec<u8>> {
let mut line = serde_json::to_vec(event).context("failed to serialize session event")?;
line.push(b'\n');
Ok(line)
}
fn rollback_session_append(writer: &mut dyn DurableSessionWriter, pre_append_len: u64) -> bool {
writer.set_len_durable(pre_append_len).is_ok()
}
fn append_event_line_durably(
writer: &mut dyn DurableSessionWriter,
line: &[u8],
pre_append_len: u64,
) -> Result<(), SessionAppendError> {
let append = |writer: &mut dyn DurableSessionWriter, error: std::io::Error, stage: &str| {
let message = format!("failed to {stage} session JSONL event: {error}");
if rollback_session_append(writer, pre_append_len) {
SessionAppendError::RolledBack(message)
} else {
SessionAppendError::Uncertain(message)
}
};
if let Err(error) = writer.write_all_durable(line) {
return Err(append(writer, error, "write"));
}
if let Err(error) = writer.flush_durable() {
return Err(append(writer, error, "flush"));
}
if let Err(error) = writer.sync_data_durable() {
return Err(append(writer, error, "sync"));
}
Ok(())
}
fn sync_session_parent_dir(parent: &Path) -> std::io::Result<()> {
#[cfg(unix)]
{
fs::File::open(parent)?.sync_all()?;
}
#[cfg(not(unix))]
{
let _ = parent;
}
Ok(())
}
static SESSION_APPEND_LOCKS: OnceLock<Mutex<HashMap<PathBuf, Weak<Mutex<()>>>>> = OnceLock::new();
fn session_append_lock(path: &Path) -> anyhow::Result<Arc<Mutex<()>>> {
let key = normalize_session_append_lock_path(path);
let registry = SESSION_APPEND_LOCKS.get_or_init(|| Mutex::new(HashMap::new()));
let mut locks = registry
.lock()
.map_err(|_| anyhow::anyhow!("session append lock registry was poisoned"))?;
if let Some(lock) = locks.get(&key).and_then(Weak::upgrade) {
return Ok(lock);
}
locks.retain(|_, lock| lock.strong_count() > 0);
let lock = Arc::new(Mutex::new(()));
locks.insert(key, Arc::downgrade(&lock));
Ok(lock)
}
fn normalize_session_append_lock_path(path: &Path) -> PathBuf {
if let Ok(canonical) = path.canonicalize() {
return canonical;
}
if let (Some(parent), Some(file_name)) = (path.parent(), path.file_name())
&& let Ok(parent) = parent.canonicalize()
{
return parent.join(file_name);
}
path.to_path_buf()
}
impl Session {
pub fn append(&self, event: &SessionEvent) -> anyhow::Result<()> {
self.append_with_outcome(event).map(|_| ())
}
pub fn append_with_outcome(
&self,
event: &SessionEvent,
) -> anyhow::Result<SessionAppendOutcome> {
self.append_owned_batch(vec![event.clone()])
.map_err(anyhow::Error::new)
}
fn append_owned_batch(
&self,
mut events: Vec<SessionEvent>,
) -> Result<SessionAppendOutcome, SessionAppendError> {
let not_written =
|error: &dyn std::fmt::Display| SessionAppendError::NotWritten(error.to_string());
validate_session_id(self.id.clone()).map_err(|error| not_written(&error))?;
let parent = self.path.parent().ok_or_else(|| {
SessionAppendError::NotWritten("session file has no parent".to_string())
})?;
super::store::prepare_session_root(parent).map_err(|error| not_written(&error))?;
let mut line = Vec::new();
for event in &mut events {
validate_session_id(event.session_id.clone()).map_err(|error| not_written(&error))?;
if event.session_id != self.id {
return Err(SessionAppendError::NotWritten(
"event session id does not match session".to_string(),
));
}
event.payload = redact_json_value(std::mem::take(&mut event.payload));
line.extend(serialize_session_event_line(event).map_err(|error| not_written(&error))?);
}
let append_lock = session_append_lock(&self.path).map_err(|error| not_written(&error))?;
let _append_guard = append_lock.lock().map_err(|_| {
SessionAppendError::NotWritten("session append lock was poisoned".to_string())
})?;
let _file_guard =
CrossProcessFileLock::acquire(&self.path).map_err(|error| not_written(&error))?;
let previous_metadata = read_session_metadata_for_append(self).ok().flatten();
let (mut file, is_new_session_file) =
super::store::open_primary(parent, &self.id).map_err(|error| not_written(&error))?;
let pre_append_len = file.metadata().map_err(|error| not_written(&error))?.len();
append_event_line_durably(&mut file, &line, pre_append_len)?;
if is_new_session_file {
sync_session_parent_dir(parent)
.map_err(|error| SessionAppendError::Uncertain(error.to_string()))?;
}
if let Err(error) =
update_session_metadata_after_append_batch(self, previous_metadata, &events)
{
report_session_diagnostic(
super::metadata::SessionDiagnosticOperation::Metadata,
&metadata_path_for_session(self),
&error,
);
return Ok(SessionAppendOutcome::DurableJsonlMetadataUpdateFailed);
}
Ok(SessionAppendOutcome::Durable)
}
}
pub(crate) fn try_record_session_event_batch(
session: Option<&Session>,
cwd: &Path,
event_type: impl AsRef<str>,
payloads: impl IntoIterator<Item = Value>,
) -> anyhow::Result<SessionAppendOutcome> {
let Some(session) = session else {
return Ok(SessionAppendOutcome::Durable);
};
let event_type = event_type.as_ref();
let events = payloads
.into_iter()
.map(|payload| {
SessionEvent::new(
event_type,
session.id().to_string(),
cwd.to_path_buf(),
payload,
)
})
.collect();
session
.append_owned_batch(events)
.map_err(anyhow::Error::new)
}
#[cfg(test)]
mod tests {
use super::super::manager::SessionManager;
use super::*;
use serde_json::json;
use tempfile::TempDir;
#[test]
fn session_jsonl_append_and_read_round_trip() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let event = SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"hello"}),
);
session.append(&event).unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
assert_eq!(
fs::symlink_metadata(session.path().parent().unwrap())
.unwrap()
.mode()
& 0o777,
0o700
);
assert_eq!(
fs::symlink_metadata(session.path()).unwrap().mode() & 0o777,
0o600
);
}
let events = session.read_events().unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].event_type, "user_input");
assert_eq!(events[0].payload["text"], "hello");
assert_eq!(events[0].session_path.as_deref(), Some(temp.path()));
}
#[test]
fn compaction_rotates_history_and_keeps_active_suffix_bounded() {
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":"before"}),
))
.unwrap();
let original = fs::read(session.path()).unwrap();
record_session_compaction(&session, temp.path(), "summary", "p", "m", 1).unwrap();
let history = temp
.path()
.join("sessions/.history")
.join(session.id())
.join("1.jsonl");
assert_eq!(fs::read(history).unwrap(), original);
let active = session.read_events().unwrap();
assert_eq!(active[0].event_type, "compaction");
assert_eq!(active[0].payload["cutoff_event_count"], 0);
assert_eq!(active.len(), 1);
session
.append(&SessionEvent::new(
"user_input",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"text":"after"}),
))
.unwrap();
record_session_compaction(&session, temp.path(), "summary 2", "p", "m", 2).unwrap();
assert!(
temp.path()
.join("sessions/.history")
.join(session.id())
.join("2.jsonl")
.exists()
);
assert!(
manager
.list()
.unwrap()
.iter()
.all(|item| item.id() == session.id())
);
}
#[test]
fn snapshot_rejects_replacement_with_shared_prefix_and_growth() {
let temp = TempDir::new().unwrap();
let path = temp.path().join("session.jsonl");
fs::write(&path, b"shared prefix\noriginal generation\n").unwrap();
crate::sessions::store::secure_test_session_root(temp.path());
let snapshot = capture_session_snapshot(&path).unwrap();
let backup = temp.path().join("session.backup");
fs::rename(&path, backup).unwrap();
fs::write(
&path,
b"shared prefix\nreplacement generation with growth\n",
)
.unwrap();
crate::sessions::store::secure_test_session_root(temp.path());
let error = validate_session_snapshot(&path, Some(snapshot.as_ref().unwrap())).unwrap_err();
assert!(error.to_string().contains("generation mismatch"));
}
#[test]
fn session_append_lock_is_shared_per_jsonl_path() {
let temp = TempDir::new().unwrap();
let path = temp.path().join("session.jsonl");
let first = session_append_lock(&path).unwrap();
let second = session_append_lock(&path).unwrap();
let other = session_append_lock(&temp.path().join("other.jsonl")).unwrap();
assert!(Arc::ptr_eq(&first, &second));
assert!(!Arc::ptr_eq(&first, &other));
}
#[test]
fn concurrent_session_appends_remain_parseable_jsonl() {
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..24 {
let session = session.clone();
let cwd = temp.path().to_path_buf();
handles.push(std::thread::spawn(move || {
let event = SessionEvent::new(
"diagnostic",
session.id().to_string(),
cwd,
json!({"index": index}),
);
session.append(&event).unwrap();
}));
}
for handle in handles {
handle.join().unwrap();
}
let tolerant = session.read_events_tolerant().unwrap();
assert!(
tolerant.diagnostics.is_empty(),
"{:?}",
tolerant.diagnostics
);
assert_eq!(tolerant.events.len(), 24);
}
#[test]
fn session_append_rejects_mismatched_or_unsafe_event_ids() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.open("safe").unwrap();
let mismatched = SessionEvent::new(
"user_input",
"other".to_string(),
temp.path().to_path_buf(),
json!({"text":"hello"}),
);
assert!(session.append(&mismatched).is_err());
let unsafe_event = SessionEvent::new(
"user_input",
"../escape".to_string(),
temp.path().to_path_buf(),
json!({"text":"hello"}),
);
assert!(session.append(&unsafe_event).is_err());
}
#[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 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 redacted_session_events_remain_readable_for_resume() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
let secret = "sk-testSecret123456";
session
.append(&SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({"tool":"bash", "success": true, "output": secret}),
))
.unwrap();
let events = session.read_events().unwrap();
assert_eq!(events[0].event_type, "tool_result");
assert_eq!(events[0].payload["tool"], "bash");
assert_eq!(events[0].payload["success"], true);
assert!(!events[0].payload.to_string().contains(secret));
}
#[test]
fn persisted_redaction_preserves_generic_fields_and_hides_credentials() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new(
"tool_result",
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"code": 200,
"key": "value",
"nested": {"code": "ok", "api_key": "api-secret"},
"provider_api_key": "provider-secret"
}),
))
.unwrap();
let payload = session.read_events().unwrap().remove(0).payload;
assert_eq!(payload["code"], 200);
assert_eq!(payload["key"], "value");
assert_eq!(payload["nested"]["code"], "ok");
assert_eq!(payload["nested"]["api_key"], "<redacted>");
assert_eq!(payload["provider_api_key"], "<redacted>");
let persisted = fs::read_to_string(session.path()).unwrap();
assert!(persisted.contains("\"code\":200"));
assert!(persisted.contains("\"key\":\"value\""));
assert!(!persisted.contains("api-secret"));
assert!(!persisted.contains("provider-secret"));
}
#[test]
fn redaction_preserves_numeric_token_usage_fields() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new(
"usage",
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"input_tokens": 12,
"output_tokens": 34,
"cache_creation_input_tokens": 56,
"token_estimate": 78,
"access_token": "secret-token",
}),
))
.unwrap();
let payload = session.read_events().unwrap().remove(0).payload;
assert_eq!(payload["input_tokens"], 12);
assert_eq!(payload["output_tokens"], 34);
assert_eq!(payload["cache_creation_input_tokens"], 56);
assert_eq!(payload["token_estimate"], 78);
assert_eq!(payload["access_token"], "<redacted>");
}
#[test]
fn provider_stream_trace_session_event_is_redacted_before_write() {
let temp = TempDir::new().unwrap();
let manager = SessionManager::new(temp.path().join("sessions"));
let session = manager.create().unwrap();
session
.append(&SessionEvent::new_kind(
SessionEventKind::ProviderStreamTrace,
session.id().to_string(),
temp.path().to_path_buf(),
json!({
"schema_version": 1,
"provider": "anthropic",
"api_key": "sk-secret123456",
"recent_events": [{"seq": 1, "usage": {"output_tokens": 7}}],
"pending_tools": [{
"index": 1,
"id": "toolu_1",
"name": "read",
"argument_bytes": 13,
"argument_sha256": "a".repeat(64)
}]
}),
))
.unwrap();
let payload = session.read_events().unwrap().remove(0).payload;
assert_eq!(payload["api_key"], "<redacted>");
assert_eq!(payload["schema_version"], 1);
assert_eq!(payload["recent_events"][0]["seq"], 1);
assert_eq!(payload["recent_events"][0]["usage"]["output_tokens"], 7);
assert_eq!(payload["pending_tools"][0]["argument_bytes"], 13);
}
#[test]
fn redaction_caps_deep_json_nesting() {
let mut value = json!({"api_key":"sk-testSecret123456"});
for _ in 0..(MAX_REDACTION_DEPTH + 8) {
value = json!([value]);
}
let redacted = redact_json_value(value);
let text = redacted.to_string();
assert!(text.contains("<redacted: max depth>"), "{text}");
assert!(!text.contains("sk-testSecret"), "{text}");
}
#[test]
fn session_append_failure_rolls_back_partial_jsonl_tail() {
struct PartialFailWriter {
bytes: Vec<u8>,
}
impl DurableSessionWriter for PartialFailWriter {
fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
self.bytes.extend_from_slice(&bytes[..bytes.len() / 2]);
Err(std::io::Error::other("simulated partial write"))
}
fn flush_durable(&mut self) -> std::io::Result<()> {
Ok(())
}
fn sync_data_durable(&mut self) -> std::io::Result<()> {
Ok(())
}
fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
self.bytes.truncate(usize::try_from(len).unwrap());
Ok(())
}
}
let original = b"{\"event_type\":\"existing\"}\n".to_vec();
let mut writer = PartialFailWriter {
bytes: original.clone(),
};
let error = append_event_line_durably(
&mut writer,
b"{\"event_type\":\"new\",\"payload\":{}}\n",
original.len() as u64,
)
.unwrap_err();
assert!(matches!(error, SessionAppendError::RolledBack(_)));
assert!(
error
.to_string()
.contains("failed to write session JSONL event")
);
assert_eq!(writer.bytes, original);
}
#[test]
fn session_append_reports_uncertain_state_when_rollback_fails() {
struct UnrecoverableWriter;
impl DurableSessionWriter for UnrecoverableWriter {
fn write_all_durable(&mut self, _: &[u8]) -> std::io::Result<()> {
Err(std::io::Error::other("write"))
}
fn flush_durable(&mut self) -> std::io::Result<()> {
Ok(())
}
fn sync_data_durable(&mut self) -> std::io::Result<()> {
Ok(())
}
fn set_len_durable(&mut self, _: u64) -> std::io::Result<()> {
Err(std::io::Error::other("rollback"))
}
}
let mut writer = UnrecoverableWriter;
let error = append_event_line_durably(&mut writer, b"line\n", 0).unwrap_err();
assert!(matches!(error, SessionAppendError::Uncertain(_)));
assert!(error.to_string().contains("on-disk state uncertain"));
}
#[test]
fn session_append_flushes_and_syncs_event_line() {
#[derive(Default)]
struct RecordingWriter {
bytes: Vec<u8>,
flushed: bool,
synced: bool,
}
impl DurableSessionWriter for RecordingWriter {
fn write_all_durable(&mut self, bytes: &[u8]) -> std::io::Result<()> {
self.bytes.extend_from_slice(bytes);
Ok(())
}
fn flush_durable(&mut self) -> std::io::Result<()> {
self.flushed = true;
Ok(())
}
fn sync_data_durable(&mut self) -> std::io::Result<()> {
self.synced = true;
Ok(())
}
fn set_len_durable(&mut self, len: u64) -> std::io::Result<()> {
self.bytes.truncate(usize::try_from(len).unwrap());
Ok(())
}
}
let temp = TempDir::new().unwrap();
let event = SessionEvent::new(
"user_input",
"safe".to_string(),
temp.path().to_path_buf(),
json!({"text":"hello"}),
);
let mut writer = RecordingWriter::default();
let line = serialize_session_event_line(&event).unwrap();
append_event_line_durably(&mut writer, &line, 0).unwrap();
assert!(writer.flushed);
assert!(writer.synced);
assert!(writer.bytes.ends_with(b"\n"));
let parsed: SessionEvent = serde_json::from_slice(&writer.bytes[..writer.bytes.len() - 1])
.expect("durable append writes one valid JSONL event");
assert_eq!(parsed.event_type, "user_input");
}
#[test]
fn try_record_session_event_returns_append_failure() {
let temp = TempDir::new().unwrap();
let invalid_session =
Session::unchecked_for_test("../unsafe".to_string(), temp.path().join("unsafe.jsonl"));
let error = try_record_session_event(
Some(&invalid_session),
temp.path(),
"assistant_chunk",
json!({"text":"best effort"}),
)
.unwrap_err();
assert!(error.to_string().contains("must not contain '..'"));
assert!(matches!(
error.downcast_ref::<SessionAppendError>(),
Some(SessionAppendError::NotWritten(_))
));
assert!(!invalid_session.path().exists());
}
}