1use 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
32pub(crate) const MIN_GC_AGE: Duration = Duration::from_secs(3600);
35
36const MAX_HANDLE_BYTES: usize = 64;
38
39#[derive(Debug, Clone, Copy)]
40pub struct ConversationLimits {
41 pub max: usize,
44 pub idle: Duration,
46}
47
48#[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 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: HashMap<String, u32>,
77}
78
79#[derive(Debug)]
81pub(crate) struct ConversationStore {
82 limits: ConversationLimits,
83 marker_dir: Option<PathBuf>,
84 owner: Option<ProcessIdentity>,
85 inner: Mutex<Inner>,
86 changed: Notify,
88}
89
90#[derive(Debug, Serialize, Deserialize)]
92struct Marker {
93 owner: ProcessIdentity,
94 agent: String,
95 handle: String,
96}
97
98#[derive(Debug)]
101pub(crate) struct TurnGuard {
102 store: Arc<ConversationStore>,
103 pub(crate) handle: String,
104 pub(crate) turn: u32,
105 pub(crate) vendor: Option<String>,
107 attachment: Option<Attachment>,
110 finished: bool,
111}
112
113pub(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
127pub(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 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 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 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 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 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 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 #[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 pub(crate) fn attachment(&self) -> Option<&Attachment> {
362 self.attachment.as_ref()
363 }
364
365 pub(crate) fn attach(&mut self, attachment: Attachment) {
367 self.attachment = Some(attachment);
368 }
369
370 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 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 self.store.changed.notify_waiters();
435 }
436}
437
438#[derive(Debug, Default, Clone, PartialEq, Eq)]
440pub struct GcReport {
441 pub removed: Vec<PathBuf>,
443 pub bytes: u64,
444 pub kept_live: usize,
446}
447
448pub 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
495pub(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
519fn 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
576fn 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 let _ = std::fs::remove_dir(dir);
589 }
590}
591
592pub 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;