use crate::event::AgentEvent;
use crate::security::{redact_memory, redact_snapshot_events};
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use std::fs;
use std::io::Read;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
use tokio::task;
static NEXT_SESSION_SEQUENCE: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionId(String);
impl SessionId {
pub fn new(id: String) -> Self {
Self(id)
}
pub fn as_str(&self) -> &str {
&self.0
}
pub fn into_inner(self) -> String {
self.0
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ProjectMemory {
pub project_hash: String,
pub entries: Vec<MemoryEntry>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MemoryEntry {
pub created_at: u64,
pub summary: String,
pub session_id: String,
}
pub fn session_title_from_events(events: &[AgentEvent]) -> Option<String> {
events
.iter()
.find_map(|event| match event {
AgentEvent::ModelOutput { text, .. } => title_from_model_text(text),
_ => None,
})
.or_else(|| {
events.iter().find_map(|event| match event {
AgentEvent::UserTaskSubmitted { text, .. } => title_from_user_text(text),
_ => None,
})
})
}
fn title_from_model_text(text: &str) -> Option<String> {
let heading = text.lines().find_map(|line| {
let trimmed = line.trim();
if trimmed.starts_with('#') {
Some(trimmed.trim_start_matches('#').trim())
} else {
None
}
});
heading
.and_then(clean_session_title)
.or_else(|| text.lines().find_map(clean_session_title))
}
fn title_from_user_text(text: &str) -> Option<String> {
clean_session_title(text)
}
pub fn clean_session_title(text: &str) -> Option<String> {
let cleaned = text
.trim()
.trim_matches('`')
.trim_matches('"')
.trim_matches('\'')
.trim_start_matches(['#', '-', '*', '>'])
.split_whitespace()
.collect::<Vec<_>>()
.join(" ");
if cleaned.is_empty() {
return None;
}
Some(
cleaned
.chars()
.take(80)
.collect::<String>()
.trim()
.to_string(),
)
}
impl ProjectMemory {
pub fn recent_entries(&self, max: usize) -> &[MemoryEntry] {
let start = self.entries.len().saturating_sub(max);
&self.entries[start..]
}
pub fn format_injection(&self, max: usize) -> Option<String> {
let entries = self.recent_entries(max);
if entries.is_empty() {
return None;
}
let mut parts = Vec::new();
for entry in entries {
parts.push(format!(
"[Session {} — {}]\n{}",
entry.session_id,
format_timestamp(entry.created_at),
entry.summary
));
}
Some(format!(
"Previous session context (summarized):\n\n{}",
parts.join("\n\n")
))
}
}
fn format_timestamp(unix_secs: u64) -> String {
let days = unix_secs / 86400;
let hours = (unix_secs % 86400) / 3600;
let minutes = (unix_secs % 3600) / 60;
format!("day {days} {hours:02}:{minutes:02}")
}
fn project_hash(project_dir: &Path) -> String {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
project_dir.hash(&mut hasher);
format!("{:016x}", hasher.finish())
}
#[derive(Debug, Clone)]
pub struct SessionStore {
root: PathBuf,
data_dir: PathBuf,
redact_secrets: bool,
}
fn default_session_version() -> u32 {
1
}
#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
pub struct SessionUsageSnapshot {
#[serde(default)]
pub input_tokens: u64,
#[serde(default)]
pub output_tokens: u64,
#[serde(default)]
pub cost_usd: f64,
#[serde(default)]
pub cost_known: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub credits_spent: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub credit_unit: Option<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionSnapshot {
#[serde(default = "default_session_version")]
pub version: u32,
pub id: SessionId,
#[serde(default)]
pub title: Option<String>,
pub project: PathBuf,
#[serde(default)]
pub created_at: u64,
#[serde(default)]
pub updated_at: u64,
pub events: Vec<AgentEvent>,
#[serde(default)]
pub memory: Option<ProjectMemory>,
#[serde(default)]
pub goal: Option<SessionGoal>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub usage: Option<SessionUsageSnapshot>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionSnapshotInfo {
pub id: SessionId,
#[serde(default)]
pub title: Option<String>,
pub project: PathBuf,
#[serde(default)]
pub created_at: u64,
#[serde(default)]
pub updated_at: u64,
}
impl SessionSnapshot {
pub const CURRENT_VERSION: u32 = 1;
}
fn is_session_snapshot_file(path: &Path) -> bool {
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
return false;
};
if !name.ends_with(".json") || name.ends_with(".ui.json") {
return false;
}
if name.contains(".tmp") || name.ends_with(".bak") {
return false;
}
path.extension().and_then(|e| e.to_str()) == Some("json")
}
fn read_session_info(path: &Path) -> Result<SessionSnapshotInfo> {
const METADATA_READ_LIMIT: usize = 64 * 1024;
let mut file =
fs::File::open(path).with_context(|| format!("failed to open {}", path.display()))?;
let mut buffer = vec![0; METADATA_READ_LIMIT];
let bytes_read = file
.read(&mut buffer)
.with_context(|| format!("failed to read {}", path.display()))?;
buffer.truncate(bytes_read);
let prefix = std::str::from_utf8(&buffer)
.with_context(|| format!("failed to decode metadata prefix from {}", path.display()))?;
if let Some(events_index) = prefix.find("\"events\"")
&& let Some(comma_index) = prefix[..events_index].rfind(',')
{
let metadata_json = format!("{}\n}}", &prefix[..comma_index]);
return serde_json::from_str::<SessionSnapshotInfo>(&metadata_json)
.with_context(|| format!("failed to parse metadata from {}", path.display()));
}
let content =
fs::read_to_string(path).with_context(|| format!("failed to read {}", path.display()))?;
serde_json::from_str::<SessionSnapshotInfo>(&content)
.with_context(|| format!("failed to parse metadata from {}", path.display()))
}
impl SessionStore {
pub fn new(data_dir: PathBuf) -> Self {
Self::with_redaction(data_dir, true)
}
pub fn with_redaction(data_dir: PathBuf, redact_secrets: bool) -> Self {
Self {
root: data_dir.join("sessions"),
data_dir,
redact_secrets,
}
}
pub fn root(&self) -> &PathBuf {
&self.root
}
pub fn create_id() -> SessionId {
let millis = current_unix_millis();
let pid = std::process::id();
let sequence = NEXT_SESSION_SEQUENCE.fetch_add(1, Ordering::Relaxed);
let nonce = fastrand::u64(..);
SessionId::new(format!("session-{millis}-{pid}-{sequence}-{nonce:016x}"))
}
pub fn save(&self, snapshot: &SessionSnapshot) -> Result<PathBuf> {
fs::create_dir_all(&self.root)
.with_context(|| format!("failed to create {}", self.root.display()))?;
crate::fs_util::set_private_dir_permissions(&self.root)?;
let path = self.root.join(format!("{}.json", snapshot.id.as_str()));
let snapshot = if self.redact_secrets {
SessionSnapshot {
version: snapshot.version,
id: snapshot.id.clone(),
title: snapshot.title.clone(),
project: snapshot.project.clone(),
created_at: snapshot.created_at,
updated_at: snapshot.updated_at,
goal: snapshot.goal.clone(),
events: redact_snapshot_events(&snapshot.events),
memory: snapshot.memory.as_ref().map(redact_memory),
usage: snapshot.usage.clone(),
}
} else {
snapshot.clone()
};
let data = serde_json::to_vec_pretty(&snapshot)?;
let tmp = path.with_extension("json.tmp");
fs::write(&tmp, &data).with_context(|| format!("failed to write {}", tmp.display()))?;
crate::fs_util::set_private_file_permissions(&tmp)?;
fs::rename(&tmp, &path).with_context(|| {
format!(
"failed to replace {} with {}",
path.display(),
tmp.display()
)
})?;
crate::fs_util::set_private_file_permissions(&path)?;
Ok(path)
}
pub async fn save_async(&self, snapshot: SessionSnapshot) -> Result<PathBuf> {
let store = self.clone();
task::spawn_blocking(move || store.save(&snapshot))
.await
.map_err(|err| anyhow::anyhow!("save_async join error: {err}"))?
}
pub fn list(&self) -> Vec<SessionSnapshot> {
let mut sessions = Vec::new();
if let Ok(entries) = fs::read_dir(&self.root) {
for entry in entries.flatten() {
let path = entry.path();
if !is_session_snapshot_file(&path) {
continue;
}
if let Ok(content) = fs::read_to_string(&path)
&& let Ok(snapshot) = serde_json::from_str::<SessionSnapshot>(&content)
{
sessions.push(snapshot);
}
}
}
sessions.sort_by(|a, b| {
b.updated_at
.cmp(&a.updated_at)
.then_with(|| b.id.as_str().cmp(a.id.as_str()))
});
sessions
}
pub fn list_info(&self) -> Vec<SessionSnapshotInfo> {
let mut sessions = Vec::new();
if let Ok(entries) = fs::read_dir(&self.root) {
for entry in entries.flatten() {
let path = entry.path();
if !is_session_snapshot_file(&path) {
continue;
}
if let Ok(info) = read_session_info(&path) {
sessions.push(info);
}
}
}
sessions.sort_by(|a, b| {
b.updated_at
.cmp(&a.updated_at)
.then_with(|| b.id.as_str().cmp(a.id.as_str()))
});
sessions
}
pub async fn list_async(&self) -> Vec<SessionSnapshot> {
let store = self.clone();
task::spawn_blocking(move || store.list())
.await
.unwrap_or_default()
}
pub async fn list_info_async(&self) -> Vec<SessionSnapshotInfo> {
let store = self.clone();
task::spawn_blocking(move || store.list_info())
.await
.unwrap_or_default()
}
pub fn load(&self, session_id: &str) -> Result<SessionSnapshot> {
let path = self.root.join(format!("{session_id}.json"));
let content = fs::read_to_string(&path)
.with_context(|| format!("failed to read {}", path.display()))?;
let snapshot: SessionSnapshot = serde_json::from_str(&content)
.with_context(|| format!("failed to parse {}", path.display()))?;
if snapshot.version > SessionSnapshot::CURRENT_VERSION {
return Err(anyhow::anyhow!(
"session snapshot version {} is newer than supported version {}",
snapshot.version,
SessionSnapshot::CURRENT_VERSION
));
}
Ok(snapshot)
}
pub async fn load_async(&self, session_id: String) -> Result<SessionSnapshot> {
let store = self.clone();
task::spawn_blocking(move || store.load(&session_id))
.await
.map_err(|err| anyhow::anyhow!("load_async join error: {err}"))?
}
pub fn delete(&self, session_id: &str) -> Result<bool> {
let path = self.root.join(format!("{session_id}.json"));
match fs::remove_file(&path) {
Ok(()) => Ok(true),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(err) => Err(err).with_context(|| format!("failed to delete {}", path.display())),
}
}
pub fn rename(&self, session_id: &str, title: &str) -> Result<bool> {
let title = title.trim();
if title.is_empty() {
return Err(anyhow::anyhow!("session title cannot be empty"));
}
let path = self.root.join(format!("{session_id}.json"));
if !path.exists() {
return Ok(false);
}
let mut snapshot = self.load(session_id)?;
snapshot.title = Some(title.to_string());
snapshot.updated_at = current_unix_timestamp();
self.save(&snapshot)?;
Ok(true)
}
pub async fn rename_async(&self, session_id: String, title: String) -> Result<bool> {
let store = self.clone();
task::spawn_blocking(move || store.rename(&session_id, &title))
.await
.map_err(|err| anyhow::anyhow!("rename_async join error: {err}"))?
}
pub async fn delete_async(&self, session_id: String) -> Result<bool> {
let store = self.clone();
task::spawn_blocking(move || store.delete(&session_id))
.await
.map_err(|err| anyhow::anyhow!("delete_async join error: {err}"))?
}
pub fn save_memory(&self, project_dir: &Path, memory: &ProjectMemory) -> Result<PathBuf> {
let memory_dir = self.data_dir.join("memory");
fs::create_dir_all(&memory_dir)
.with_context(|| format!("failed to create {}", memory_dir.display()))?;
crate::fs_util::set_private_dir_permissions(&memory_dir)?;
let hash = project_hash(project_dir);
let path = memory_dir.join(format!("{hash}.json"));
let data = serde_json::to_vec_pretty(memory)?;
fs::write(&path, data).with_context(|| format!("failed to write {}", path.display()))?;
crate::fs_util::set_private_file_permissions(&path)?;
Ok(path)
}
pub async fn save_memory_async(
&self,
project_dir: PathBuf,
memory: ProjectMemory,
) -> Result<PathBuf> {
let store = self.clone();
task::spawn_blocking(move || store.save_memory(&project_dir, &memory))
.await
.map_err(|err| anyhow::anyhow!("save_memory_async join error: {err}"))?
}
pub fn load_memory(&self, project_dir: &Path) -> Option<ProjectMemory> {
let hash = project_hash(project_dir);
let path = self.data_dir.join("memory").join(format!("{hash}.json"));
let content = fs::read_to_string(&path).ok()?;
match serde_json::from_str(&content) {
Ok(memory) => Some(memory),
Err(err) => {
tracing::warn!(
path = %path.display(),
error = %err,
"failed to parse project memory file"
);
None
}
}
}
pub async fn load_memory_async(&self, project_dir: PathBuf) -> Option<ProjectMemory> {
let store = self.clone();
task::spawn_blocking(move || store.load_memory(&project_dir))
.await
.ok()
.flatten()
}
pub fn add_memory_entry(
&self,
project_dir: &Path,
session_id: &SessionId,
summary: String,
) -> Result<PathBuf> {
let hash = project_hash(project_dir);
let path = self.data_dir.join("memory").join(format!("{hash}.json"));
for _attempt in 0..3 {
let mtime_before = fs::metadata(&path).ok().and_then(|m| m.modified().ok());
let mut memory = self.load_memory(project_dir).unwrap_or(ProjectMemory {
project_hash: project_hash(project_dir),
entries: Vec::new(),
});
memory.entries.push(crate::session::MemoryEntry {
created_at: current_unix_timestamp(),
summary: summary.clone(),
session_id: session_id.as_str().to_string(),
});
let data = serde_json::to_vec_pretty(&memory)?;
let mtime_after = fs::metadata(&path).ok().and_then(|m| m.modified().ok());
if mtime_before != mtime_after && mtime_before.is_some() {
std::thread::sleep(std::time::Duration::from_millis(50));
continue;
}
if let Some(parent) = path.parent() {
if !parent.exists() {
fs::create_dir_all(parent)?;
}
}
let tmp = path.with_extension("json.tmp");
fs::write(&tmp, data)?;
fs::rename(&tmp, &path)?;
crate::fs_util::set_private_file_permissions(&path)?;
return Ok(path);
}
anyhow::bail!("failed to add memory entry after 3 retries (concurrent write conflict)");
}
pub async fn add_memory_entry_async(
&self,
project_dir: PathBuf,
session_id: String,
summary: String,
) -> Result<PathBuf> {
let store = self.clone();
task::spawn_blocking(move || {
let sid = SessionId::new(session_id);
store.add_memory_entry(&project_dir, &sid, summary)
})
.await
.map_err(|err| anyhow::anyhow!("add_memory_entry_async join error: {err}"))?
}
}
pub fn current_unix_timestamp() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_secs())
.unwrap_or_default()
}
fn current_unix_millis() -> u128 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|duration| duration.as_millis())
.unwrap_or_default()
}
use crate::goal::types::SessionGoal;
use crate::model::ContentPart;
pub struct Submission {
pub task: String,
pub content_parts: Vec<ContentPart>,
pub response_tx: tokio::sync::oneshot::Sender<Result<String>>,
}
pub enum SessionCommand {
Turn(Submission),
TruncateToUserTurns {
keep_user_turns: usize,
response_tx: tokio::sync::oneshot::Sender<Result<usize>>,
},
Compact {
response_tx: tokio::sync::oneshot::Sender<Result<crate::compact::CompactOutcome>>,
},
}
pub fn truncate_messages_to_user_turns(
messages: &mut Vec<crate::model::ModelMessage>,
keep_user_turns: usize,
) {
use crate::model::ModelRole;
let mut seen_users = 0usize;
let mut cut: Option<usize> = None;
for (i, msg) in messages.iter().enumerate() {
if msg.role == ModelRole::User {
if seen_users == keep_user_turns {
cut = Some(i);
break;
}
seen_users += 1;
}
}
if let Some(i) = cut {
messages.truncate(i);
}
}
#[derive(Clone)]
pub struct SessionRuntime {
pub submission_tx: tokio::sync::mpsc::UnboundedSender<SessionCommand>,
}
impl SessionRuntime {
pub fn spawn(
ctx: std::sync::Arc<crate::turn::TurnContext>,
policy: crate::harness::HarnessPolicy,
initial_messages: Vec<crate::model::ModelMessage>,
_memory_injection: Option<String>,
) -> Self {
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<SessionCommand>();
tokio::spawn(async move {
let mut messages = initial_messages;
while let Some(command) = rx.recv().await {
match command {
SessionCommand::Turn(submission) => {
if submission.content_parts.is_empty() {
messages.push(crate::model::ModelMessage::user(submission.task));
} else {
messages.push(crate::model::ModelMessage::user_multimodal(
submission.task,
submission.content_parts,
));
}
let res = crate::turn::run_turn(&ctx, &mut messages, policy).await;
let _ = submission.response_tx.send(res);
}
SessionCommand::TruncateToUserTurns {
keep_user_turns,
response_tx,
} => {
truncate_messages_to_user_turns(&mut messages, keep_user_turns);
let _ = response_tx.send(Ok(messages.len()));
}
SessionCommand::Compact { response_tx } => {
let result = force_compact_session_messages(&ctx, &mut messages).await;
let _ = response_tx.send(result);
}
}
}
});
Self { submission_tx: tx }
}
}
async fn force_compact_session_messages(
ctx: &std::sync::Arc<crate::turn::TurnContext>,
messages: &mut Vec<crate::model::ModelMessage>,
) -> Result<crate::compact::CompactOutcome> {
if let Some(ref tx) = ctx.event_tx {
let _ = tx.send(crate::event::AgentEvent::AutoCompactStarted);
}
let provider = ctx.active_model_provider();
let model = ctx.active_model_name();
let mut state = ctx.compact_state.lock().await;
match ctx
.components
.compaction
.force_compact(
&mut state,
messages,
provider.as_ref(),
&model,
&ctx.harness_config,
)
.await
{
Ok(Some(outcome)) => {
if let Some(ref tx) = ctx.event_tx {
let _ = tx.send(crate::event::AgentEvent::AutoCompactCompleted {
tokens_saved: outcome.tokens_saved,
summary: outcome.summary.clone(),
kept_recent_messages: outcome.kept_recent_messages,
});
}
Ok(outcome)
}
Ok(None) => {
let reason = "nothing to compact".to_string();
if let Some(ref tx) = ctx.event_tx {
let _ = tx.send(crate::event::AgentEvent::AutoCompactFailed {
reason: reason.clone(),
});
}
Err(anyhow::anyhow!(reason))
}
Err(e) => {
if let Some(ref tx) = ctx.event_tx {
let _ = tx.send(crate::event::AgentEvent::AutoCompactFailed {
reason: e.to_string(),
});
}
Err(e)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::tool::{ToolInvocation, ToolResult};
#[test]
fn save_writes_session_snapshot() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("test-session".to_string()),
title: Some("Test session".to_string()),
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: Vec::new(),
memory: None,
goal: None,
usage: None,
};
let path = store.save(&snapshot).expect("save session");
assert!(path.exists());
assert_eq!(path.file_name().unwrap(), "test-session.json");
assert!(!path.with_extension("json.tmp").exists());
}
#[test]
fn save_replaces_existing_snapshot_atomically() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
let id = SessionId::new("atomic-session".to_string());
let mut snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: id.clone(),
title: Some("v1".to_string()),
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: vec![AgentEvent::UserTaskSubmitted {
text: "first".to_string(),
content_parts: vec![],
submitted_at: Some(1),
}],
memory: None,
goal: None,
usage: None,
};
store.save(&snapshot).expect("save v1");
snapshot.title = Some("v2".to_string());
snapshot.updated_at = 3;
snapshot.events.push(AgentEvent::ModelOutput {
text: "reply".to_string(),
thinking: None,
});
store.save(&snapshot).expect("save v2");
let loaded = store.load("atomic-session").expect("load");
assert_eq!(loaded.title.as_deref(), Some("v2"));
assert_eq!(loaded.events.len(), 2);
assert!(!store.root().join("atomic-session.json.tmp").exists());
}
#[cfg(unix)]
#[test]
fn save_restricts_session_file_and_directory_permissions() {
use std::os::unix::fs::PermissionsExt;
let tempdir = tempfile::tempdir().expect("tempdir");
let data_dir = tempdir.path().join("navi-data");
let store = SessionStore::new(data_dir);
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("private-session".to_string()),
title: None,
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: Vec::new(),
memory: None,
goal: None,
usage: None,
};
let path = store.save(&snapshot).expect("save session");
let dir_mode = fs::metadata(store.root())
.expect("dir metadata")
.permissions()
.mode()
& 0o777;
let file_mode = fs::metadata(path)
.expect("file metadata")
.permissions()
.mode()
& 0o777;
assert_eq!(dir_mode, 0o700);
assert_eq!(file_mode, 0o600);
}
#[test]
fn save_redacts_secret_like_event_content() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("redacted-session".to_string()),
title: None,
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: vec![AgentEvent::UserTaskSubmitted {
text: "OPENAI_API_KEY=sk-proj-1234567890abcdef".to_string(),
content_parts: vec![],
submitted_at: None,
}],
memory: None,
goal: None,
usage: None,
};
let path = store.save(&snapshot).expect("save session");
let content = fs::read_to_string(path).expect("read session");
assert!(content.contains("OPENAI_API_KEY=<redacted>"));
assert!(!content.contains("sk-proj-1234567890abcdef"));
}
#[test]
fn save_redacts_secret_like_memory_summaries() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("redacted-memory-session".to_string()),
title: None,
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: Vec::new(),
memory: Some(ProjectMemory {
project_hash: "abc".to_string(),
entries: vec![MemoryEntry {
created_at: 1_700_000_000,
summary: "Configured with OPENAI_API_KEY=sk-proj-abcdef0123456789".to_string(),
session_id: "session-x".to_string(),
}],
}),
goal: None,
usage: None,
};
let path = store.save(&snapshot).expect("save session");
let content = fs::read_to_string(path).expect("read session");
assert!(content.contains("OPENAI_API_KEY=<redacted>"));
assert!(!content.contains("sk-proj-abcdef0123456789"));
}
#[test]
fn save_can_preserve_event_content_when_redaction_is_disabled() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("unredacted-session".to_string()),
title: None,
project: PathBuf::from("/tmp/project"),
created_at: 1,
updated_at: 2,
events: vec![AgentEvent::UserTaskSubmitted {
text: "OPENAI_API_KEY=sk-proj-1234567890abcdef".to_string(),
content_parts: vec![],
submitted_at: None,
}],
memory: None,
goal: None,
usage: None,
};
let path = store.save(&snapshot).expect("save session");
let content = fs::read_to_string(path).expect("read session");
assert!(content.contains("sk-proj-1234567890abcdef"));
}
struct MockProvider;
#[async_trait::async_trait]
impl crate::model::ModelProvider for MockProvider {
fn stream(&self, _request: crate::model::ModelRequest) -> crate::model::ModelStream {
Box::pin(futures_util::stream::iter(vec![
Ok(crate::model::ModelStreamEvent::TextDelta {
text: "mock task response".to_string(),
}),
Ok(crate::model::ModelStreamEvent::Done),
]))
}
}
#[tokio::test]
async fn test_session_runtime_background_loop() {
let tempdir = tempfile::tempdir().unwrap();
let security_policy = crate::SecurityPolicy::new(
tempdir.path().to_path_buf(),
tempdir.path().to_path_buf(),
crate::SecurityConfig::default(),
)
.unwrap();
let tool_executor = std::sync::Arc::new(crate::ToolExecutor::new(security_policy));
let ctx = std::sync::Arc::new(crate::turn::TurnContext {
model_provider: std::sync::Arc::new(std::sync::RwLock::new(std::sync::Arc::new(
MockProvider,
))),
tool_executor,
project_dir: tempdir.path().to_path_buf(),
data_dir: tempdir.path().join("data"),
model_name: std::sync::Arc::new(std::sync::RwLock::new("test-model".to_string())),
event_tx: None,
approval_resolver: crate::runtime::ApprovalResolver::new_for_test(),
question_resolver: crate::runtime::QuestionResolver::new_for_test(),
plan_review_resolver: crate::runtime::PlanReviewResolver::new_for_test(),
sudo_password_resolver: crate::runtime::SudoPasswordResolver::new_for_test(),
compact_state: std::sync::Arc::new(tokio::sync::Mutex::new(
crate::compact::CompactState::new(128_000),
)),
harness_config: crate::config::HarnessConfig::default(),
include_tool_prompt_manifest: false,
context_packets: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
available_skills: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
skill_pools: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
active_skills: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
prompt_cache: std::sync::Arc::new(crate::prompt::PromptCache::new()),
instructions: std::sync::Arc::new(std::sync::RwLock::new(None)),
prompt_prefix: std::sync::Arc::new(std::sync::Mutex::new(None)),
components: crate::RuntimeComponents::default(),
cancel_token: crate::cancel::CancelToken::new(),
config: std::sync::Arc::new(std::sync::RwLock::new(
crate::config::NaviConfig::default(),
)),
memory_injection: None,
compaction_provider: None,
agent_mode: crate::plan_mode::AgentMode::Default,
compaction_model_name: None,
session_id: "test-session".to_string(),
allowed_tool_names: None,
is_subagent: false,
memory_manager: std::sync::Arc::new(std::sync::Mutex::new(None)),
harness_card: None,
});
let policy = crate::harness::policy_for_profile(
&crate::config::HarnessConfig {
observation_bytes_small: 1000,
..crate::config::HarnessConfig::default()
},
crate::config::HarnessProfile::Small,
);
let runtime = SessionRuntime::spawn(ctx, policy, Vec::new(), None);
let (tx, rx) = tokio::sync::oneshot::channel();
let submission = SessionCommand::Turn(Submission {
task: "hello world".to_string(),
content_parts: Vec::new(),
response_tx: tx,
});
runtime.submission_tx.send(submission).unwrap();
let result = rx.await.unwrap().unwrap();
assert_eq!(result, "mock task response");
}
#[test]
fn truncate_messages_keeps_preamble_and_prior_turns() {
use crate::model::{ModelMessage, ModelRole};
let mut messages = vec![
ModelMessage::system("sys"),
ModelMessage::developer("dev"),
ModelMessage::user("u1"),
ModelMessage {
role: ModelRole::Assistant,
content: "a1".into(),
content_parts: vec![],
tool_call_id: None,
tool_name: None,
tool_calls: vec![],
created_at: None,
thinking_content: None,
},
ModelMessage::user("u2"),
ModelMessage {
role: ModelRole::Assistant,
content: "a2".into(),
content_parts: vec![],
tool_call_id: None,
tool_name: None,
tool_calls: vec![],
created_at: None,
thinking_content: None,
},
ModelMessage::user("u3"),
];
truncate_messages_to_user_turns(&mut messages, 1);
assert_eq!(messages.len(), 4);
assert_eq!(messages[2].content, "u1");
assert_eq!(messages[3].content, "a1");
truncate_messages_to_user_turns(&mut messages, 0);
assert_eq!(messages.len(), 2);
assert!(matches!(messages[0].role, ModelRole::System));
assert!(matches!(messages[1].role, ModelRole::Developer));
}
#[test]
fn project_memory_format_injection_returns_none_when_empty() {
let memory = ProjectMemory {
project_hash: "abc".to_string(),
entries: Vec::new(),
};
assert!(memory.format_injection(3).is_none());
}
#[test]
fn project_memory_format_injection_returns_latest_entries() {
let memory = ProjectMemory {
project_hash: "abc".to_string(),
entries: vec![
MemoryEntry {
created_at: 1000,
summary: "First session".to_string(),
session_id: "session-1".to_string(),
},
MemoryEntry {
created_at: 2000,
summary: "Second session".to_string(),
session_id: "session-2".to_string(),
},
MemoryEntry {
created_at: 3000,
summary: "Third session".to_string(),
session_id: "session-3".to_string(),
},
MemoryEntry {
created_at: 4000,
summary: "Fourth session".to_string(),
session_id: "session-4".to_string(),
},
],
};
let injection = memory.format_injection(2).unwrap();
assert!(injection.contains("Third session"));
assert!(injection.contains("Fourth session"));
assert!(!injection.contains("First session"));
assert!(!injection.contains("Second session"));
}
#[test]
fn save_and_load_memory_roundtrip() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let project_dir = PathBuf::from("/tmp/test-project");
let memory = ProjectMemory {
project_hash: project_hash(&project_dir),
entries: vec![MemoryEntry {
created_at: 12345,
summary: "Worked on auth module".to_string(),
session_id: "session-test".to_string(),
}],
};
store
.save_memory(&project_dir, &memory)
.expect("save memory");
let loaded = store.load_memory(&project_dir).expect("load memory");
assert_eq!(loaded.entries.len(), 1);
assert_eq!(loaded.entries[0].summary, "Worked on auth module");
}
#[test]
fn add_memory_entry_appends_to_existing() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let project_dir = PathBuf::from("/tmp/test-project-2");
let session_id = SessionId::new("session-1".to_string());
store
.add_memory_entry(&project_dir, &session_id, "First summary".to_string())
.expect("add entry 1");
let session_id2 = SessionId::new("session-2".to_string());
store
.add_memory_entry(&project_dir, &session_id2, "Second summary".to_string())
.expect("add entry 2");
let loaded = store.load_memory(&project_dir).expect("load memory");
assert_eq!(loaded.entries.len(), 2);
assert_eq!(loaded.entries[0].summary, "First summary");
assert_eq!(loaded.entries[1].summary, "Second summary");
}
fn make_snapshot(id: &str, updated_at: u64) -> SessionSnapshot {
SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new(id.to_string()),
title: Some(format!("Session {id}")),
project: PathBuf::from("/tmp/project"),
created_at: updated_at - 10,
updated_at,
events: Vec::new(),
memory: None,
goal: None,
usage: None,
}
}
#[test]
fn list_returns_sessions_sorted_by_updated_at() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
store.save(&make_snapshot("s-old", 100)).expect("save");
store.save(&make_snapshot("s-new", 300)).expect("save");
store.save(&make_snapshot("s-mid", 200)).expect("save");
let sessions = store.list();
assert_eq!(sessions.len(), 3);
assert_eq!(sessions[0].id.as_str(), "s-new");
assert_eq!(sessions[1].id.as_str(), "s-mid");
assert_eq!(sessions[2].id.as_str(), "s-old");
}
#[test]
fn list_returns_empty_when_no_sessions() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
assert!(store.list().is_empty());
}
#[test]
fn load_roundtrip_save_then_load() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
let snapshot = make_snapshot("roundtrip-1", 500);
store.save(&snapshot).expect("save");
let loaded = store.load("roundtrip-1").expect("load");
assert_eq!(loaded.id.as_str(), "roundtrip-1");
assert_eq!(loaded.title, Some("Session roundtrip-1".to_string()));
assert_eq!(loaded.updated_at, 500);
}
#[test]
fn load_rejects_unsupported_version() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
let mut snapshot = make_snapshot("future-session", 100);
snapshot.version = 999;
store.save(&snapshot).expect("save");
let result = store.load("future-session");
assert!(result.is_err());
let err = result.unwrap_err().to_string();
assert!(err.contains("version"), "expected version error: {err}");
}
#[test]
fn delete_removes_session_file() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
store.save(&make_snapshot("del-1", 100)).expect("save");
assert!(store.root().join("del-1.json").exists());
let deleted = store.delete("del-1").expect("delete");
assert!(deleted);
assert!(!store.root().join("del-1.json").exists());
}
#[test]
fn delete_returns_false_for_missing() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let deleted = store.delete("nonexistent").expect("delete");
assert!(!deleted);
}
#[test]
fn session_snapshot_serialization_roundtrip() {
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("ser-1".to_string()),
title: Some("Test".to_string()),
project: PathBuf::from("/tmp/p"),
created_at: 1000,
updated_at: 2000,
events: vec![
AgentEvent::UserTaskSubmitted {
text: "hello".to_string(),
content_parts: vec![],
submitted_at: None,
},
AgentEvent::ModelOutput {
text: "response".to_string(),
thinking: Some("reasoning".to_string()),
},
],
memory: None,
goal: None,
usage: None,
};
let json = serde_json::to_string(&snapshot).expect("serialize");
let loaded: SessionSnapshot = serde_json::from_str(&json).expect("deserialize");
assert_eq!(loaded.id.as_str(), "ser-1");
assert_eq!(loaded.events.len(), 2);
}
#[test]
fn session_title_from_events_prefers_model_heading() {
let events = vec![
AgentEvent::UserTaskSubmitted {
text: "do something".to_string(),
content_parts: vec![],
submitted_at: None,
},
AgentEvent::ModelOutput {
text: "# My Analysis\n\nSome content here".to_string(),
thinking: None,
},
];
let title = session_title_from_events(&events);
assert_eq!(title.as_deref(), Some("My Analysis"));
}
#[test]
fn session_title_from_events_falls_back_to_user_text() {
let events = vec![AgentEvent::UserTaskSubmitted {
text: "Fix the bug".to_string(),
content_parts: vec![],
submitted_at: None,
}];
let title = session_title_from_events(&events);
assert_eq!(title.as_deref(), Some("Fix the bug"));
}
#[test]
fn clean_session_title_strips_markdown_and_truncates() {
assert_eq!(clean_session_title("## Short"), Some("Short".to_string()));
assert_eq!(
clean_session_title("`code snippet`"),
Some("code snippet".to_string())
);
let long = "a".repeat(200);
let result = clean_session_title(&long).unwrap();
assert!(result.len() <= 80);
}
#[test]
fn clean_session_title_returns_none_for_empty() {
assert!(clean_session_title("").is_none());
assert!(clean_session_title("###").is_none());
}
#[test]
fn save_and_load_preserves_events() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
let snapshot = SessionSnapshot {
version: SessionSnapshot::CURRENT_VERSION,
id: SessionId::new("events-session".to_string()),
title: None,
project: PathBuf::from("/tmp/p"),
created_at: 10,
updated_at: 20,
events: vec![
AgentEvent::UserTaskSubmitted {
text: "task".to_string(),
content_parts: vec![],
submitted_at: None,
},
AgentEvent::ToolRequested(ToolInvocation {
id: "c1".to_string(),
tool_name: "read_file".to_string(),
input: serde_json::json!({"path": "x.txt"}),
}),
AgentEvent::ToolCompleted(ToolResult {
invocation_id: "c1".to_string(),
ok: true,
output: serde_json::json!("file content"),
}),
],
memory: None,
goal: None,
usage: None,
};
store.save(&snapshot).expect("save");
let loaded = store.load("events-session").expect("load");
assert_eq!(loaded.events.len(), 3);
}
#[test]
fn regression_corrupt_json_on_disk_skipped_by_list() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
store.save(&make_snapshot("valid", 100)).expect("save");
let corrupt_path = store.root().join("corrupt.json");
std::fs::write(&corrupt_path, "{invalid json!!!").expect("write corrupt");
let sessions = store.list();
assert_eq!(sessions.len(), 1);
assert_eq!(sessions[0].id.as_str(), "valid");
}
#[test]
fn regression_list_ignores_non_json_files() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
store.save(&make_snapshot("valid", 100)).expect("save");
std::fs::write(store.root().join("notes.txt"), "not a session").expect("write");
std::fs::write(store.root().join("README.md"), "# readme").expect("write");
let sessions = store.list();
assert_eq!(sessions.len(), 1);
}
#[test]
fn regression_load_missing_version_defaults_to_one() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let json = serde_json::json!({
"id": "no-version",
"title": null,
"project": "/tmp/p",
"created_at": 1,
"updated_at": 2,
"events": [],
"memory": null
});
let path = store.root().join("no-version.json");
std::fs::create_dir_all(store.root()).expect("create sessions dir");
std::fs::write(&path, serde_json::to_string(&json).unwrap()).expect("write");
let loaded = store.load("no-version").expect("load");
assert_eq!(loaded.version, 1); }
#[test]
fn regression_load_memory_malformed_json_returns_none() {
let tempdir = tempfile::tempdir().expect("tempdir");
let store = SessionStore::new(tempdir.path().to_path_buf());
let project_dir = PathBuf::from("/tmp/test-project");
let hash = {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
project_dir.hash(&mut hasher);
format!("{:016x}", hasher.finish())
};
let memory_dir = tempdir.path().join("memory");
std::fs::create_dir_all(&memory_dir).expect("create");
std::fs::write(memory_dir.join(format!("{hash}.json")), "not json!").expect("write");
let loaded = store.load_memory(&project_dir);
assert!(loaded.is_none(), "malformed memory should return None");
}
#[test]
fn regression_session_title_only_tool_events_returns_none() {
let events = vec![
AgentEvent::ToolRequested(ToolInvocation {
id: "c1".to_string(),
tool_name: "read_file".to_string(),
input: serde_json::json!({}),
}),
AgentEvent::ToolCompleted(ToolResult {
invocation_id: "c1".to_string(),
ok: true,
output: serde_json::json!("content"),
}),
];
let title = session_title_from_events(&events);
assert!(title.is_none(), "no user/model text should return None");
}
#[test]
fn regression_project_hash_is_stable() {
let path = PathBuf::from("/tmp/some/project/dir");
let hash1 = {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
path.hash(&mut hasher);
format!("{:016x}", hasher.finish())
};
let hash2 = {
use std::hash::{Hash, Hasher};
let mut hasher = std::collections::hash_map::DefaultHasher::new();
path.hash(&mut hasher);
format!("{:016x}", hasher.finish())
};
assert_eq!(hash1, hash2);
}
#[test]
fn regression_create_id_format() {
let id = SessionStore::create_id();
assert!(
id.as_str().starts_with("session-"),
"session id must start with 'session-'"
);
}
#[test]
fn create_id_is_unique_for_concurrent_agent_sessions() {
const WORKERS: usize = 8;
const IDS_PER_WORKER: usize = 128;
let ids = (0..WORKERS)
.map(|_| {
std::thread::spawn(|| {
(0..IDS_PER_WORKER)
.map(|_| SessionStore::create_id().into_inner())
.collect::<Vec<_>>()
})
})
.flat_map(|worker| worker.join().expect("session-id worker should not panic"))
.collect::<std::collections::HashSet<_>>();
assert_eq!(ids.len(), WORKERS * IDS_PER_WORKER);
}
}