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
123impl ConversationStore {
124 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 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 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 #[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 pub(crate) fn attachment(&self) -> Option<&Attachment> {
316 self.attachment.as_ref()
317 }
318
319 pub(crate) fn attach(&mut self, attachment: Attachment) {
321 self.attachment = Some(attachment);
322 }
323
324 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 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#[derive(Debug, Default, Clone, PartialEq, Eq)]
393pub struct GcReport {
394 pub removed: Vec<PathBuf>,
396 pub bytes: u64,
397 pub kept_live: usize,
399}
400
401pub 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
448pub(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
472fn 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
529fn 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 let _ = std::fs::remove_dir(dir);
542 }
543}
544
545pub 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;