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