Skip to main content

scv_tools/delegate/
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    delegate::{
23        adapters::ConversationFiles,
24        output::valid_session_id,
25        records::{ProcessIdentity, write_private_json},
26    },
27    sync::lock,
28};
29
30/// Transcripts younger than this are never removed, so a turn that is still
31/// running (and so still writing its transcript) is safe from `gc`.
32pub(crate) const MIN_GC_AGE: Duration = Duration::from_secs(3600);
33
34/// Longest handle accepted from the model.
35const MAX_HANDLE_BYTES: usize = 64;
36
37#[derive(Debug, Clone, Copy)]
38pub struct ConversationLimits {
39    /// Conversations remembered per session; starting another forgets the
40    /// least recently used idle one.
41    pub max: usize,
42    /// A conversation unused this long is forgotten.
43    pub idle: Duration,
44}
45
46/// What a live conversation keeps between turns, such as its running child
47/// process. Dropping the last reference (when the conversation is forgotten,
48/// expires, or its session ends) is what shuts the child down.
49#[derive(Clone)]
50pub(crate) struct Attachment(pub(crate) Arc<dyn Any + Send + Sync>);
51
52impl fmt::Debug for Attachment {
53    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
54        formatter.write_str("Attachment")
55    }
56}
57
58#[derive(Debug)]
59struct Conversation {
60    agent: String,
61    cwd: PathBuf,
62    /// The CLI's session ID, once known.
63    vendor: Option<String>,
64    turns: u32,
65    busy: bool,
66    last_used: Instant,
67    attachment: Option<Attachment>,
68}
69
70#[derive(Debug, Default)]
71struct Inner {
72    conversations: HashMap<String, Conversation>,
73    /// Next handle number per agent.
74    next: HashMap<String, u32>,
75}
76
77/// One session's conversations with its delegated agents.
78#[derive(Debug)]
79pub(crate) struct ConversationStore {
80    limits: ConversationLimits,
81    marker_dir: Option<PathBuf>,
82    owner: Option<ProcessIdentity>,
83    inner: Mutex<Inner>,
84}
85
86/// On-disk marker keeping a live conversation's transcript from `gc`.
87#[derive(Debug, Serialize, Deserialize)]
88struct Marker {
89    owner: ProcessIdentity,
90    agent: String,
91    handle: String,
92}
93
94/// A turn in progress. Finish it with [`TurnGuard::finish`]; dropping it
95/// instead (a cancelled or abandoned call) frees the conversation.
96#[derive(Debug)]
97pub(crate) struct TurnGuard {
98    store: Arc<ConversationStore>,
99    pub(crate) handle: String,
100    pub(crate) turn: u32,
101    /// Continuing: the known session ID. Starting: the ID SCV chose, if any.
102    pub(crate) vendor: Option<String>,
103    /// Continuing: what the conversation kept. Starting: what
104    /// [`TurnGuard::attach`] set, stored when the turn finishes.
105    attachment: Option<Attachment>,
106    finished: bool,
107}
108
109/// Whether `value` has the shape of an SCV conversation handle
110/// (`<agent>-<number>`). CLI session IDs never do.
111pub(crate) fn is_handle(value: &str) -> bool {
112    value.len() <= MAX_HANDLE_BYTES
113        && value.split_once('-').is_some_and(|(agent, number)| {
114            !agent.is_empty()
115                && agent
116                    .chars()
117                    .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
118                && !number.is_empty()
119                && number.chars().all(|c| c.is_ascii_digit())
120        })
121}
122
123impl ConversationStore {
124    /// `marker_dir` is `$SCV_HOME/state/conversations`; `None` keeps no markers.
125    pub(crate) fn new(limits: ConversationLimits, marker_dir: Option<PathBuf>) -> Self {
126        Self {
127            limits,
128            marker_dir,
129            owner: ProcessIdentity::current(),
130            inner: Mutex::new(Inner::default()),
131        }
132    }
133
134    /// Start a turn: a new conversation when `handle` is `None`, otherwise
135    /// the next turn of that conversation. `assign_id` has SCV choose the
136    /// CLI's session ID for a new conversation.
137    pub(crate) fn begin(
138        self: &Arc<Self>,
139        agent: &str,
140        handle: Option<&str>,
141        cwd: &Path,
142        assign_id: bool,
143    ) -> Result<TurnGuard, ToolError> {
144        let mut inner = lock(&self.inner);
145        let now = Instant::now();
146        let expired: Vec<String> = inner
147            .conversations
148            .iter()
149            .filter(|(_, conversation)| {
150                !conversation.busy && now.duration_since(conversation.last_used) >= self.limits.idle
151            })
152            .map(|(handle, _)| handle.clone())
153            .collect();
154        for handle in &expired {
155            if let Some(conversation) = inner.conversations.remove(handle) {
156                self.remove_marker(conversation.vendor.as_deref());
157            }
158        }
159        let Some(handle) = handle else {
160            return self.start(&mut inner, agent, cwd, assign_id, now);
161        };
162        if !is_handle(handle) {
163            return Err(ToolError::invalid_arguments(format!(
164                "session {:?} is not a conversation handle; pass the `session` value an \
165                 earlier {agent} call returned, or omit it to start a new conversation",
166                crate::args::bounded(handle, 80)
167            )));
168        }
169        let Some(conversation) = inner.conversations.get_mut(handle) else {
170            let reason = if expired.iter().any(|expired| expired == handle) {
171                format!(
172                    "was forgotten after {} seconds idle",
173                    self.limits.idle.as_secs()
174                )
175            } else {
176                "is unknown in this session".to_owned()
177            };
178            return Err(ToolError::invalid_arguments(format!(
179                "conversation {handle} {reason}; omit session to start a new conversation"
180            )));
181        };
182        if conversation.agent != agent {
183            return Err(ToolError::invalid_arguments(format!(
184                "conversation {handle} belongs to agent_{}, not agent_{agent}",
185                conversation.agent
186            )));
187        }
188        if conversation.busy {
189            return Err(ToolError::failed(format!(
190                "session busy: conversation {handle} is still running a turn"
191            )));
192        }
193        if conversation.cwd != cwd {
194            return Err(ToolError::invalid_arguments(format!(
195                "conversation {handle} runs in {:?}; continue it there or omit session to \
196                 start a new conversation in {:?}",
197                conversation.cwd, cwd
198            )));
199        }
200        let Some(vendor) = conversation.vendor.clone() else {
201            return Err(ToolError::failed(format!(
202                "conversation {handle} cannot be continued: the agent reported no session"
203            )));
204        };
205        conversation.busy = true;
206        conversation.last_used = now;
207        Ok(TurnGuard {
208            store: Arc::clone(self),
209            handle: handle.to_owned(),
210            turn: conversation.turns + 1,
211            vendor: Some(vendor),
212            attachment: conversation.attachment.clone(),
213            finished: false,
214        })
215    }
216
217    fn start(
218        self: &Arc<Self>,
219        inner: &mut Inner,
220        agent: &str,
221        cwd: &Path,
222        assign_id: bool,
223        now: Instant,
224    ) -> Result<TurnGuard, ToolError> {
225        while inner.conversations.len() >= self.limits.max.max(1) {
226            let Some(oldest) = inner
227                .conversations
228                .iter()
229                .filter(|(_, conversation)| !conversation.busy)
230                .min_by_key(|(_, conversation)| conversation.last_used)
231                .map(|(handle, _)| handle.clone())
232            else {
233                return Err(ToolError::limit(format!(
234                    "all {} conversations of this session are running a turn",
235                    inner.conversations.len()
236                )));
237            };
238            if let Some(conversation) = inner.conversations.remove(&oldest) {
239                self.remove_marker(conversation.vendor.as_deref());
240            }
241        }
242        let number = inner.next.entry(agent.to_owned()).or_insert(0);
243        *number += 1;
244        let handle = format!("{agent}-{number}");
245        let vendor = assign_id.then(|| uuid::Uuid::new_v4().to_string());
246        inner.conversations.insert(
247            handle.clone(),
248            Conversation {
249                agent: agent.to_owned(),
250                cwd: cwd.to_owned(),
251                vendor: None,
252                turns: 0,
253                busy: true,
254                last_used: now,
255                attachment: None,
256            },
257        );
258        if let Some(vendor) = &vendor {
259            self.write_marker(vendor, agent, &handle);
260        }
261        Ok(TurnGuard {
262            store: Arc::clone(self),
263            handle,
264            turn: 1,
265            vendor,
266            attachment: None,
267            finished: false,
268        })
269    }
270
271    fn write_marker(&self, vendor: &str, agent: &str, handle: &str) {
272        let (Some(dir), Some(owner)) = (&self.marker_dir, self.owner) else {
273            return;
274        };
275        let marker = Marker {
276            owner,
277            agent: agent.to_owned(),
278            handle: handle.to_owned(),
279        };
280        // A missing marker only makes `gc` less careful; it never fails a turn.
281        let _ = write_private_json(dir, &format!("{vendor}.json"), &marker);
282    }
283
284    fn remove_marker(&self, vendor: Option<&str>) {
285        if let (Some(dir), Some(vendor)) = (&self.marker_dir, vendor) {
286            let _ = std::fs::remove_file(dir.join(format!("{vendor}.json")));
287        }
288    }
289
290    /// Handles this session remembers, for tests.
291    #[cfg(test)]
292    pub(crate) fn handles(&self) -> Vec<String> {
293        let mut handles: Vec<_> = lock(&self.inner).conversations.keys().cloned().collect();
294        handles.sort();
295        handles
296    }
297}
298
299impl Drop for ConversationStore {
300    fn drop(&mut self) {
301        let inner = self.inner.get_mut().expect("conversation lock");
302        let vendors: Vec<_> = inner
303            .conversations
304            .values()
305            .filter_map(|conversation| conversation.vendor.clone())
306            .collect();
307        for vendor in vendors {
308            self.remove_marker(Some(&vendor));
309        }
310    }
311}
312
313impl TurnGuard {
314    /// What the conversation keeps between turns, if anything.
315    pub(crate) fn attachment(&self) -> Option<&Attachment> {
316        self.attachment.as_ref()
317    }
318
319    /// Keep `attachment` with the conversation once this turn finishes.
320    pub(crate) fn attach(&mut self, attachment: Attachment) {
321        self.attachment = Some(attachment);
322    }
323
324    /// End the conversation now, for one that cannot go on (its live child
325    /// exited). Its attachment is dropped.
326    pub(crate) fn forget(mut self) {
327        self.finished = true;
328        let removed = lock(&self.store.inner).conversations.remove(&self.handle);
329        self.store.remove_marker(
330            removed
331                .and_then(|conversation| conversation.vendor)
332                .as_deref(),
333        );
334        self.store.remove_marker(self.vendor.as_deref());
335    }
336
337    /// Settle the turn. `reported` is the session ID the CLI printed. Returns
338    /// the handle when the conversation can be continued; a first turn that
339    /// failed without the CLI reporting a session is forgotten.
340    pub(crate) fn finish(mut self, reported: Option<String>, completed: bool) -> Option<String> {
341        self.finished = true;
342        let reported = reported.filter(|id| valid_session_id(id));
343        let mut inner = lock(&self.store.inner);
344        let first = self.turn == 1;
345        if first && reported.is_none() && !completed {
346            inner.conversations.remove(&self.handle);
347            drop(inner);
348            self.store.remove_marker(self.vendor.as_deref());
349            return None;
350        }
351        let vendor = reported.or_else(|| self.vendor.clone());
352        let conversation = inner.conversations.get_mut(&self.handle)?;
353        if let Some(attachment) = self.attachment.take() {
354            conversation.attachment = Some(attachment);
355        }
356        conversation.busy = false;
357        conversation.turns = self.turn;
358        conversation.last_used = Instant::now();
359        let previous = std::mem::replace(&mut conversation.vendor, vendor.clone());
360        let agent = conversation.agent.clone();
361        drop(inner);
362        if previous != vendor {
363            self.store.remove_marker(previous.as_deref());
364        }
365        match vendor {
366            Some(vendor) => {
367                self.store.write_marker(&vendor, &agent, &self.handle);
368                Some(self.handle.clone())
369            }
370            None => None,
371        }
372    }
373}
374
375impl Drop for TurnGuard {
376    fn drop(&mut self) {
377        if self.finished {
378            return;
379        }
380        let mut inner = lock(&self.store.inner);
381        if self.turn == 1 {
382            inner.conversations.remove(&self.handle);
383            drop(inner);
384            self.store.remove_marker(self.vendor.as_deref());
385        } else if let Some(conversation) = inner.conversations.get_mut(&self.handle) {
386            conversation.busy = false;
387        }
388    }
389}
390
391/// What `scv agents gc` found for one agent.
392#[derive(Debug, Default, Clone, PartialEq, Eq)]
393pub struct GcReport {
394    /// Transcripts removed, or that would be with `dry_run`.
395    pub removed: Vec<PathBuf>,
396    pub bytes: u64,
397    /// Old transcripts kept because a live conversation still uses them.
398    pub kept_live: usize,
399}
400
401/// Remove transcripts under `adapter_home` last modified at least
402/// `older_than` ago, keeping any whose name carries the session ID of a live
403/// conversation (a marker in `marker_dir` whose owning process runs).
404/// Symlinks are never followed. Markers of dead owners are removed unless
405/// `dry_run`.
406pub fn collect_garbage(
407    adapter_home: &Path,
408    files: ConversationFiles,
409    marker_dir: &Path,
410    older_than: Duration,
411    dry_run: bool,
412) -> std::io::Result<GcReport> {
413    let older_than = older_than.max(MIN_GC_AGE);
414    let live = live_sessions(marker_dir, dry_run);
415    let root = adapter_home.join(files.dir);
416    let mut report = GcReport::default();
417    let Some(cutoff) = SystemTime::now().checked_sub(older_than) else {
418        return Ok(report);
419    };
420    let mut transcripts = Vec::new();
421    walk(&root, files.extension, &mut transcripts)?;
422    transcripts.sort();
423    for (path, modified, bytes) in transcripts {
424        if modified > cutoff {
425            continue;
426        }
427        let relative = path.strip_prefix(&root).unwrap_or(&path);
428        let in_use = relative.components().any(|component| {
429            let name = component.as_os_str().to_string_lossy();
430            live.iter().any(|id| name.contains(id.as_str()))
431        });
432        if in_use {
433            report.kept_live += 1;
434            continue;
435        }
436        if !dry_run {
437            std::fs::remove_file(&path)?;
438        }
439        report.bytes += bytes;
440        report.removed.push(path);
441    }
442    if !dry_run {
443        remove_empty_dirs(&root, &root);
444    }
445    Ok(report)
446}
447
448/// Remove markers whose owning SCV process no longer runs, returning how
449/// many were removed. A short-lived `scv exec` leaves its markers behind when
450/// it exits; the daemon's reconcile pass clears them.
451pub(crate) fn remove_stale_markers(marker_dir: &Path) -> usize {
452    let before = marker_count(marker_dir);
453    let live = live_sessions(marker_dir, false).len();
454    before.saturating_sub(live)
455}
456
457fn marker_count(marker_dir: &Path) -> usize {
458    std::fs::read_dir(marker_dir).map_or(0, |entries| {
459        entries
460            .flatten()
461            .filter(|entry| {
462                entry
463                    .file_name()
464                    .to_str()
465                    .and_then(|name| name.strip_suffix(".json"))
466                    .is_some_and(valid_session_id)
467            })
468            .count()
469    })
470}
471
472/// Session IDs of conversations whose owning SCV process still runs.
473fn live_sessions(marker_dir: &Path, dry_run: bool) -> Vec<String> {
474    let Ok(entries) = std::fs::read_dir(marker_dir) else {
475        return Vec::new();
476    };
477    let mut live = Vec::new();
478    for entry in entries.flatten() {
479        let path = entry.path();
480        let Some(id) = path
481            .file_name()
482            .and_then(|name| name.to_str())
483            .and_then(|name| name.strip_suffix(".json"))
484            .filter(|id| valid_session_id(id))
485            .map(str::to_owned)
486        else {
487            continue;
488        };
489        let alive = std::fs::read(&path)
490            .ok()
491            .and_then(|bytes| serde_json::from_slice::<Marker>(&bytes).ok())
492            .is_some_and(|marker| marker.owner.is_alive());
493        if alive {
494            live.push(id);
495        } else if !dry_run {
496            let _ = std::fs::remove_file(&path);
497        }
498    }
499    live
500}
501
502fn walk(
503    dir: &Path,
504    extension: &str,
505    found: &mut Vec<(PathBuf, SystemTime, u64)>,
506) -> std::io::Result<()> {
507    let entries = match std::fs::read_dir(dir) {
508        Ok(entries) => entries,
509        Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
510        Err(error) => return Err(error),
511    };
512    for entry in entries {
513        let entry = entry?;
514        let metadata = std::fs::symlink_metadata(entry.path())?;
515        if metadata.is_dir() {
516            walk(&entry.path(), extension, found)?;
517        } else if metadata.is_file()
518            && entry
519                .path()
520                .extension()
521                .is_some_and(|actual| actual == extension)
522        {
523            found.push((entry.path(), metadata.modified()?, metadata.len()));
524        }
525    }
526    Ok(())
527}
528
529/// Remove directories below `root` left empty, deepest first.
530fn remove_empty_dirs(dir: &Path, root: &Path) {
531    let Ok(entries) = std::fs::read_dir(dir) else {
532        return;
533    };
534    for entry in entries.flatten() {
535        if std::fs::symlink_metadata(entry.path()).is_ok_and(|metadata| metadata.is_dir()) {
536            remove_empty_dirs(&entry.path(), root);
537        }
538    }
539    if dir != root {
540        // Fails, harmlessly, unless the directory is empty.
541        let _ = std::fs::remove_dir(dir);
542    }
543}
544
545/// Parse an age such as `30d`, `12h`, `90m`, `45s`, or plain seconds.
546pub fn parse_age(value: &str) -> Result<Duration, String> {
547    let value = value.trim();
548    let (number, unit) = match value.find(|c: char| !c.is_ascii_digit()) {
549        Some(index) => value.split_at(index),
550        None => (value, "s"),
551    };
552    let number: u64 = number
553        .parse()
554        .map_err(|_| format!("invalid age {value:?}; use a number with s, m, h, or d"))?;
555    let unit = match unit {
556        "s" => 1,
557        "m" => 60,
558        "h" => 3600,
559        "d" => 86400,
560        _ => {
561            return Err(format!(
562                "invalid age {value:?}; use a number with s, m, h, or d"
563            ));
564        }
565    };
566    number
567        .checked_mul(unit)
568        .map(Duration::from_secs)
569        .ok_or_else(|| format!("age {value:?} is too large"))
570}
571
572#[cfg(test)]
573mod tests;