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};
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
30pub(crate) const MIN_GC_AGE: Duration = Duration::from_secs(3600);
33
34const MAX_HANDLE_BYTES: usize = 64;
36
37#[derive(Debug, Clone, Copy)]
38pub struct ConversationLimits {
39 pub max: usize,
42 pub idle: Duration,
44}
45
46#[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 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: HashMap<String, u32>,
75}
76
77#[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#[derive(Debug, Serialize, Deserialize)]
88struct Marker {
89 owner: ProcessIdentity,
90 agent: String,
91 handle: String,
92}
93
94#[derive(Debug)]
97pub(crate) struct TurnGuard {
98 store: Arc<ConversationStore>,
99 pub(crate) handle: String,
100 pub(crate) turn: u32,
101 pub(crate) vendor: Option<String>,
103 attachment: Option<Attachment>,
106 finished: bool,
107}
108
109pub(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
123pub(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 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 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 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 #[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 pub(crate) fn attachment(&self) -> Option<&Attachment> {
323 self.attachment.as_ref()
324 }
325
326 pub(crate) fn attach(&mut self, attachment: Attachment) {
328 self.attachment = Some(attachment);
329 }
330
331 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 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#[derive(Debug, Default, Clone, PartialEq, Eq)]
400pub struct GcReport {
401 pub removed: Vec<PathBuf>,
403 pub bytes: u64,
404 pub kept_live: usize,
406}
407
408pub 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
455pub(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
479fn 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
536fn 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 let _ = std::fs::remove_dir(dir);
549 }
550}
551
552pub 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;