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