Skip to main content

navi_core/
session.rs

1use crate::event::AgentEvent;
2use crate::security::{redact_memory, redact_snapshot_events};
3use anyhow::{Context, Result};
4use serde::{Deserialize, Serialize};
5use std::fs;
6use std::io::Read;
7use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicU64, Ordering};
9use std::time::{SystemTime, UNIX_EPOCH};
10use tokio::task;
11
12/// Per-process sequence used alongside time, PID, and entropy when allocating
13/// a session id. A timestamp by itself collides when two NAVI processes start
14/// in the same millisecond.
15static NEXT_SESSION_SEQUENCE: AtomicU64 = AtomicU64::new(1);
16
17/// Unique identifier for a session, wrapping a string id like
18/// `"session-1719612345000-1234-1-a1b2c3d4e5f60708"`.
19#[derive(Debug, Clone, Serialize, Deserialize)]
20pub struct SessionId(String);
21
22impl SessionId {
23    /// Creates a new `SessionId` from the given string.
24    pub fn new(id: String) -> Self {
25        Self(id)
26    }
27
28    /// Returns the id as a string slice.
29    pub fn as_str(&self) -> &str {
30        &self.0
31    }
32
33    /// Consumes the id and returns the inner string.
34    pub fn into_inner(self) -> String {
35        self.0
36    }
37}
38
39/// Accumulated session memory for a project, used to inject past context into new sessions.
40#[derive(Debug, Clone, Serialize, Deserialize)]
41pub struct ProjectMemory {
42    /// Hash identifying the project directory.
43    pub project_hash: String,
44    /// Ordered memory entries from past sessions.
45    pub entries: Vec<MemoryEntry>,
46}
47
48/// A single memory entry from a completed session.
49#[derive(Debug, Clone, Serialize, Deserialize)]
50pub struct MemoryEntry {
51    /// Unix timestamp (seconds) when this entry was created.
52    pub created_at: u64,
53    /// Short summary of what the session covered.
54    pub summary: String,
55    /// Identifier of the originating session.
56    pub session_id: String,
57}
58
59/// Derives a short title from session events by looking for a markdown heading
60/// in the first model output, falling back to a cleaned version of the first
61/// user task text.
62///
63/// Returns `None` if no suitable text is found.
64pub fn session_title_from_events(events: &[AgentEvent]) -> Option<String> {
65    events
66        .iter()
67        .find_map(|event| match event {
68            AgentEvent::ModelOutput { text, .. } => title_from_model_text(text),
69            _ => None,
70        })
71        .or_else(|| {
72            events.iter().find_map(|event| match event {
73                AgentEvent::UserTaskSubmitted { text, .. } => title_from_user_text(text),
74                _ => None,
75            })
76        })
77}
78
79fn title_from_model_text(text: &str) -> Option<String> {
80    let heading = text.lines().find_map(|line| {
81        let trimmed = line.trim();
82        if trimmed.starts_with('#') {
83            Some(trimmed.trim_start_matches('#').trim())
84        } else {
85            None
86        }
87    });
88
89    heading
90        .and_then(clean_session_title)
91        .or_else(|| text.lines().find_map(clean_session_title))
92}
93
94fn title_from_user_text(text: &str) -> Option<String> {
95    clean_session_title(text)
96}
97
98/// Sanitizes text into a short session title by trimming whitespace, quotes,
99/// markdown markers, and truncating to 80 characters.
100///
101/// Returns `None` if the cleaned text is empty.
102pub fn clean_session_title(text: &str) -> Option<String> {
103    let cleaned = text
104        .trim()
105        .trim_matches('`')
106        .trim_matches('"')
107        .trim_matches('\'')
108        .trim_start_matches(['#', '-', '*', '>'])
109        .split_whitespace()
110        .collect::<Vec<_>>()
111        .join(" ");
112
113    if cleaned.is_empty() {
114        return None;
115    }
116
117    Some(
118        cleaned
119            .chars()
120            .take(80)
121            .collect::<String>()
122            .trim()
123            .to_string(),
124    )
125}
126
127impl ProjectMemory {
128    /// Returns at most `max` of the most recent memory entries.
129    pub fn recent_entries(&self, max: usize) -> &[MemoryEntry] {
130        let start = self.entries.len().saturating_sub(max);
131        &self.entries[start..]
132    }
133
134    /// Formats up to `max` recent entries into a text block suitable for
135    /// injection into the system prompt. Returns `None` if there are no entries.
136    pub fn format_injection(&self, max: usize) -> Option<String> {
137        let entries = self.recent_entries(max);
138        if entries.is_empty() {
139            return None;
140        }
141        let mut parts = Vec::new();
142        for entry in entries {
143            parts.push(format!(
144                "[Session {} — {}]\n{}",
145                entry.session_id,
146                format_timestamp(entry.created_at),
147                entry.summary
148            ));
149        }
150        Some(format!(
151            "Previous session context (summarized):\n\n{}",
152            parts.join("\n\n")
153        ))
154    }
155}
156
157fn format_timestamp(unix_secs: u64) -> String {
158    let days = unix_secs / 86400;
159    let hours = (unix_secs % 86400) / 3600;
160    let minutes = (unix_secs % 3600) / 60;
161    format!("day {days} {hours:02}:{minutes:02}")
162}
163
164fn project_hash(project_dir: &Path) -> String {
165    use std::hash::{Hash, Hasher};
166    let mut hasher = std::collections::hash_map::DefaultHasher::new();
167    project_dir.hash(&mut hasher);
168    format!("{:016x}", hasher.finish())
169}
170
171/// Persists [`SessionSnapshot`] JSON files to disk under `<data_dir>/sessions/`.
172///
173/// By default, secret redaction is enabled so API keys and tokens are scrubbed
174/// from saved event text.
175#[derive(Debug, Clone)]
176pub struct SessionStore {
177    root: PathBuf,
178    data_dir: PathBuf,
179    redact_secrets: bool,
180}
181
182fn default_session_version() -> u32 {
183    1
184}
185
186/// Accumulated token/cost usage for a session (persisted with the snapshot).
187#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
188pub struct SessionUsageSnapshot {
189    /// Cumulative prompt/context tokens billed this session.
190    #[serde(default)]
191    pub input_tokens: u64,
192    /// Cumulative completion tokens billed this session.
193    #[serde(default)]
194    pub output_tokens: u64,
195    /// Estimated spend in USD from list rates × tokens (when known).
196    #[serde(default)]
197    pub cost_usd: f64,
198    /// True once at least one turn had usable list pricing.
199    #[serde(default)]
200    pub cost_known: bool,
201    /// Estimated prepaid credits spent (e.g. Hypercredits = USD / $0.05).
202    #[serde(default, skip_serializing_if = "Option::is_none")]
203    pub credits_spent: Option<f64>,
204    /// Credit unit label when `credits_spent` is set (e.g. `hypercredits`).
205    #[serde(default, skip_serializing_if = "Option::is_none")]
206    pub credit_unit: Option<String>,
207}
208
209/// A serializable snapshot of a complete session, persisted to disk as JSON.
210#[derive(Debug, Clone, Serialize, Deserialize)]
211pub struct SessionSnapshot {
212    /// Snapshot schema version; currently `1`.
213    #[serde(default = "default_session_version")]
214    pub version: u32,
215    /// Unique session identifier.
216    pub id: SessionId,
217    /// Short human-readable title, derived from the first user/assistant message.
218    #[serde(default)]
219    pub title: Option<String>,
220    /// Project directory this session belongs to.
221    pub project: PathBuf,
222    /// Unix timestamp (seconds) when the session was created.
223    #[serde(default)]
224    pub created_at: u64,
225    /// Unix timestamp (seconds) when the session was last updated.
226    #[serde(default)]
227    pub updated_at: u64,
228    /// All agent events recorded during the session.
229    pub events: Vec<AgentEvent>,
230    /// Optional project memory snapshot co-persisted with the session.
231    #[serde(default)]
232    pub memory: Option<ProjectMemory>,
233    /// Optional session goal co-persisted with the session.
234    #[serde(default)]
235    pub goal: Option<SessionGoal>,
236    /// Token and estimated cost usage for this session (restored on reload).
237    #[serde(default, skip_serializing_if = "Option::is_none")]
238    pub usage: Option<SessionUsageSnapshot>,
239}
240
241/// Lightweight metadata for listing saved sessions without loading event history.
242#[derive(Debug, Clone, Serialize, Deserialize)]
243pub struct SessionSnapshotInfo {
244    /// Unique session identifier.
245    pub id: SessionId,
246    /// Short human-readable title, derived when the snapshot was saved.
247    #[serde(default)]
248    pub title: Option<String>,
249    /// Project directory this session belongs to.
250    pub project: PathBuf,
251    /// Unix timestamp (seconds) when the session was created.
252    #[serde(default)]
253    pub created_at: u64,
254    /// Unix timestamp (seconds) when the session was last updated.
255    #[serde(default)]
256    pub updated_at: u64,
257}
258
259impl SessionSnapshot {
260    /// Current snapshot schema version.
261    pub const CURRENT_VERSION: u32 = 1;
262}
263
264/// True only for persisted session snapshots (`{session_id}.json`).
265///
266/// Rejects desktop UI paint sidecars (`{session_id}.ui.json`) and other
267/// auxiliary `*.json` files that share the sessions directory. Without this,
268/// UIs list ghost sessions that fail to load.
269fn is_session_snapshot_file(path: &Path) -> bool {
270    let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
271        return false;
272    };
273    if !name.ends_with(".json") || name.ends_with(".ui.json") {
274        return false;
275    }
276    // e.g. foo.json.tmp written mid-save
277    if name.contains(".tmp") || name.ends_with(".bak") {
278        return false;
279    }
280    path.extension().and_then(|e| e.to_str()) == Some("json")
281}
282
283fn read_session_info(path: &Path) -> Result<SessionSnapshotInfo> {
284    const METADATA_READ_LIMIT: usize = 64 * 1024;
285
286    let mut file =
287        fs::File::open(path).with_context(|| format!("failed to open {}", path.display()))?;
288    let mut buffer = vec![0; METADATA_READ_LIMIT];
289    let bytes_read = file
290        .read(&mut buffer)
291        .with_context(|| format!("failed to read {}", path.display()))?;
292    buffer.truncate(bytes_read);
293
294    let prefix = std::str::from_utf8(&buffer)
295        .with_context(|| format!("failed to decode metadata prefix from {}", path.display()))?;
296
297    if let Some(events_index) = prefix.find("\"events\"")
298        && let Some(comma_index) = prefix[..events_index].rfind(',')
299    {
300        let metadata_json = format!("{}\n}}", &prefix[..comma_index]);
301        return serde_json::from_str::<SessionSnapshotInfo>(&metadata_json)
302            .with_context(|| format!("failed to parse metadata from {}", path.display()));
303    }
304
305    let content =
306        fs::read_to_string(path).with_context(|| format!("failed to read {}", path.display()))?;
307    serde_json::from_str::<SessionSnapshotInfo>(&content)
308        .with_context(|| format!("failed to parse metadata from {}", path.display()))
309}
310
311impl SessionStore {
312    /// Creates a new store with secret redaction enabled.
313    pub fn new(data_dir: PathBuf) -> Self {
314        Self::with_redaction(data_dir, true)
315    }
316
317    /// Creates a new store with the given redaction setting.
318    pub fn with_redaction(data_dir: PathBuf, redact_secrets: bool) -> Self {
319        Self {
320            root: data_dir.join("sessions"),
321            data_dir,
322            redact_secrets,
323        }
324    }
325
326    /// Returns the directory where session JSON files are stored.
327    pub fn root(&self) -> &PathBuf {
328        &self.root
329    }
330
331    /// Generates a session id that is unique across concurrent NAVI instances.
332    ///
333    /// The timestamp keeps ids naturally ordered, while the process id,
334    /// in-process sequence, and random nonce prevent two agents started in the
335    /// same millisecond from sharing a provider cache-affinity identity.
336    pub fn create_id() -> SessionId {
337        let millis = current_unix_millis();
338        let pid = std::process::id();
339        let sequence = NEXT_SESSION_SEQUENCE.fetch_add(1, Ordering::Relaxed);
340        let nonce = fastrand::u64(..);
341        SessionId::new(format!("session-{millis}-{pid}-{sequence}-{nonce:016x}"))
342    }
343
344    /// Serializes and saves a snapshot to disk, creating the sessions directory
345    /// if needed. Applies secret redaction unless disabled.
346    ///
347    /// This is the blocking implementation used internally and in tests. Use
348    /// [`Self::save_async`] from async contexts to avoid blocking the Tokio
349    /// runtime.
350    pub fn save(&self, snapshot: &SessionSnapshot) -> Result<PathBuf> {
351        fs::create_dir_all(&self.root)
352            .with_context(|| format!("failed to create {}", self.root.display()))?;
353        crate::fs_util::set_private_dir_permissions(&self.root)?;
354
355        let path = self.root.join(format!("{}.json", snapshot.id.as_str()));
356        let snapshot = if self.redact_secrets {
357            SessionSnapshot {
358                version: snapshot.version,
359                id: snapshot.id.clone(),
360                title: snapshot.title.clone(),
361                project: snapshot.project.clone(),
362                created_at: snapshot.created_at,
363                updated_at: snapshot.updated_at,
364                goal: snapshot.goal.clone(),
365                events: redact_snapshot_events(&snapshot.events),
366                memory: snapshot.memory.as_ref().map(redact_memory),
367                usage: snapshot.usage.clone(),
368            }
369        } else {
370            snapshot.clone()
371        };
372        let data = serde_json::to_vec_pretty(&snapshot)?;
373        // Atomic replace: write temp then rename so a crash mid-save cannot leave
374        // a truncated/corrupt session JSON (readers only see the previous file
375        // or the complete new one).
376        let tmp = path.with_extension("json.tmp");
377        fs::write(&tmp, &data).with_context(|| format!("failed to write {}", tmp.display()))?;
378        crate::fs_util::set_private_file_permissions(&tmp)?;
379        fs::rename(&tmp, &path).with_context(|| {
380            format!(
381                "failed to replace {} with {}",
382                path.display(),
383                tmp.display()
384            )
385        })?;
386        // Re-assert final path perms (rename may inherit dir default ACLs).
387        crate::fs_util::set_private_file_permissions(&path)?;
388
389        Ok(path)
390    }
391
392    /// Async wrapper around [`Self::save`] that runs the blocking filesystem
393    /// operations on the Tokio blocking thread pool.
394    pub async fn save_async(&self, snapshot: SessionSnapshot) -> Result<PathBuf> {
395        let store = self.clone();
396        task::spawn_blocking(move || store.save(&snapshot))
397            .await
398            .map_err(|err| anyhow::anyhow!("save_async join error: {err}"))?
399    }
400
401    /// Loads all saved sessions from disk, sorted by most recently updated first.
402    pub fn list(&self) -> Vec<SessionSnapshot> {
403        let mut sessions = Vec::new();
404        if let Ok(entries) = fs::read_dir(&self.root) {
405            for entry in entries.flatten() {
406                let path = entry.path();
407                // Only real session snapshots (`{id}.json`). Desktop may write
408                // `{id}.ui.json` paint caches next to them — those must not list.
409                if !is_session_snapshot_file(&path) {
410                    continue;
411                }
412                if let Ok(content) = fs::read_to_string(&path)
413                    && let Ok(snapshot) = serde_json::from_str::<SessionSnapshot>(&content)
414                {
415                    sessions.push(snapshot);
416                }
417            }
418        }
419        sessions.sort_by(|a, b| {
420            b.updated_at
421                .cmp(&a.updated_at)
422                .then_with(|| b.id.as_str().cmp(a.id.as_str()))
423        });
424        sessions
425    }
426
427    /// Loads only session metadata from disk, sorted by most recently updated first.
428    pub fn list_info(&self) -> Vec<SessionSnapshotInfo> {
429        let mut sessions = Vec::new();
430        if let Ok(entries) = fs::read_dir(&self.root) {
431            for entry in entries.flatten() {
432                let path = entry.path();
433                if !is_session_snapshot_file(&path) {
434                    continue;
435                }
436                if let Ok(info) = read_session_info(&path) {
437                    sessions.push(info);
438                }
439            }
440        }
441        sessions.sort_by(|a, b| {
442            b.updated_at
443                .cmp(&a.updated_at)
444                .then_with(|| b.id.as_str().cmp(a.id.as_str()))
445        });
446        sessions
447    }
448
449    /// Async wrapper around [`Self::list`] that runs the blocking filesystem
450    /// operations on the Tokio blocking thread pool.
451    pub async fn list_async(&self) -> Vec<SessionSnapshot> {
452        let store = self.clone();
453        task::spawn_blocking(move || store.list())
454            .await
455            .unwrap_or_default()
456    }
457
458    /// Async wrapper around [`Self::list_info`] that avoids blocking the async runtime.
459    pub async fn list_info_async(&self) -> Vec<SessionSnapshotInfo> {
460        let store = self.clone();
461        task::spawn_blocking(move || store.list_info())
462            .await
463            .unwrap_or_default()
464    }
465
466    /// Loads a single session by id. Returns an error if the file is missing or
467    /// the snapshot version is newer than supported.
468    pub fn load(&self, session_id: &str) -> Result<SessionSnapshot> {
469        let path = self.root.join(format!("{session_id}.json"));
470        let content = fs::read_to_string(&path)
471            .with_context(|| format!("failed to read {}", path.display()))?;
472        let snapshot: SessionSnapshot = serde_json::from_str(&content)
473            .with_context(|| format!("failed to parse {}", path.display()))?;
474        if snapshot.version > SessionSnapshot::CURRENT_VERSION {
475            return Err(anyhow::anyhow!(
476                "session snapshot version {} is newer than supported version {}",
477                snapshot.version,
478                SessionSnapshot::CURRENT_VERSION
479            ));
480        }
481        Ok(snapshot)
482    }
483
484    /// Async wrapper around [`Self::load`] that runs the blocking filesystem
485    /// operations on the Tokio blocking thread pool.
486    pub async fn load_async(&self, session_id: String) -> Result<SessionSnapshot> {
487        let store = self.clone();
488        task::spawn_blocking(move || store.load(&session_id))
489            .await
490            .map_err(|err| anyhow::anyhow!("load_async join error: {err}"))?
491    }
492
493    /// Deletes the session file. Returns `true` if the file existed and was removed.
494    pub fn delete(&self, session_id: &str) -> Result<bool> {
495        let path = self.root.join(format!("{session_id}.json"));
496        match fs::remove_file(&path) {
497            Ok(()) => Ok(true),
498            Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(false),
499            Err(err) => Err(err).with_context(|| format!("failed to delete {}", path.display())),
500        }
501    }
502
503    /// Renames a saved session by updating its title field in the snapshot.
504    /// Returns `true` if the session existed and was updated.
505    pub fn rename(&self, session_id: &str, title: &str) -> Result<bool> {
506        let title = title.trim();
507        if title.is_empty() {
508            return Err(anyhow::anyhow!("session title cannot be empty"));
509        }
510        let path = self.root.join(format!("{session_id}.json"));
511        if !path.exists() {
512            return Ok(false);
513        }
514        let mut snapshot = self.load(session_id)?;
515        snapshot.title = Some(title.to_string());
516        snapshot.updated_at = current_unix_timestamp();
517        self.save(&snapshot)?;
518        Ok(true)
519    }
520
521    /// Async wrapper around [`Self::rename`].
522    pub async fn rename_async(&self, session_id: String, title: String) -> Result<bool> {
523        let store = self.clone();
524        task::spawn_blocking(move || store.rename(&session_id, &title))
525            .await
526            .map_err(|err| anyhow::anyhow!("rename_async join error: {err}"))?
527    }
528
529    /// Async wrapper around [`Self::delete`] that runs the blocking filesystem
530    /// operations on the Tokio blocking thread pool.
531    pub async fn delete_async(&self, session_id: String) -> Result<bool> {
532        let store = self.clone();
533        task::spawn_blocking(move || store.delete(&session_id))
534            .await
535            .map_err(|err| anyhow::anyhow!("delete_async join error: {err}"))?
536    }
537
538    /// Persists project memory to `<data_dir>/memory/<hash>.json`.
539    pub fn save_memory(&self, project_dir: &Path, memory: &ProjectMemory) -> Result<PathBuf> {
540        let memory_dir = self.data_dir.join("memory");
541        fs::create_dir_all(&memory_dir)
542            .with_context(|| format!("failed to create {}", memory_dir.display()))?;
543        crate::fs_util::set_private_dir_permissions(&memory_dir)?;
544
545        let hash = project_hash(project_dir);
546        let path = memory_dir.join(format!("{hash}.json"));
547        let data = serde_json::to_vec_pretty(memory)?;
548        fs::write(&path, data).with_context(|| format!("failed to write {}", path.display()))?;
549        crate::fs_util::set_private_file_permissions(&path)?;
550
551        Ok(path)
552    }
553
554    /// Async wrapper around [`Self::save_memory`] that runs the blocking
555    /// filesystem operations on the Tokio blocking thread pool.
556    pub async fn save_memory_async(
557        &self,
558        project_dir: PathBuf,
559        memory: ProjectMemory,
560    ) -> Result<PathBuf> {
561        let store = self.clone();
562        task::spawn_blocking(move || store.save_memory(&project_dir, &memory))
563            .await
564            .map_err(|err| anyhow::anyhow!("save_memory_async join error: {err}"))?
565    }
566
567    /// Loads project memory from disk, returning `None` if no memory file exists.
568    pub fn load_memory(&self, project_dir: &Path) -> Option<ProjectMemory> {
569        let hash = project_hash(project_dir);
570        let path = self.data_dir.join("memory").join(format!("{hash}.json"));
571        let content = fs::read_to_string(&path).ok()?;
572        match serde_json::from_str(&content) {
573            Ok(memory) => Some(memory),
574            Err(err) => {
575                tracing::warn!(
576                    path = %path.display(),
577                    error = %err,
578                    "failed to parse project memory file"
579                );
580                None
581            }
582        }
583    }
584
585    /// Async wrapper around [`Self::load_memory`] that runs the blocking
586    /// filesystem operations on the Tokio blocking thread pool.
587    pub async fn load_memory_async(&self, project_dir: PathBuf) -> Option<ProjectMemory> {
588        let store = self.clone();
589        task::spawn_blocking(move || store.load_memory(&project_dir))
590            .await
591            .ok()
592            .flatten()
593    }
594
595    /// Appends a new memory entry for the project and persists it to disk.
596    pub fn add_memory_entry(
597        &self,
598        project_dir: &Path,
599        session_id: &SessionId,
600        summary: String,
601    ) -> Result<PathBuf> {
602        let hash = project_hash(project_dir);
603        let path = self.data_dir.join("memory").join(format!("{hash}.json"));
604
605        // Retry loop to handle concurrent writes from other NAVI instances.
606        // If the file changes between load and save (detected via mtime), we
607        // reload and retry instead of overwriting and losing entries.
608        for _attempt in 0..3 {
609            let mtime_before = fs::metadata(&path).ok().and_then(|m| m.modified().ok());
610            let mut memory = self.load_memory(project_dir).unwrap_or(ProjectMemory {
611                project_hash: project_hash(project_dir),
612                entries: Vec::new(),
613            });
614            memory.entries.push(crate::session::MemoryEntry {
615                created_at: current_unix_timestamp(),
616                summary: summary.clone(),
617                session_id: session_id.as_str().to_string(),
618            });
619            let data = serde_json::to_vec_pretty(&memory)?;
620
621            // Check if the file changed while we were preparing
622            let mtime_after = fs::metadata(&path).ok().and_then(|m| m.modified().ok());
623            if mtime_before != mtime_after && mtime_before.is_some() {
624                // File was modified by another process — retry
625                std::thread::sleep(std::time::Duration::from_millis(50));
626                continue;
627            }
628
629            // Atomic write via temp file + rename
630            if let Some(parent) = path.parent() {
631                if !parent.exists() {
632                    fs::create_dir_all(parent)?;
633                }
634            }
635            let tmp = path.with_extension("json.tmp");
636            fs::write(&tmp, data)?;
637            fs::rename(&tmp, &path)?;
638            crate::fs_util::set_private_file_permissions(&path)?;
639            return Ok(path);
640        }
641        anyhow::bail!("failed to add memory entry after 3 retries (concurrent write conflict)");
642    }
643
644    /// Async wrapper around [`Self::add_memory_entry`] that runs the blocking
645    /// filesystem operations on the Tokio blocking thread pool.
646    pub async fn add_memory_entry_async(
647        &self,
648        project_dir: PathBuf,
649        session_id: String,
650        summary: String,
651    ) -> Result<PathBuf> {
652        let store = self.clone();
653        task::spawn_blocking(move || {
654            let sid = SessionId::new(session_id);
655            store.add_memory_entry(&project_dir, &sid, summary)
656        })
657        .await
658        .map_err(|err| anyhow::anyhow!("add_memory_entry_async join error: {err}"))?
659    }
660}
661
662/// Returns the current time as a Unix timestamp in seconds.
663pub fn current_unix_timestamp() -> u64 {
664    SystemTime::now()
665        .duration_since(UNIX_EPOCH)
666        .map(|duration| duration.as_secs())
667        .unwrap_or_default()
668}
669
670fn current_unix_millis() -> u128 {
671    SystemTime::now()
672        .duration_since(UNIX_EPOCH)
673        .map(|duration| duration.as_millis())
674        .unwrap_or_default()
675}
676
677use crate::goal::types::SessionGoal;
678use crate::model::ContentPart;
679
680/// A user task submission sent to the session background loop.
681pub struct Submission {
682    /// The user's task text.
683    pub task: String,
684    /// Optional multimodal content parts (images + text).
685    /// When non-empty, the session loop creates a multimodal user message.
686    pub content_parts: Vec<ContentPart>,
687    /// Channel to send the assistant's response back to the caller.
688    pub response_tx: tokio::sync::oneshot::Sender<Result<String>>,
689}
690
691/// Commands accepted by the session background loop.
692pub enum SessionCommand {
693    /// Run a full agent turn for a user message.
694    Turn(Submission),
695    /// Drop conversation history after `keep_user_turns` user messages
696    /// (and the assistant/tool messages belonging to those turns).
697    /// Used when the UI edits a past user message and re-sends.
698    TruncateToUserTurns {
699        keep_user_turns: usize,
700        response_tx: tokio::sync::oneshot::Sender<Result<usize>>,
701    },
702    /// Force-compact live conversation history with the session model.
703    Compact {
704        response_tx: tokio::sync::oneshot::Sender<Result<crate::compact::CompactOutcome>>,
705    },
706}
707
708/// Truncate model history so only the first `keep_user_turns` user turns remain.
709///
710/// System/developer preamble is always kept. The cut point is the start of the
711/// `(keep_user_turns + 1)`-th user message (0-based count of user messages kept).
712pub fn truncate_messages_to_user_turns(
713    messages: &mut Vec<crate::model::ModelMessage>,
714    keep_user_turns: usize,
715) {
716    use crate::model::ModelRole;
717    let mut seen_users = 0usize;
718    let mut cut: Option<usize> = None;
719    for (i, msg) in messages.iter().enumerate() {
720        if msg.role == ModelRole::User {
721            if seen_users == keep_user_turns {
722                cut = Some(i);
723                break;
724            }
725            seen_users += 1;
726        }
727    }
728    if let Some(i) = cut {
729        messages.truncate(i);
730    }
731}
732
733/// A handle to a background session loop that accepts [`SessionCommand`]s and
734/// runs them through the turn pipeline.
735#[derive(Clone)]
736pub struct SessionRuntime {
737    /// Channel for sending commands to the background loop.
738    pub submission_tx: tokio::sync::mpsc::UnboundedSender<SessionCommand>,
739}
740
741impl SessionRuntime {
742    /// Spawns a background tokio task that processes submissions sequentially
743    /// through the turn pipeline, maintaining conversation history.
744    pub fn spawn(
745        ctx: std::sync::Arc<crate::turn::TurnContext>,
746        policy: crate::harness::HarnessPolicy,
747        initial_messages: Vec<crate::model::ModelMessage>,
748        _memory_injection: Option<String>,
749    ) -> Self {
750        let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<SessionCommand>();
751
752        tokio::spawn(async move {
753            let mut messages = initial_messages;
754
755            while let Some(command) = rx.recv().await {
756                match command {
757                    SessionCommand::Turn(submission) => {
758                        if submission.content_parts.is_empty() {
759                            messages.push(crate::model::ModelMessage::user(submission.task));
760                        } else {
761                            messages.push(crate::model::ModelMessage::user_multimodal(
762                                submission.task,
763                                submission.content_parts,
764                            ));
765                        }
766                        let res = crate::turn::run_turn(&ctx, &mut messages, policy).await;
767                        let _ = submission.response_tx.send(res);
768                    }
769                    SessionCommand::TruncateToUserTurns {
770                        keep_user_turns,
771                        response_tx,
772                    } => {
773                        truncate_messages_to_user_turns(&mut messages, keep_user_turns);
774                        let _ = response_tx.send(Ok(messages.len()));
775                    }
776                    SessionCommand::Compact { response_tx } => {
777                        let result = force_compact_session_messages(&ctx, &mut messages).await;
778                        let _ = response_tx.send(result);
779                    }
780                }
781            }
782        });
783
784        Self { submission_tx: tx }
785    }
786}
787
788/// Force-compact the live session message list using the active session model.
789async fn force_compact_session_messages(
790    ctx: &std::sync::Arc<crate::turn::TurnContext>,
791    messages: &mut Vec<crate::model::ModelMessage>,
792) -> Result<crate::compact::CompactOutcome> {
793    if let Some(ref tx) = ctx.event_tx {
794        let _ = tx.send(crate::event::AgentEvent::AutoCompactStarted);
795    }
796
797    let provider = ctx.active_model_provider();
798    let model = ctx.active_model_name();
799    let mut state = ctx.compact_state.lock().await;
800    match ctx
801        .components
802        .compaction
803        .force_compact(
804            &mut state,
805            messages,
806            provider.as_ref(),
807            &model,
808            &ctx.harness_config,
809        )
810        .await
811    {
812        Ok(Some(outcome)) => {
813            if let Some(ref tx) = ctx.event_tx {
814                let _ = tx.send(crate::event::AgentEvent::AutoCompactCompleted {
815                    tokens_saved: outcome.tokens_saved,
816                    summary: outcome.summary.clone(),
817                    kept_recent_messages: outcome.kept_recent_messages,
818                });
819            }
820            Ok(outcome)
821        }
822        Ok(None) => {
823            let reason = "nothing to compact".to_string();
824            if let Some(ref tx) = ctx.event_tx {
825                let _ = tx.send(crate::event::AgentEvent::AutoCompactFailed {
826                    reason: reason.clone(),
827                });
828            }
829            Err(anyhow::anyhow!(reason))
830        }
831        Err(e) => {
832            if let Some(ref tx) = ctx.event_tx {
833                let _ = tx.send(crate::event::AgentEvent::AutoCompactFailed {
834                    reason: e.to_string(),
835                });
836            }
837            Err(e)
838        }
839    }
840}
841
842#[cfg(test)]
843mod tests {
844    use super::*;
845    use crate::tool::{ToolInvocation, ToolResult};
846
847    #[test]
848    fn save_writes_session_snapshot() {
849        let tempdir = tempfile::tempdir().expect("tempdir");
850        let store = SessionStore::new(tempdir.path().to_path_buf());
851        let snapshot = SessionSnapshot {
852            version: SessionSnapshot::CURRENT_VERSION,
853            id: SessionId::new("test-session".to_string()),
854            title: Some("Test session".to_string()),
855            project: PathBuf::from("/tmp/project"),
856            created_at: 1,
857            updated_at: 2,
858            events: Vec::new(),
859            memory: None,
860            goal: None,
861            usage: None,
862        };
863
864        let path = store.save(&snapshot).expect("save session");
865        assert!(path.exists());
866        assert_eq!(path.file_name().unwrap(), "test-session.json");
867        // Temp file must not linger after a successful atomic replace.
868        assert!(!path.with_extension("json.tmp").exists());
869    }
870
871    #[test]
872    fn save_replaces_existing_snapshot_atomically() {
873        let tempdir = tempfile::tempdir().expect("tempdir");
874        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
875        let id = SessionId::new("atomic-session".to_string());
876        let mut snapshot = SessionSnapshot {
877            version: SessionSnapshot::CURRENT_VERSION,
878            id: id.clone(),
879            title: Some("v1".to_string()),
880            project: PathBuf::from("/tmp/project"),
881            created_at: 1,
882            updated_at: 2,
883            events: vec![AgentEvent::UserTaskSubmitted {
884                text: "first".to_string(),
885                content_parts: vec![],
886                submitted_at: Some(1),
887            }],
888            memory: None,
889            goal: None,
890            usage: None,
891        };
892        store.save(&snapshot).expect("save v1");
893        snapshot.title = Some("v2".to_string());
894        snapshot.updated_at = 3;
895        snapshot.events.push(AgentEvent::ModelOutput {
896            text: "reply".to_string(),
897            thinking: None,
898        });
899        store.save(&snapshot).expect("save v2");
900
901        let loaded = store.load("atomic-session").expect("load");
902        assert_eq!(loaded.title.as_deref(), Some("v2"));
903        assert_eq!(loaded.events.len(), 2);
904        assert!(!store.root().join("atomic-session.json.tmp").exists());
905    }
906
907    #[cfg(unix)]
908    #[test]
909    fn save_restricts_session_file_and_directory_permissions() {
910        use std::os::unix::fs::PermissionsExt;
911
912        let tempdir = tempfile::tempdir().expect("tempdir");
913        let data_dir = tempdir.path().join("navi-data");
914        let store = SessionStore::new(data_dir);
915        let snapshot = SessionSnapshot {
916            version: SessionSnapshot::CURRENT_VERSION,
917            id: SessionId::new("private-session".to_string()),
918            title: None,
919            project: PathBuf::from("/tmp/project"),
920            created_at: 1,
921            updated_at: 2,
922            events: Vec::new(),
923            memory: None,
924            goal: None,
925            usage: None,
926        };
927
928        let path = store.save(&snapshot).expect("save session");
929        let dir_mode = fs::metadata(store.root())
930            .expect("dir metadata")
931            .permissions()
932            .mode()
933            & 0o777;
934        let file_mode = fs::metadata(path)
935            .expect("file metadata")
936            .permissions()
937            .mode()
938            & 0o777;
939
940        assert_eq!(dir_mode, 0o700);
941        assert_eq!(file_mode, 0o600);
942    }
943
944    #[test]
945    fn save_redacts_secret_like_event_content() {
946        let tempdir = tempfile::tempdir().expect("tempdir");
947        let store = SessionStore::new(tempdir.path().to_path_buf());
948        let snapshot = SessionSnapshot {
949            version: SessionSnapshot::CURRENT_VERSION,
950            id: SessionId::new("redacted-session".to_string()),
951            title: None,
952            project: PathBuf::from("/tmp/project"),
953            created_at: 1,
954            updated_at: 2,
955            events: vec![AgentEvent::UserTaskSubmitted {
956                text: "OPENAI_API_KEY=sk-proj-1234567890abcdef".to_string(),
957                content_parts: vec![],
958                submitted_at: None,
959            }],
960            memory: None,
961            goal: None,
962            usage: None,
963        };
964
965        let path = store.save(&snapshot).expect("save session");
966        let content = fs::read_to_string(path).expect("read session");
967
968        assert!(content.contains("OPENAI_API_KEY=<redacted>"));
969        assert!(!content.contains("sk-proj-1234567890abcdef"));
970    }
971
972    #[test]
973    fn save_redacts_secret_like_memory_summaries() {
974        let tempdir = tempfile::tempdir().expect("tempdir");
975        let store = SessionStore::new(tempdir.path().to_path_buf());
976        let snapshot = SessionSnapshot {
977            version: SessionSnapshot::CURRENT_VERSION,
978            id: SessionId::new("redacted-memory-session".to_string()),
979            title: None,
980            project: PathBuf::from("/tmp/project"),
981            created_at: 1,
982            updated_at: 2,
983            events: Vec::new(),
984            memory: Some(ProjectMemory {
985                project_hash: "abc".to_string(),
986                entries: vec![MemoryEntry {
987                    created_at: 1_700_000_000,
988                    summary: "Configured with OPENAI_API_KEY=sk-proj-abcdef0123456789".to_string(),
989                    session_id: "session-x".to_string(),
990                }],
991            }),
992            goal: None,
993            usage: None,
994        };
995
996        let path = store.save(&snapshot).expect("save session");
997        let content = fs::read_to_string(path).expect("read session");
998
999        assert!(content.contains("OPENAI_API_KEY=<redacted>"));
1000        assert!(!content.contains("sk-proj-abcdef0123456789"));
1001    }
1002
1003    #[test]
1004    fn save_can_preserve_event_content_when_redaction_is_disabled() {
1005        let tempdir = tempfile::tempdir().expect("tempdir");
1006        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1007        let snapshot = SessionSnapshot {
1008            version: SessionSnapshot::CURRENT_VERSION,
1009            id: SessionId::new("unredacted-session".to_string()),
1010            title: None,
1011            project: PathBuf::from("/tmp/project"),
1012            created_at: 1,
1013            updated_at: 2,
1014            events: vec![AgentEvent::UserTaskSubmitted {
1015                text: "OPENAI_API_KEY=sk-proj-1234567890abcdef".to_string(),
1016                content_parts: vec![],
1017                submitted_at: None,
1018            }],
1019            memory: None,
1020            goal: None,
1021            usage: None,
1022        };
1023
1024        let path = store.save(&snapshot).expect("save session");
1025        let content = fs::read_to_string(path).expect("read session");
1026
1027        assert!(content.contains("sk-proj-1234567890abcdef"));
1028    }
1029
1030    struct MockProvider;
1031
1032    #[async_trait::async_trait]
1033    impl crate::model::ModelProvider for MockProvider {
1034        fn stream(&self, _request: crate::model::ModelRequest) -> crate::model::ModelStream {
1035            Box::pin(futures_util::stream::iter(vec![
1036                Ok(crate::model::ModelStreamEvent::TextDelta {
1037                    text: "mock task response".to_string(),
1038                }),
1039                Ok(crate::model::ModelStreamEvent::Done),
1040            ]))
1041        }
1042    }
1043
1044    #[tokio::test]
1045    async fn test_session_runtime_background_loop() {
1046        let tempdir = tempfile::tempdir().unwrap();
1047        let security_policy = crate::SecurityPolicy::new(
1048            tempdir.path().to_path_buf(),
1049            tempdir.path().to_path_buf(),
1050            crate::SecurityConfig::default(),
1051        )
1052        .unwrap();
1053        let tool_executor = std::sync::Arc::new(crate::ToolExecutor::new(security_policy));
1054
1055        let ctx = std::sync::Arc::new(crate::turn::TurnContext {
1056            model_provider: std::sync::Arc::new(std::sync::RwLock::new(std::sync::Arc::new(
1057                MockProvider,
1058            ))),
1059            tool_executor,
1060            project_dir: tempdir.path().to_path_buf(),
1061            data_dir: tempdir.path().join("data"),
1062            model_name: std::sync::Arc::new(std::sync::RwLock::new("test-model".to_string())),
1063            event_tx: None,
1064            approval_resolver: crate::runtime::ApprovalResolver::new_for_test(),
1065            question_resolver: crate::runtime::QuestionResolver::new_for_test(),
1066            plan_review_resolver: crate::runtime::PlanReviewResolver::new_for_test(),
1067            sudo_password_resolver: crate::runtime::SudoPasswordResolver::new_for_test(),
1068            compact_state: std::sync::Arc::new(tokio::sync::Mutex::new(
1069                crate::compact::CompactState::new(128_000),
1070            )),
1071            harness_config: crate::config::HarnessConfig::default(),
1072            include_tool_prompt_manifest: false,
1073            context_packets: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
1074            available_skills: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
1075            skill_pools: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
1076            active_skills: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
1077            prompt_cache: std::sync::Arc::new(crate::prompt::PromptCache::new()),
1078            instructions: std::sync::Arc::new(std::sync::RwLock::new(None)),
1079            prompt_prefix: std::sync::Arc::new(std::sync::Mutex::new(None)),
1080            components: crate::RuntimeComponents::default(),
1081            cancel_token: crate::cancel::CancelToken::new(),
1082            config: std::sync::Arc::new(std::sync::RwLock::new(
1083                crate::config::NaviConfig::default(),
1084            )),
1085            memory_injection: None,
1086            compaction_provider: None,
1087            agent_mode: crate::plan_mode::AgentMode::Default,
1088            compaction_model_name: None,
1089            session_id: "test-session".to_string(),
1090            allowed_tool_names: None,
1091            is_subagent: false,
1092            memory_manager: std::sync::Arc::new(std::sync::Mutex::new(None)),
1093            harness_card: None,
1094        });
1095
1096        let policy = crate::harness::policy_for_profile(
1097            &crate::config::HarnessConfig {
1098                observation_bytes_small: 1000,
1099                ..crate::config::HarnessConfig::default()
1100            },
1101            crate::config::HarnessProfile::Small,
1102        );
1103
1104        let runtime = SessionRuntime::spawn(ctx, policy, Vec::new(), None);
1105
1106        let (tx, rx) = tokio::sync::oneshot::channel();
1107        let submission = SessionCommand::Turn(Submission {
1108            task: "hello world".to_string(),
1109            content_parts: Vec::new(),
1110            response_tx: tx,
1111        });
1112
1113        runtime.submission_tx.send(submission).unwrap();
1114
1115        let result = rx.await.unwrap().unwrap();
1116        assert_eq!(result, "mock task response");
1117    }
1118
1119    #[test]
1120    fn truncate_messages_keeps_preamble_and_prior_turns() {
1121        use crate::model::{ModelMessage, ModelRole};
1122        let mut messages = vec![
1123            ModelMessage::system("sys"),
1124            ModelMessage::developer("dev"),
1125            ModelMessage::user("u1"),
1126            ModelMessage {
1127                role: ModelRole::Assistant,
1128                content: "a1".into(),
1129                content_parts: vec![],
1130                tool_call_id: None,
1131                tool_name: None,
1132                tool_calls: vec![],
1133                created_at: None,
1134                thinking_content: None,
1135            },
1136            ModelMessage::user("u2"),
1137            ModelMessage {
1138                role: ModelRole::Assistant,
1139                content: "a2".into(),
1140                content_parts: vec![],
1141                tool_call_id: None,
1142                tool_name: None,
1143                tool_calls: vec![],
1144                created_at: None,
1145                thinking_content: None,
1146            },
1147            ModelMessage::user("u3"),
1148        ];
1149        truncate_messages_to_user_turns(&mut messages, 1);
1150        assert_eq!(messages.len(), 4);
1151        assert_eq!(messages[2].content, "u1");
1152        assert_eq!(messages[3].content, "a1");
1153
1154        truncate_messages_to_user_turns(&mut messages, 0);
1155        assert_eq!(messages.len(), 2);
1156        assert!(matches!(messages[0].role, ModelRole::System));
1157        assert!(matches!(messages[1].role, ModelRole::Developer));
1158    }
1159
1160    #[test]
1161    fn project_memory_format_injection_returns_none_when_empty() {
1162        let memory = ProjectMemory {
1163            project_hash: "abc".to_string(),
1164            entries: Vec::new(),
1165        };
1166        assert!(memory.format_injection(3).is_none());
1167    }
1168
1169    #[test]
1170    fn project_memory_format_injection_returns_latest_entries() {
1171        let memory = ProjectMemory {
1172            project_hash: "abc".to_string(),
1173            entries: vec![
1174                MemoryEntry {
1175                    created_at: 1000,
1176                    summary: "First session".to_string(),
1177                    session_id: "session-1".to_string(),
1178                },
1179                MemoryEntry {
1180                    created_at: 2000,
1181                    summary: "Second session".to_string(),
1182                    session_id: "session-2".to_string(),
1183                },
1184                MemoryEntry {
1185                    created_at: 3000,
1186                    summary: "Third session".to_string(),
1187                    session_id: "session-3".to_string(),
1188                },
1189                MemoryEntry {
1190                    created_at: 4000,
1191                    summary: "Fourth session".to_string(),
1192                    session_id: "session-4".to_string(),
1193                },
1194            ],
1195        };
1196        let injection = memory.format_injection(2).unwrap();
1197        assert!(injection.contains("Third session"));
1198        assert!(injection.contains("Fourth session"));
1199        assert!(!injection.contains("First session"));
1200        assert!(!injection.contains("Second session"));
1201    }
1202
1203    #[test]
1204    fn save_and_load_memory_roundtrip() {
1205        let tempdir = tempfile::tempdir().expect("tempdir");
1206        let store = SessionStore::new(tempdir.path().to_path_buf());
1207        let project_dir = PathBuf::from("/tmp/test-project");
1208
1209        let memory = ProjectMemory {
1210            project_hash: project_hash(&project_dir),
1211            entries: vec![MemoryEntry {
1212                created_at: 12345,
1213                summary: "Worked on auth module".to_string(),
1214                session_id: "session-test".to_string(),
1215            }],
1216        };
1217
1218        store
1219            .save_memory(&project_dir, &memory)
1220            .expect("save memory");
1221        let loaded = store.load_memory(&project_dir).expect("load memory");
1222        assert_eq!(loaded.entries.len(), 1);
1223        assert_eq!(loaded.entries[0].summary, "Worked on auth module");
1224    }
1225
1226    #[test]
1227    fn add_memory_entry_appends_to_existing() {
1228        let tempdir = tempfile::tempdir().expect("tempdir");
1229        let store = SessionStore::new(tempdir.path().to_path_buf());
1230        let project_dir = PathBuf::from("/tmp/test-project-2");
1231
1232        let session_id = SessionId::new("session-1".to_string());
1233        store
1234            .add_memory_entry(&project_dir, &session_id, "First summary".to_string())
1235            .expect("add entry 1");
1236
1237        let session_id2 = SessionId::new("session-2".to_string());
1238        store
1239            .add_memory_entry(&project_dir, &session_id2, "Second summary".to_string())
1240            .expect("add entry 2");
1241
1242        let loaded = store.load_memory(&project_dir).expect("load memory");
1243        assert_eq!(loaded.entries.len(), 2);
1244        assert_eq!(loaded.entries[0].summary, "First summary");
1245        assert_eq!(loaded.entries[1].summary, "Second summary");
1246    }
1247
1248    fn make_snapshot(id: &str, updated_at: u64) -> SessionSnapshot {
1249        SessionSnapshot {
1250            version: SessionSnapshot::CURRENT_VERSION,
1251            id: SessionId::new(id.to_string()),
1252            title: Some(format!("Session {id}")),
1253            project: PathBuf::from("/tmp/project"),
1254            created_at: updated_at - 10,
1255            updated_at,
1256            events: Vec::new(),
1257            memory: None,
1258            goal: None,
1259            usage: None,
1260        }
1261    }
1262
1263    #[test]
1264    fn list_returns_sessions_sorted_by_updated_at() {
1265        let tempdir = tempfile::tempdir().expect("tempdir");
1266        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1267        store.save(&make_snapshot("s-old", 100)).expect("save");
1268        store.save(&make_snapshot("s-new", 300)).expect("save");
1269        store.save(&make_snapshot("s-mid", 200)).expect("save");
1270
1271        let sessions = store.list();
1272        assert_eq!(sessions.len(), 3);
1273        assert_eq!(sessions[0].id.as_str(), "s-new");
1274        assert_eq!(sessions[1].id.as_str(), "s-mid");
1275        assert_eq!(sessions[2].id.as_str(), "s-old");
1276    }
1277
1278    #[test]
1279    fn list_returns_empty_when_no_sessions() {
1280        let tempdir = tempfile::tempdir().expect("tempdir");
1281        let store = SessionStore::new(tempdir.path().to_path_buf());
1282        assert!(store.list().is_empty());
1283    }
1284
1285    #[test]
1286    fn load_roundtrip_save_then_load() {
1287        let tempdir = tempfile::tempdir().expect("tempdir");
1288        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1289        let snapshot = make_snapshot("roundtrip-1", 500);
1290        store.save(&snapshot).expect("save");
1291
1292        let loaded = store.load("roundtrip-1").expect("load");
1293        assert_eq!(loaded.id.as_str(), "roundtrip-1");
1294        assert_eq!(loaded.title, Some("Session roundtrip-1".to_string()));
1295        assert_eq!(loaded.updated_at, 500);
1296    }
1297
1298    #[test]
1299    fn load_rejects_unsupported_version() {
1300        let tempdir = tempfile::tempdir().expect("tempdir");
1301        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1302        let mut snapshot = make_snapshot("future-session", 100);
1303        snapshot.version = 999;
1304        store.save(&snapshot).expect("save");
1305
1306        let result = store.load("future-session");
1307        assert!(result.is_err());
1308        let err = result.unwrap_err().to_string();
1309        assert!(err.contains("version"), "expected version error: {err}");
1310    }
1311
1312    #[test]
1313    fn delete_removes_session_file() {
1314        let tempdir = tempfile::tempdir().expect("tempdir");
1315        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1316        store.save(&make_snapshot("del-1", 100)).expect("save");
1317        assert!(store.root().join("del-1.json").exists());
1318
1319        let deleted = store.delete("del-1").expect("delete");
1320        assert!(deleted);
1321        assert!(!store.root().join("del-1.json").exists());
1322    }
1323
1324    #[test]
1325    fn delete_returns_false_for_missing() {
1326        let tempdir = tempfile::tempdir().expect("tempdir");
1327        let store = SessionStore::new(tempdir.path().to_path_buf());
1328        let deleted = store.delete("nonexistent").expect("delete");
1329        assert!(!deleted);
1330    }
1331
1332    #[test]
1333    fn session_snapshot_serialization_roundtrip() {
1334        let snapshot = SessionSnapshot {
1335            version: SessionSnapshot::CURRENT_VERSION,
1336            id: SessionId::new("ser-1".to_string()),
1337            title: Some("Test".to_string()),
1338            project: PathBuf::from("/tmp/p"),
1339            created_at: 1000,
1340            updated_at: 2000,
1341            events: vec![
1342                AgentEvent::UserTaskSubmitted {
1343                    text: "hello".to_string(),
1344                    content_parts: vec![],
1345                    submitted_at: None,
1346                },
1347                AgentEvent::ModelOutput {
1348                    text: "response".to_string(),
1349                    thinking: Some("reasoning".to_string()),
1350                },
1351            ],
1352            memory: None,
1353            goal: None,
1354            usage: None,
1355        };
1356        let json = serde_json::to_string(&snapshot).expect("serialize");
1357        let loaded: SessionSnapshot = serde_json::from_str(&json).expect("deserialize");
1358        assert_eq!(loaded.id.as_str(), "ser-1");
1359        assert_eq!(loaded.events.len(), 2);
1360    }
1361
1362    #[test]
1363    fn session_title_from_events_prefers_model_heading() {
1364        let events = vec![
1365            AgentEvent::UserTaskSubmitted {
1366                text: "do something".to_string(),
1367                content_parts: vec![],
1368                submitted_at: None,
1369            },
1370            AgentEvent::ModelOutput {
1371                text: "# My Analysis\n\nSome content here".to_string(),
1372                thinking: None,
1373            },
1374        ];
1375        let title = session_title_from_events(&events);
1376        assert_eq!(title.as_deref(), Some("My Analysis"));
1377    }
1378
1379    #[test]
1380    fn session_title_from_events_falls_back_to_user_text() {
1381        let events = vec![AgentEvent::UserTaskSubmitted {
1382            text: "Fix the bug".to_string(),
1383            content_parts: vec![],
1384            submitted_at: None,
1385        }];
1386        let title = session_title_from_events(&events);
1387        assert_eq!(title.as_deref(), Some("Fix the bug"));
1388    }
1389
1390    #[test]
1391    fn clean_session_title_strips_markdown_and_truncates() {
1392        assert_eq!(clean_session_title("## Short"), Some("Short".to_string()));
1393        assert_eq!(
1394            clean_session_title("`code snippet`"),
1395            Some("code snippet".to_string())
1396        );
1397        let long = "a".repeat(200);
1398        let result = clean_session_title(&long).unwrap();
1399        assert!(result.len() <= 80);
1400    }
1401
1402    #[test]
1403    fn clean_session_title_returns_none_for_empty() {
1404        assert!(clean_session_title("").is_none());
1405        assert!(clean_session_title("###").is_none());
1406    }
1407
1408    #[test]
1409    fn save_and_load_preserves_events() {
1410        let tempdir = tempfile::tempdir().expect("tempdir");
1411        let store = SessionStore::with_redaction(tempdir.path().to_path_buf(), false);
1412        let snapshot = SessionSnapshot {
1413            version: SessionSnapshot::CURRENT_VERSION,
1414            id: SessionId::new("events-session".to_string()),
1415            title: None,
1416            project: PathBuf::from("/tmp/p"),
1417            created_at: 10,
1418            updated_at: 20,
1419            events: vec![
1420                AgentEvent::UserTaskSubmitted {
1421                    text: "task".to_string(),
1422                    content_parts: vec![],
1423                    submitted_at: None,
1424                },
1425                AgentEvent::ToolRequested(ToolInvocation {
1426                    id: "c1".to_string(),
1427                    tool_name: "read_file".to_string(),
1428                    input: serde_json::json!({"path": "x.txt"}),
1429                }),
1430                AgentEvent::ToolCompleted(ToolResult {
1431                    invocation_id: "c1".to_string(),
1432                    ok: true,
1433                    output: serde_json::json!("file content"),
1434                }),
1435            ],
1436            memory: None,
1437            goal: None,
1438            usage: None,
1439        };
1440        store.save(&snapshot).expect("save");
1441        let loaded = store.load("events-session").expect("load");
1442        assert_eq!(loaded.events.len(), 3);
1443    }
1444
1445    // ── Regression tests ──────────────────────────────────────────────────────
1446
1447    #[test]
1448    fn regression_corrupt_json_on_disk_skipped_by_list() {
1449        let tempdir = tempfile::tempdir().expect("tempdir");
1450        let store = SessionStore::new(tempdir.path().to_path_buf());
1451
1452        // Write a valid session
1453        store.save(&make_snapshot("valid", 100)).expect("save");
1454
1455        // Write a corrupt JSON file
1456        let corrupt_path = store.root().join("corrupt.json");
1457        std::fs::write(&corrupt_path, "{invalid json!!!").expect("write corrupt");
1458
1459        // list() should skip the corrupt file and return only the valid one
1460        let sessions = store.list();
1461        assert_eq!(sessions.len(), 1);
1462        assert_eq!(sessions[0].id.as_str(), "valid");
1463    }
1464
1465    #[test]
1466    fn regression_list_ignores_non_json_files() {
1467        let tempdir = tempfile::tempdir().expect("tempdir");
1468        let store = SessionStore::new(tempdir.path().to_path_buf());
1469
1470        store.save(&make_snapshot("valid", 100)).expect("save");
1471
1472        // Write non-json files
1473        std::fs::write(store.root().join("notes.txt"), "not a session").expect("write");
1474        std::fs::write(store.root().join("README.md"), "# readme").expect("write");
1475
1476        let sessions = store.list();
1477        assert_eq!(sessions.len(), 1);
1478    }
1479
1480    #[test]
1481    fn regression_load_missing_version_defaults_to_one() {
1482        let tempdir = tempfile::tempdir().expect("tempdir");
1483        let store = SessionStore::new(tempdir.path().to_path_buf());
1484
1485        // Write a snapshot JSON without the "version" field
1486        // SessionId serializes as a plain string
1487        let json = serde_json::json!({
1488            "id": "no-version",
1489            "title": null,
1490            "project": "/tmp/p",
1491            "created_at": 1,
1492            "updated_at": 2,
1493            "events": [],
1494            "memory": null
1495        });
1496        let path = store.root().join("no-version.json");
1497        std::fs::create_dir_all(store.root()).expect("create sessions dir");
1498        std::fs::write(&path, serde_json::to_string(&json).unwrap()).expect("write");
1499
1500        let loaded = store.load("no-version").expect("load");
1501        assert_eq!(loaded.version, 1); // default_session_version
1502    }
1503
1504    #[test]
1505    fn regression_load_memory_malformed_json_returns_none() {
1506        let tempdir = tempfile::tempdir().expect("tempdir");
1507        let store = SessionStore::new(tempdir.path().to_path_buf());
1508        let project_dir = PathBuf::from("/tmp/test-project");
1509
1510        // Write a corrupt memory file
1511        let hash = {
1512            use std::hash::{Hash, Hasher};
1513            let mut hasher = std::collections::hash_map::DefaultHasher::new();
1514            project_dir.hash(&mut hasher);
1515            format!("{:016x}", hasher.finish())
1516        };
1517        let memory_dir = tempdir.path().join("memory");
1518        std::fs::create_dir_all(&memory_dir).expect("create");
1519        std::fs::write(memory_dir.join(format!("{hash}.json")), "not json!").expect("write");
1520
1521        let loaded = store.load_memory(&project_dir);
1522        assert!(loaded.is_none(), "malformed memory should return None");
1523    }
1524
1525    #[test]
1526    fn regression_session_title_only_tool_events_returns_none() {
1527        let events = vec![
1528            AgentEvent::ToolRequested(ToolInvocation {
1529                id: "c1".to_string(),
1530                tool_name: "read_file".to_string(),
1531                input: serde_json::json!({}),
1532            }),
1533            AgentEvent::ToolCompleted(ToolResult {
1534                invocation_id: "c1".to_string(),
1535                ok: true,
1536                output: serde_json::json!("content"),
1537            }),
1538        ];
1539        let title = session_title_from_events(&events);
1540        assert!(title.is_none(), "no user/model text should return None");
1541    }
1542
1543    #[test]
1544    fn regression_project_hash_is_stable() {
1545        let path = PathBuf::from("/tmp/some/project/dir");
1546        let hash1 = {
1547            use std::hash::{Hash, Hasher};
1548            let mut hasher = std::collections::hash_map::DefaultHasher::new();
1549            path.hash(&mut hasher);
1550            format!("{:016x}", hasher.finish())
1551        };
1552        let hash2 = {
1553            use std::hash::{Hash, Hasher};
1554            let mut hasher = std::collections::hash_map::DefaultHasher::new();
1555            path.hash(&mut hasher);
1556            format!("{:016x}", hasher.finish())
1557        };
1558        assert_eq!(hash1, hash2);
1559    }
1560
1561    #[test]
1562    fn regression_create_id_format() {
1563        let id = SessionStore::create_id();
1564        assert!(
1565            id.as_str().starts_with("session-"),
1566            "session id must start with 'session-'"
1567        );
1568    }
1569
1570    #[test]
1571    fn create_id_is_unique_for_concurrent_agent_sessions() {
1572        const WORKERS: usize = 8;
1573        const IDS_PER_WORKER: usize = 128;
1574
1575        let ids = (0..WORKERS)
1576            .map(|_| {
1577                std::thread::spawn(|| {
1578                    (0..IDS_PER_WORKER)
1579                        .map(|_| SessionStore::create_id().into_inner())
1580                        .collect::<Vec<_>>()
1581                })
1582            })
1583            .flat_map(|worker| worker.join().expect("session-id worker should not panic"))
1584            .collect::<std::collections::HashSet<_>>();
1585
1586        assert_eq!(ids.len(), WORKERS * IDS_PER_WORKER);
1587    }
1588}