Skip to main content

scv_tools/
conversation.rs

1//! Multi-turn conversations with delegated agents.
2//!
3//! A session's conversations live in one [`ConversationStore`] shared by its
4//! `agent_*` tools, and end with the session. The model sees only SCV-issued
5//! handles such as `codex-2`; the CLI's own session IDs stay here and are
6//! never accepted from the model. While a conversation exists, a marker file
7//! named after its session ID tells `scv agents gc` to keep its transcript.
8
9use std::{
10    any::Any,
11    collections::HashMap,
12    fmt,
13    path::{Path, PathBuf},
14    sync::{Arc, Mutex},
15    time::{Duration, Instant, SystemTime},
16};
17
18use scv_core::ToolError;
19use serde::{Deserialize, Serialize};
20
21use crate::{
22    adapters::ConversationFiles,
23    agent_output::valid_session_id,
24    delegation::{ProcessIdentity, write_private_json},
25};
26
27/// Transcripts younger than this are never removed, so a turn that is still
28/// running (and so still writing its transcript) is safe from `gc`.
29pub const MIN_GC_AGE: Duration = Duration::from_secs(3600);
30
31/// Longest handle accepted from the model.
32const MAX_HANDLE_BYTES: usize = 64;
33
34#[derive(Debug, Clone, Copy)]
35pub struct ConversationLimits {
36    /// Conversations remembered per session; starting another forgets the
37    /// least recently used idle one.
38    pub max: usize,
39    /// A conversation unused this long is forgotten.
40    pub idle: Duration,
41}
42
43/// What a live conversation keeps between turns, such as its running child
44/// process. Dropping the last reference (when the conversation is forgotten,
45/// expires, or its session ends) is what shuts the child down.
46#[derive(Clone)]
47pub(crate) struct Attachment(pub Arc<dyn Any + Send + Sync>);
48
49impl fmt::Debug for Attachment {
50    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
51        formatter.write_str("Attachment")
52    }
53}
54
55#[derive(Debug)]
56struct Conversation {
57    agent: String,
58    cwd: PathBuf,
59    /// The CLI's session ID, once known.
60    vendor: Option<String>,
61    turns: u32,
62    busy: bool,
63    last_used: Instant,
64    attachment: Option<Attachment>,
65}
66
67#[derive(Debug, Default)]
68struct Inner {
69    conversations: HashMap<String, Conversation>,
70    /// Next handle number per agent.
71    next: HashMap<String, u32>,
72}
73
74/// One session's conversations with its delegated agents.
75#[derive(Debug)]
76pub struct ConversationStore {
77    limits: ConversationLimits,
78    marker_dir: Option<PathBuf>,
79    owner: Option<ProcessIdentity>,
80    inner: Mutex<Inner>,
81}
82
83/// On-disk marker keeping a live conversation's transcript from `gc`.
84#[derive(Debug, Serialize, Deserialize)]
85struct Marker {
86    owner: ProcessIdentity,
87    agent: String,
88    handle: String,
89}
90
91/// A turn in progress. Finish it with [`TurnGuard::finish`]; dropping it
92/// instead (a cancelled or abandoned call) frees the conversation.
93#[derive(Debug)]
94pub(crate) struct TurnGuard {
95    store: Arc<ConversationStore>,
96    pub handle: String,
97    pub turn: u32,
98    /// Continuing: the known session ID. Starting: the ID SCV chose, if any.
99    pub vendor: Option<String>,
100    /// Continuing: what the conversation kept. Starting: what
101    /// [`TurnGuard::attach`] set, stored when the turn finishes.
102    attachment: Option<Attachment>,
103    finished: bool,
104}
105
106/// Whether `value` has the shape of an SCV conversation handle
107/// (`<agent>-<number>`). CLI session IDs never do.
108pub fn is_handle(value: &str) -> bool {
109    value.len() <= MAX_HANDLE_BYTES
110        && value.split_once('-').is_some_and(|(agent, number)| {
111            !agent.is_empty()
112                && agent
113                    .chars()
114                    .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
115                && !number.is_empty()
116                && number.chars().all(|c| c.is_ascii_digit())
117        })
118}
119
120impl ConversationStore {
121    /// `marker_dir` is `$SCV_HOME/state/conversations`; `None` keeps no markers.
122    pub fn new(limits: ConversationLimits, marker_dir: Option<PathBuf>) -> Self {
123        Self {
124            limits,
125            marker_dir,
126            owner: ProcessIdentity::current(),
127            inner: Mutex::new(Inner::default()),
128        }
129    }
130
131    /// Start a turn: a new conversation when `handle` is `None`, otherwise
132    /// the next turn of that conversation. `assign_id` has SCV choose the
133    /// CLI's session ID for a new conversation.
134    pub(crate) fn begin(
135        self: &Arc<Self>,
136        agent: &str,
137        handle: Option<&str>,
138        cwd: &Path,
139        assign_id: bool,
140    ) -> Result<TurnGuard, ToolError> {
141        let mut inner = self.inner.lock().expect("conversation lock");
142        let now = Instant::now();
143        let expired: Vec<String> = inner
144            .conversations
145            .iter()
146            .filter(|(_, conversation)| {
147                !conversation.busy && now.duration_since(conversation.last_used) >= self.limits.idle
148            })
149            .map(|(handle, _)| handle.clone())
150            .collect();
151        for handle in &expired {
152            if let Some(conversation) = inner.conversations.remove(handle) {
153                self.remove_marker(conversation.vendor.as_deref());
154            }
155        }
156        let Some(handle) = handle else {
157            return self.start(&mut inner, agent, cwd, assign_id, now);
158        };
159        if !is_handle(handle) {
160            return Err(ToolError(format!(
161                "session {:?} is not a conversation handle; pass the `session` value an \
162                 earlier {agent} call returned, or omit it to start a new conversation",
163                crate::bounded(handle, 80)
164            )));
165        }
166        let Some(conversation) = inner.conversations.get_mut(handle) else {
167            let reason = if expired.iter().any(|expired| expired == handle) {
168                format!(
169                    "was forgotten after {} seconds idle",
170                    self.limits.idle.as_secs()
171                )
172            } else {
173                "is unknown in this session".to_owned()
174            };
175            return Err(ToolError(format!(
176                "conversation {handle} {reason}; omit session to start a new conversation"
177            )));
178        };
179        if conversation.agent != agent {
180            return Err(ToolError(format!(
181                "conversation {handle} belongs to agent_{}, not agent_{agent}",
182                conversation.agent
183            )));
184        }
185        if conversation.busy {
186            return Err(ToolError(format!(
187                "session busy: conversation {handle} is still running a turn"
188            )));
189        }
190        if conversation.cwd != cwd {
191            return Err(ToolError(format!(
192                "conversation {handle} runs in {:?}; continue it there or omit session to \
193                 start a new conversation in {:?}",
194                conversation.cwd, cwd
195            )));
196        }
197        let Some(vendor) = conversation.vendor.clone() else {
198            return Err(ToolError(format!(
199                "conversation {handle} cannot be continued: the agent reported no session"
200            )));
201        };
202        conversation.busy = true;
203        conversation.last_used = now;
204        Ok(TurnGuard {
205            store: Arc::clone(self),
206            handle: handle.to_owned(),
207            turn: conversation.turns + 1,
208            vendor: Some(vendor),
209            attachment: conversation.attachment.clone(),
210            finished: false,
211        })
212    }
213
214    fn start(
215        self: &Arc<Self>,
216        inner: &mut Inner,
217        agent: &str,
218        cwd: &Path,
219        assign_id: bool,
220        now: Instant,
221    ) -> Result<TurnGuard, ToolError> {
222        while inner.conversations.len() >= self.limits.max.max(1) {
223            let Some(oldest) = inner
224                .conversations
225                .iter()
226                .filter(|(_, conversation)| !conversation.busy)
227                .min_by_key(|(_, conversation)| conversation.last_used)
228                .map(|(handle, _)| handle.clone())
229            else {
230                return Err(ToolError(format!(
231                    "all {} conversations of this session are running a turn",
232                    inner.conversations.len()
233                )));
234            };
235            if let Some(conversation) = inner.conversations.remove(&oldest) {
236                self.remove_marker(conversation.vendor.as_deref());
237            }
238        }
239        let number = inner.next.entry(agent.to_owned()).or_insert(0);
240        *number += 1;
241        let handle = format!("{agent}-{number}");
242        let vendor = assign_id.then(|| uuid::Uuid::new_v4().to_string());
243        inner.conversations.insert(
244            handle.clone(),
245            Conversation {
246                agent: agent.to_owned(),
247                cwd: cwd.to_owned(),
248                vendor: None,
249                turns: 0,
250                busy: true,
251                last_used: now,
252                attachment: None,
253            },
254        );
255        if let Some(vendor) = &vendor {
256            self.write_marker(vendor, agent, &handle);
257        }
258        Ok(TurnGuard {
259            store: Arc::clone(self),
260            handle,
261            turn: 1,
262            vendor,
263            attachment: None,
264            finished: false,
265        })
266    }
267
268    fn write_marker(&self, vendor: &str, agent: &str, handle: &str) {
269        let (Some(dir), Some(owner)) = (&self.marker_dir, self.owner) else {
270            return;
271        };
272        let marker = Marker {
273            owner,
274            agent: agent.to_owned(),
275            handle: handle.to_owned(),
276        };
277        // A missing marker only makes `gc` less careful; it never fails a turn.
278        let _ = write_private_json(dir, &format!("{vendor}.json"), &marker);
279    }
280
281    fn remove_marker(&self, vendor: Option<&str>) {
282        if let (Some(dir), Some(vendor)) = (&self.marker_dir, vendor) {
283            let _ = std::fs::remove_file(dir.join(format!("{vendor}.json")));
284        }
285    }
286
287    /// Handles this session remembers, for tests and diagnostics.
288    pub fn handles(&self) -> Vec<String> {
289        let mut handles: Vec<_> = self
290            .inner
291            .lock()
292            .expect("conversation lock")
293            .conversations
294            .keys()
295            .cloned()
296            .collect();
297        handles.sort();
298        handles
299    }
300}
301
302impl Drop for ConversationStore {
303    fn drop(&mut self) {
304        let inner = self.inner.get_mut().expect("conversation lock");
305        let vendors: Vec<_> = inner
306            .conversations
307            .values()
308            .filter_map(|conversation| conversation.vendor.clone())
309            .collect();
310        for vendor in vendors {
311            self.remove_marker(Some(&vendor));
312        }
313    }
314}
315
316impl TurnGuard {
317    /// What the conversation keeps between turns, if anything.
318    pub(crate) fn attachment(&self) -> Option<&Attachment> {
319        self.attachment.as_ref()
320    }
321
322    /// Keep `attachment` with the conversation once this turn finishes.
323    pub(crate) fn attach(&mut self, attachment: Attachment) {
324        self.attachment = Some(attachment);
325    }
326
327    /// End the conversation now, for one that cannot go on (its live child
328    /// exited). Its attachment is dropped.
329    pub(crate) fn forget(mut self) {
330        self.finished = true;
331        let removed = self
332            .store
333            .inner
334            .lock()
335            .expect("conversation lock")
336            .conversations
337            .remove(&self.handle);
338        self.store.remove_marker(
339            removed
340                .and_then(|conversation| conversation.vendor)
341                .as_deref(),
342        );
343        self.store.remove_marker(self.vendor.as_deref());
344    }
345
346    /// Settle the turn. `reported` is the session ID the CLI printed. Returns
347    /// the handle when the conversation can be continued; a first turn that
348    /// failed without the CLI reporting a session is forgotten.
349    pub(crate) fn finish(mut self, reported: Option<String>, completed: bool) -> Option<String> {
350        self.finished = true;
351        let reported = reported.filter(|id| valid_session_id(id));
352        let mut inner = self.store.inner.lock().expect("conversation lock");
353        let first = self.turn == 1;
354        if first && reported.is_none() && !completed {
355            inner.conversations.remove(&self.handle);
356            drop(inner);
357            self.store.remove_marker(self.vendor.as_deref());
358            return None;
359        }
360        let vendor = reported.or_else(|| self.vendor.clone());
361        let conversation = inner.conversations.get_mut(&self.handle)?;
362        if let Some(attachment) = self.attachment.take() {
363            conversation.attachment = Some(attachment);
364        }
365        conversation.busy = false;
366        conversation.turns = self.turn;
367        conversation.last_used = Instant::now();
368        let previous = std::mem::replace(&mut conversation.vendor, vendor.clone());
369        let agent = conversation.agent.clone();
370        drop(inner);
371        if previous != vendor {
372            self.store.remove_marker(previous.as_deref());
373        }
374        match vendor {
375            Some(vendor) => {
376                self.store.write_marker(&vendor, &agent, &self.handle);
377                Some(self.handle.clone())
378            }
379            None => None,
380        }
381    }
382}
383
384impl Drop for TurnGuard {
385    fn drop(&mut self) {
386        if self.finished {
387            return;
388        }
389        let mut inner = self.store.inner.lock().expect("conversation lock");
390        if self.turn == 1 {
391            inner.conversations.remove(&self.handle);
392            drop(inner);
393            self.store.remove_marker(self.vendor.as_deref());
394        } else if let Some(conversation) = inner.conversations.get_mut(&self.handle) {
395            conversation.busy = false;
396        }
397    }
398}
399
400/// What `scv agents gc` found for one agent.
401#[derive(Debug, Default, Clone, PartialEq, Eq)]
402pub struct GcReport {
403    /// Transcripts removed, or that would be with `dry_run`.
404    pub removed: Vec<PathBuf>,
405    pub bytes: u64,
406    /// Old transcripts kept because a live conversation still uses them.
407    pub kept_live: usize,
408}
409
410/// Remove transcripts under `adapter_home` last modified at least
411/// `older_than` ago, keeping any whose name carries the session ID of a live
412/// conversation (a marker in `marker_dir` whose owning process runs).
413/// Symlinks are never followed. Markers of dead owners are removed unless
414/// `dry_run`.
415pub fn collect_garbage(
416    adapter_home: &Path,
417    files: ConversationFiles,
418    marker_dir: &Path,
419    older_than: Duration,
420    dry_run: bool,
421) -> std::io::Result<GcReport> {
422    let older_than = older_than.max(MIN_GC_AGE);
423    let live = live_sessions(marker_dir, dry_run);
424    let root = adapter_home.join(files.dir);
425    let mut report = GcReport::default();
426    let Some(cutoff) = SystemTime::now().checked_sub(older_than) else {
427        return Ok(report);
428    };
429    let mut transcripts = Vec::new();
430    walk(&root, files.extension, &mut transcripts)?;
431    transcripts.sort();
432    for (path, modified, bytes) in transcripts {
433        if modified > cutoff {
434            continue;
435        }
436        let relative = path.strip_prefix(&root).unwrap_or(&path);
437        let in_use = relative.components().any(|component| {
438            let name = component.as_os_str().to_string_lossy();
439            live.iter().any(|id| name.contains(id.as_str()))
440        });
441        if in_use {
442            report.kept_live += 1;
443            continue;
444        }
445        if !dry_run {
446            std::fs::remove_file(&path)?;
447        }
448        report.bytes += bytes;
449        report.removed.push(path);
450    }
451    if !dry_run {
452        remove_empty_dirs(&root, &root);
453    }
454    Ok(report)
455}
456
457/// Remove markers whose owning SCV process no longer runs, returning how
458/// many were removed. A short-lived `scv exec` leaves its markers behind when
459/// it exits; the daemon's reconcile pass clears them.
460pub fn remove_stale_markers(marker_dir: &Path) -> usize {
461    let before = marker_count(marker_dir);
462    let live = live_sessions(marker_dir, false).len();
463    before.saturating_sub(live)
464}
465
466fn marker_count(marker_dir: &Path) -> usize {
467    std::fs::read_dir(marker_dir).map_or(0, |entries| {
468        entries
469            .flatten()
470            .filter(|entry| {
471                entry
472                    .file_name()
473                    .to_str()
474                    .and_then(|name| name.strip_suffix(".json"))
475                    .is_some_and(valid_session_id)
476            })
477            .count()
478    })
479}
480
481/// Session IDs of conversations whose owning SCV process still runs.
482fn live_sessions(marker_dir: &Path, dry_run: bool) -> Vec<String> {
483    let Ok(entries) = std::fs::read_dir(marker_dir) else {
484        return Vec::new();
485    };
486    let mut live = Vec::new();
487    for entry in entries.flatten() {
488        let path = entry.path();
489        let Some(id) = path
490            .file_name()
491            .and_then(|name| name.to_str())
492            .and_then(|name| name.strip_suffix(".json"))
493            .filter(|id| valid_session_id(id))
494            .map(str::to_owned)
495        else {
496            continue;
497        };
498        let alive = std::fs::read(&path)
499            .ok()
500            .and_then(|bytes| serde_json::from_slice::<Marker>(&bytes).ok())
501            .is_some_and(|marker| marker.owner.is_alive());
502        if alive {
503            live.push(id);
504        } else if !dry_run {
505            let _ = std::fs::remove_file(&path);
506        }
507    }
508    live
509}
510
511fn walk(
512    dir: &Path,
513    extension: &str,
514    found: &mut Vec<(PathBuf, SystemTime, u64)>,
515) -> std::io::Result<()> {
516    let entries = match std::fs::read_dir(dir) {
517        Ok(entries) => entries,
518        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
519        Err(error) => return Err(error),
520    };
521    for entry in entries {
522        let entry = entry?;
523        let metadata = std::fs::symlink_metadata(entry.path())?;
524        if metadata.is_dir() {
525            walk(&entry.path(), extension, found)?;
526        } else if metadata.is_file()
527            && entry
528                .path()
529                .extension()
530                .is_some_and(|actual| actual == extension)
531        {
532            found.push((entry.path(), metadata.modified()?, metadata.len()));
533        }
534    }
535    Ok(())
536}
537
538/// Remove directories below `root` left empty, deepest first.
539fn remove_empty_dirs(dir: &Path, root: &Path) {
540    let Ok(entries) = std::fs::read_dir(dir) else {
541        return;
542    };
543    for entry in entries.flatten() {
544        if std::fs::symlink_metadata(entry.path()).is_ok_and(|metadata| metadata.is_dir()) {
545            remove_empty_dirs(&entry.path(), root);
546        }
547    }
548    if dir != root {
549        // Fails, harmlessly, unless the directory is empty.
550        let _ = std::fs::remove_dir(dir);
551    }
552}
553
554/// Parse an age such as `30d`, `12h`, `90m`, `45s`, or plain seconds.
555pub fn parse_age(value: &str) -> Result<Duration, String> {
556    let value = value.trim();
557    let (number, unit) = match value.find(|c: char| !c.is_ascii_digit()) {
558        Some(index) => value.split_at(index),
559        None => (value, "s"),
560    };
561    let number: u64 = number
562        .parse()
563        .map_err(|_| format!("invalid age {value:?}; use a number with s, m, h, or d"))?;
564    let unit = match unit {
565        "s" => 1,
566        "m" => 60,
567        "h" => 3600,
568        "d" => 86400,
569        _ => {
570            return Err(format!(
571                "invalid age {value:?}; use a number with s, m, h, or d"
572            ));
573        }
574    };
575    number
576        .checked_mul(unit)
577        .map(Duration::from_secs)
578        .ok_or_else(|| format!("age {value:?} is too large"))
579}
580
581#[cfg(test)]
582mod tests {
583    use super::*;
584
585    fn store(max: usize, idle: Duration, markers: Option<&Path>) -> Arc<ConversationStore> {
586        Arc::new(ConversationStore::new(
587            ConversationLimits { max, idle },
588            markers.map(Path::to_path_buf),
589        ))
590    }
591
592    const DAY: Duration = Duration::from_secs(86400);
593
594    #[test]
595    fn handles_are_issued_per_agent_and_vendor_ids_are_not_handles() {
596        assert!(is_handle("codex-2"));
597        assert!(is_handle("pi-10"));
598        for not_handle in [
599            "01a0cd5a-7195-7b31-a503-e235d5da7b45",
600            "b514bbf5-a5b7-4bbe-83f9-5824ab41c35c",
601            "codex",
602            "codex-",
603            "-1",
604            "Codex-1",
605            "codex-1a",
606            "../codex-1",
607        ] {
608            assert!(!is_handle(not_handle), "{not_handle}");
609        }
610        let store = store(8, DAY, None);
611        let cwd = Path::new("/w");
612        let first = store.begin("codex", None, cwd, false).unwrap();
613        assert_eq!((first.handle.as_str(), first.turn), ("codex-1", 1));
614        assert_eq!(first.vendor, None);
615        assert_eq!(
616            first.finish(Some("t-1".into()), true).as_deref(),
617            Some("codex-1")
618        );
619        let claude = store.begin("claude", None, cwd, true).unwrap();
620        assert_eq!(claude.handle, "claude-1");
621        assert!(claude.vendor.is_some());
622        claude.finish(None, true);
623        let second = store.begin("codex", None, cwd, false).unwrap();
624        assert_eq!(second.handle, "codex-2");
625    }
626
627    #[test]
628    fn continuing_pins_agent_and_cwd_and_counts_turns() {
629        let store = store(8, DAY, None);
630        let cwd = Path::new("/w/scv");
631        let turn = store.begin("codex", None, cwd, false).unwrap();
632        turn.finish(Some("thread-a".into()), true);
633        let next = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
634        assert_eq!(next.turn, 2);
635        assert_eq!(next.vendor.as_deref(), Some("thread-a"));
636        // One turn at a time.
637        let busy = store
638            .begin("codex", Some("codex-1"), cwd, false)
639            .unwrap_err();
640        assert!(busy.0.starts_with("session busy"), "{}", busy.0);
641        next.finish(Some("thread-a".into()), true);
642        let moved = store
643            .begin("codex", Some("codex-1"), Path::new("/w/other"), false)
644            .unwrap_err();
645        assert!(moved.0.contains("runs in"), "{}", moved.0);
646        let other_agent = store
647            .begin("claude", Some("codex-1"), cwd, false)
648            .unwrap_err();
649        assert!(
650            other_agent.0.contains("belongs to agent_codex"),
651            "{}",
652            other_agent.0
653        );
654        let third = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
655        assert_eq!(third.turn, 3);
656    }
657
658    #[test]
659    fn vendor_ids_and_unknown_handles_are_rejected() {
660        let store = store(8, DAY, None);
661        let cwd = Path::new("/w");
662        store
663            .begin("codex", None, cwd, false)
664            .unwrap()
665            .finish(Some("01a0cd5a-7195-7b31".into()), true);
666        let vendor = store
667            .begin("codex", Some("01a0cd5a-7195-7b31"), cwd, false)
668            .unwrap_err();
669        assert!(
670            vendor.0.contains("not a conversation handle"),
671            "{}",
672            vendor.0
673        );
674        let unknown = store
675            .begin("codex", Some("codex-9"), cwd, false)
676            .unwrap_err();
677        assert!(
678            unknown.0.contains("unknown in this session"),
679            "{}",
680            unknown.0
681        );
682        // Another session's store knows nothing of this one's handles.
683        let other = super::tests::store(8, DAY, None);
684        assert!(other.begin("codex", Some("codex-1"), cwd, false).is_err());
685    }
686
687    #[test]
688    fn timed_out_turns_stay_resumable_but_failed_first_turns_are_forgotten() {
689        let store = store(8, DAY, None);
690        let cwd = Path::new("/w");
691        // A first turn that timed out after the CLI reported its session.
692        let turn = store.begin("codex", None, cwd, false).unwrap();
693        assert_eq!(
694            turn.finish(Some("t".into()), false).as_deref(),
695            Some("codex-1")
696        );
697        assert_eq!(
698            store
699                .begin("codex", Some("codex-1"), cwd, false)
700                .unwrap()
701                .turn,
702            2
703        );
704        // A first turn that failed before the CLI reported anything.
705        let failed = store.begin("codex", None, cwd, false).unwrap();
706        assert_eq!(failed.finish(None, false), None);
707        assert!(!store.handles().contains(&"codex-2".to_owned()));
708        // An abandoned first turn is forgotten; an abandoned later turn frees the conversation.
709        drop(store.begin("codex", None, cwd, false).unwrap());
710        assert_eq!(store.handles(), vec!["codex-1".to_owned()]);
711    }
712
713    #[test]
714    fn limits_forget_the_least_recently_used_and_idle_conversations() {
715        let store = store(2, DAY, None);
716        let cwd = Path::new("/w");
717        for id in ["a", "b"] {
718            store
719                .begin("codex", None, cwd, false)
720                .unwrap()
721                .finish(Some(id.into()), true);
722        }
723        // codex-1 was used more recently than codex-2.
724        store
725            .begin("codex", Some("codex-1"), cwd, false)
726            .unwrap()
727            .finish(Some("a".into()), true);
728        store
729            .begin("codex", None, cwd, false)
730            .unwrap()
731            .finish(Some("c".into()), true);
732        assert_eq!(
733            store.handles(),
734            vec!["codex-1".to_owned(), "codex-3".to_owned()]
735        );
736        // Every remembered conversation busy: no room for another.
737        let busy_a = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
738        let busy_b = store.begin("codex", Some("codex-3"), cwd, false).unwrap();
739        assert!(store.begin("codex", None, cwd, false).is_err());
740        drop((busy_a, busy_b));
741
742        let idle = super::tests::store(8, Duration::ZERO, None);
743        idle.begin("pi", None, cwd, true)
744            .unwrap()
745            .finish(None, true);
746        let expired = idle.begin("pi", Some("pi-1"), cwd, true).unwrap_err();
747        assert!(
748            expired.0.contains("forgotten after 0 seconds idle"),
749            "{}",
750            expired.0
751        );
752    }
753
754    #[test]
755    fn markers_follow_the_conversation_and_gc_keeps_live_transcripts() {
756        let home = tempfile::tempdir().unwrap();
757        let markers = home.path().join("state/conversations");
758        let adapter = home.path().join("agents/codex");
759        let day = adapter.join("sessions/2026/01/02");
760        std::fs::create_dir_all(&day).unwrap();
761        let old = SystemTime::now() - Duration::from_secs(10 * 86400);
762        let transcript = |id: &str| {
763            let path = day.join(format!("rollout-2026-01-02T00-00-00-{id}.jsonl"));
764            std::fs::write(&path, "{}\n").unwrap();
765            std::fs::File::options()
766                .write(true)
767                .open(&path)
768                .unwrap()
769                .set_modified(old)
770                .unwrap();
771            path
772        };
773        let live = transcript("live-id");
774        let stale = transcript("stale-id");
775        let recent = day.join("rollout-recent-id.jsonl");
776        std::fs::write(&recent, "{}\n").unwrap();
777        // A link out of the tree is never followed or removed.
778        let outside = tempfile::NamedTempFile::new().unwrap();
779        std::os::unix::fs::symlink(outside.path(), day.join("link.jsonl")).unwrap();
780
781        let store = store(8, DAY, Some(&markers));
782        store
783            .begin("codex", None, Path::new("/w"), false)
784            .unwrap()
785            .finish(Some("live-id".into()), true);
786        assert!(markers.join("live-id.json").is_file());
787        // A marker left by a process that no longer runs does not protect anything.
788        write_private_json(
789            &markers,
790            "stale-id.json",
791            &Marker {
792                owner: ProcessIdentity {
793                    pid: u32::MAX - 1,
794                    start_time: 1,
795                },
796                agent: "codex".into(),
797                handle: "codex-9".into(),
798            },
799        )
800        .unwrap();
801        let files = ConversationFiles {
802            dir: "sessions",
803            extension: "jsonl",
804        };
805        let dry = collect_garbage(&adapter, files, &markers, DAY, true).unwrap();
806        assert_eq!(dry.removed, vec![stale.clone()]);
807        assert_eq!(dry.kept_live, 1);
808        assert!(stale.exists() && markers.join("stale-id.json").exists());
809        let report = collect_garbage(&adapter, files, &markers, DAY, false).unwrap();
810        assert_eq!(report.removed, vec![stale.clone()]);
811        assert!(!stale.exists() && live.exists() && recent.exists());
812        assert!(outside.path().exists());
813        assert!(!markers.join("stale-id.json").exists());
814        // Ending the session releases its transcripts.
815        drop(store);
816        assert!(!markers.join("live-id.json").exists());
817        // Never younger than the minimum age, whatever was asked.
818        let report = collect_garbage(&adapter, files, &markers, Duration::ZERO, false).unwrap();
819        assert_eq!(report.removed, vec![live]);
820        assert!(recent.exists());
821    }
822
823    #[test]
824    fn ages_parse_with_units() {
825        assert_eq!(parse_age("30d"), Ok(Duration::from_secs(30 * 86400)));
826        assert_eq!(parse_age("12h"), Ok(Duration::from_secs(12 * 3600)));
827        assert_eq!(parse_age("90m"), Ok(Duration::from_secs(5400)));
828        assert_eq!(parse_age("45"), Ok(Duration::from_secs(45)));
829        assert!(parse_age("3w").is_err());
830        assert!(parse_age("d").is_err());
831        assert!(parse_age("-1d").is_err());
832    }
833}