1use std::{
10 collections::HashMap,
11 path::{Path, PathBuf},
12 sync::{Arc, Mutex},
13 time::{Duration, Instant, SystemTime},
14};
15
16use scv_core::ToolError;
17use serde::{Deserialize, Serialize};
18
19use crate::{
20 adapters::ConversationFiles,
21 agent_output::valid_session_id,
22 delegation::{ProcessIdentity, write_private_json},
23};
24
25pub const MIN_GC_AGE: Duration = Duration::from_secs(3600);
28
29const MAX_HANDLE_BYTES: usize = 64;
31
32#[derive(Debug, Clone, Copy)]
33pub struct ConversationLimits {
34 pub max: usize,
37 pub idle: Duration,
39}
40
41#[derive(Debug)]
42struct Conversation {
43 agent: String,
44 cwd: PathBuf,
45 vendor: Option<String>,
47 turns: u32,
48 busy: bool,
49 last_used: Instant,
50}
51
52#[derive(Debug, Default)]
53struct Inner {
54 conversations: HashMap<String, Conversation>,
55 next: HashMap<String, u32>,
57}
58
59#[derive(Debug)]
61pub struct ConversationStore {
62 limits: ConversationLimits,
63 marker_dir: Option<PathBuf>,
64 owner: Option<ProcessIdentity>,
65 inner: Mutex<Inner>,
66}
67
68#[derive(Debug, Serialize, Deserialize)]
70struct Marker {
71 owner: ProcessIdentity,
72 agent: String,
73 handle: String,
74}
75
76#[derive(Debug)]
79pub(crate) struct TurnGuard {
80 store: Arc<ConversationStore>,
81 pub handle: String,
82 pub turn: u32,
83 pub vendor: Option<String>,
85 finished: bool,
86}
87
88pub fn is_handle(value: &str) -> bool {
91 value.len() <= MAX_HANDLE_BYTES
92 && value.split_once('-').is_some_and(|(agent, number)| {
93 !agent.is_empty()
94 && agent
95 .chars()
96 .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit())
97 && !number.is_empty()
98 && number.chars().all(|c| c.is_ascii_digit())
99 })
100}
101
102impl ConversationStore {
103 pub fn new(limits: ConversationLimits, marker_dir: Option<PathBuf>) -> Self {
105 Self {
106 limits,
107 marker_dir,
108 owner: ProcessIdentity::current(),
109 inner: Mutex::new(Inner::default()),
110 }
111 }
112
113 pub(crate) fn begin(
117 self: &Arc<Self>,
118 agent: &str,
119 handle: Option<&str>,
120 cwd: &Path,
121 assign_id: bool,
122 ) -> Result<TurnGuard, ToolError> {
123 let mut inner = self.inner.lock().expect("conversation lock");
124 let now = Instant::now();
125 let expired: Vec<String> = inner
126 .conversations
127 .iter()
128 .filter(|(_, conversation)| {
129 !conversation.busy && now.duration_since(conversation.last_used) >= self.limits.idle
130 })
131 .map(|(handle, _)| handle.clone())
132 .collect();
133 for handle in &expired {
134 if let Some(conversation) = inner.conversations.remove(handle) {
135 self.remove_marker(conversation.vendor.as_deref());
136 }
137 }
138 let Some(handle) = handle else {
139 return self.start(&mut inner, agent, cwd, assign_id, now);
140 };
141 if !is_handle(handle) {
142 return Err(ToolError(format!(
143 "session {:?} is not a conversation handle; pass the `session` value an \
144 earlier {agent} call returned, or omit it to start a new conversation",
145 crate::bounded(handle, 80)
146 )));
147 }
148 let Some(conversation) = inner.conversations.get_mut(handle) else {
149 let reason = if expired.iter().any(|expired| expired == handle) {
150 format!(
151 "was forgotten after {} seconds idle",
152 self.limits.idle.as_secs()
153 )
154 } else {
155 "is unknown in this session".to_owned()
156 };
157 return Err(ToolError(format!(
158 "conversation {handle} {reason}; omit session to start a new conversation"
159 )));
160 };
161 if conversation.agent != agent {
162 return Err(ToolError(format!(
163 "conversation {handle} belongs to agent_{}, not agent_{agent}",
164 conversation.agent
165 )));
166 }
167 if conversation.busy {
168 return Err(ToolError(format!(
169 "session busy: conversation {handle} is still running a turn"
170 )));
171 }
172 if conversation.cwd != cwd {
173 return Err(ToolError(format!(
174 "conversation {handle} runs in {:?}; continue it there or omit session to \
175 start a new conversation in {:?}",
176 conversation.cwd, cwd
177 )));
178 }
179 let Some(vendor) = conversation.vendor.clone() else {
180 return Err(ToolError(format!(
181 "conversation {handle} cannot be continued: the agent reported no session"
182 )));
183 };
184 conversation.busy = true;
185 conversation.last_used = now;
186 Ok(TurnGuard {
187 store: Arc::clone(self),
188 handle: handle.to_owned(),
189 turn: conversation.turns + 1,
190 vendor: Some(vendor),
191 finished: false,
192 })
193 }
194
195 fn start(
196 self: &Arc<Self>,
197 inner: &mut Inner,
198 agent: &str,
199 cwd: &Path,
200 assign_id: bool,
201 now: Instant,
202 ) -> Result<TurnGuard, ToolError> {
203 while inner.conversations.len() >= self.limits.max.max(1) {
204 let Some(oldest) = inner
205 .conversations
206 .iter()
207 .filter(|(_, conversation)| !conversation.busy)
208 .min_by_key(|(_, conversation)| conversation.last_used)
209 .map(|(handle, _)| handle.clone())
210 else {
211 return Err(ToolError(format!(
212 "all {} conversations of this session are running a turn",
213 inner.conversations.len()
214 )));
215 };
216 if let Some(conversation) = inner.conversations.remove(&oldest) {
217 self.remove_marker(conversation.vendor.as_deref());
218 }
219 }
220 let number = inner.next.entry(agent.to_owned()).or_insert(0);
221 *number += 1;
222 let handle = format!("{agent}-{number}");
223 let vendor = assign_id.then(|| uuid::Uuid::new_v4().to_string());
224 inner.conversations.insert(
225 handle.clone(),
226 Conversation {
227 agent: agent.to_owned(),
228 cwd: cwd.to_owned(),
229 vendor: None,
230 turns: 0,
231 busy: true,
232 last_used: now,
233 },
234 );
235 if let Some(vendor) = &vendor {
236 self.write_marker(vendor, agent, &handle);
237 }
238 Ok(TurnGuard {
239 store: Arc::clone(self),
240 handle,
241 turn: 1,
242 vendor,
243 finished: false,
244 })
245 }
246
247 fn write_marker(&self, vendor: &str, agent: &str, handle: &str) {
248 let (Some(dir), Some(owner)) = (&self.marker_dir, self.owner) else {
249 return;
250 };
251 let marker = Marker {
252 owner,
253 agent: agent.to_owned(),
254 handle: handle.to_owned(),
255 };
256 let _ = write_private_json(dir, &format!("{vendor}.json"), &marker);
258 }
259
260 fn remove_marker(&self, vendor: Option<&str>) {
261 if let (Some(dir), Some(vendor)) = (&self.marker_dir, vendor) {
262 let _ = std::fs::remove_file(dir.join(format!("{vendor}.json")));
263 }
264 }
265
266 pub fn handles(&self) -> Vec<String> {
268 let mut handles: Vec<_> = self
269 .inner
270 .lock()
271 .expect("conversation lock")
272 .conversations
273 .keys()
274 .cloned()
275 .collect();
276 handles.sort();
277 handles
278 }
279}
280
281impl Drop for ConversationStore {
282 fn drop(&mut self) {
283 let inner = self.inner.get_mut().expect("conversation lock");
284 let vendors: Vec<_> = inner
285 .conversations
286 .values()
287 .filter_map(|conversation| conversation.vendor.clone())
288 .collect();
289 for vendor in vendors {
290 self.remove_marker(Some(&vendor));
291 }
292 }
293}
294
295impl TurnGuard {
296 pub(crate) fn finish(mut self, reported: Option<String>, completed: bool) -> Option<String> {
300 self.finished = true;
301 let reported = reported.filter(|id| valid_session_id(id));
302 let mut inner = self.store.inner.lock().expect("conversation lock");
303 let first = self.turn == 1;
304 if first && reported.is_none() && !completed {
305 inner.conversations.remove(&self.handle);
306 drop(inner);
307 self.store.remove_marker(self.vendor.as_deref());
308 return None;
309 }
310 let vendor = reported.or_else(|| self.vendor.clone());
311 let conversation = inner.conversations.get_mut(&self.handle)?;
312 conversation.busy = false;
313 conversation.turns = self.turn;
314 conversation.last_used = Instant::now();
315 let previous = std::mem::replace(&mut conversation.vendor, vendor.clone());
316 let agent = conversation.agent.clone();
317 drop(inner);
318 if previous != vendor {
319 self.store.remove_marker(previous.as_deref());
320 }
321 match vendor {
322 Some(vendor) => {
323 self.store.write_marker(&vendor, &agent, &self.handle);
324 Some(self.handle.clone())
325 }
326 None => None,
327 }
328 }
329}
330
331impl Drop for TurnGuard {
332 fn drop(&mut self) {
333 if self.finished {
334 return;
335 }
336 let mut inner = self.store.inner.lock().expect("conversation lock");
337 if self.turn == 1 {
338 inner.conversations.remove(&self.handle);
339 drop(inner);
340 self.store.remove_marker(self.vendor.as_deref());
341 } else if let Some(conversation) = inner.conversations.get_mut(&self.handle) {
342 conversation.busy = false;
343 }
344 }
345}
346
347#[derive(Debug, Default, Clone, PartialEq, Eq)]
349pub struct GcReport {
350 pub removed: Vec<PathBuf>,
352 pub bytes: u64,
353 pub kept_live: usize,
355}
356
357pub fn collect_garbage(
363 adapter_home: &Path,
364 files: ConversationFiles,
365 marker_dir: &Path,
366 older_than: Duration,
367 dry_run: bool,
368) -> std::io::Result<GcReport> {
369 let older_than = older_than.max(MIN_GC_AGE);
370 let live = live_sessions(marker_dir, dry_run);
371 let root = adapter_home.join(files.dir);
372 let mut report = GcReport::default();
373 let Some(cutoff) = SystemTime::now().checked_sub(older_than) else {
374 return Ok(report);
375 };
376 let mut transcripts = Vec::new();
377 walk(&root, files.extension, &mut transcripts)?;
378 transcripts.sort();
379 for (path, modified, bytes) in transcripts {
380 if modified > cutoff {
381 continue;
382 }
383 let relative = path.strip_prefix(&root).unwrap_or(&path);
384 let in_use = relative.components().any(|component| {
385 let name = component.as_os_str().to_string_lossy();
386 live.iter().any(|id| name.contains(id.as_str()))
387 });
388 if in_use {
389 report.kept_live += 1;
390 continue;
391 }
392 if !dry_run {
393 std::fs::remove_file(&path)?;
394 }
395 report.bytes += bytes;
396 report.removed.push(path);
397 }
398 if !dry_run {
399 remove_empty_dirs(&root, &root);
400 }
401 Ok(report)
402}
403
404pub fn remove_stale_markers(marker_dir: &Path) -> usize {
408 let before = marker_count(marker_dir);
409 let live = live_sessions(marker_dir, false).len();
410 before.saturating_sub(live)
411}
412
413fn marker_count(marker_dir: &Path) -> usize {
414 std::fs::read_dir(marker_dir).map_or(0, |entries| {
415 entries
416 .flatten()
417 .filter(|entry| {
418 entry
419 .file_name()
420 .to_str()
421 .and_then(|name| name.strip_suffix(".json"))
422 .is_some_and(valid_session_id)
423 })
424 .count()
425 })
426}
427
428fn live_sessions(marker_dir: &Path, dry_run: bool) -> Vec<String> {
430 let Ok(entries) = std::fs::read_dir(marker_dir) else {
431 return Vec::new();
432 };
433 let mut live = Vec::new();
434 for entry in entries.flatten() {
435 let path = entry.path();
436 let Some(id) = path
437 .file_name()
438 .and_then(|name| name.to_str())
439 .and_then(|name| name.strip_suffix(".json"))
440 .filter(|id| valid_session_id(id))
441 .map(str::to_owned)
442 else {
443 continue;
444 };
445 let alive = std::fs::read(&path)
446 .ok()
447 .and_then(|bytes| serde_json::from_slice::<Marker>(&bytes).ok())
448 .is_some_and(|marker| marker.owner.is_alive());
449 if alive {
450 live.push(id);
451 } else if !dry_run {
452 let _ = std::fs::remove_file(&path);
453 }
454 }
455 live
456}
457
458fn walk(
459 dir: &Path,
460 extension: &str,
461 found: &mut Vec<(PathBuf, SystemTime, u64)>,
462) -> std::io::Result<()> {
463 let entries = match std::fs::read_dir(dir) {
464 Ok(entries) => entries,
465 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
466 Err(error) => return Err(error),
467 };
468 for entry in entries {
469 let entry = entry?;
470 let metadata = std::fs::symlink_metadata(entry.path())?;
471 if metadata.is_dir() {
472 walk(&entry.path(), extension, found)?;
473 } else if metadata.is_file()
474 && entry
475 .path()
476 .extension()
477 .is_some_and(|actual| actual == extension)
478 {
479 found.push((entry.path(), metadata.modified()?, metadata.len()));
480 }
481 }
482 Ok(())
483}
484
485fn remove_empty_dirs(dir: &Path, root: &Path) {
487 let Ok(entries) = std::fs::read_dir(dir) else {
488 return;
489 };
490 for entry in entries.flatten() {
491 if std::fs::symlink_metadata(entry.path()).is_ok_and(|metadata| metadata.is_dir()) {
492 remove_empty_dirs(&entry.path(), root);
493 }
494 }
495 if dir != root {
496 let _ = std::fs::remove_dir(dir);
498 }
499}
500
501pub fn parse_age(value: &str) -> Result<Duration, String> {
503 let value = value.trim();
504 let (number, unit) = match value.find(|c: char| !c.is_ascii_digit()) {
505 Some(index) => value.split_at(index),
506 None => (value, "s"),
507 };
508 let number: u64 = number
509 .parse()
510 .map_err(|_| format!("invalid age {value:?}; use a number with s, m, h, or d"))?;
511 let unit = match unit {
512 "s" => 1,
513 "m" => 60,
514 "h" => 3600,
515 "d" => 86400,
516 _ => {
517 return Err(format!(
518 "invalid age {value:?}; use a number with s, m, h, or d"
519 ));
520 }
521 };
522 number
523 .checked_mul(unit)
524 .map(Duration::from_secs)
525 .ok_or_else(|| format!("age {value:?} is too large"))
526}
527
528#[cfg(test)]
529mod tests {
530 use super::*;
531
532 fn store(max: usize, idle: Duration, markers: Option<&Path>) -> Arc<ConversationStore> {
533 Arc::new(ConversationStore::new(
534 ConversationLimits { max, idle },
535 markers.map(Path::to_path_buf),
536 ))
537 }
538
539 const DAY: Duration = Duration::from_secs(86400);
540
541 #[test]
542 fn handles_are_issued_per_agent_and_vendor_ids_are_not_handles() {
543 assert!(is_handle("codex-2"));
544 assert!(is_handle("pi-10"));
545 for not_handle in [
546 "01a0cd5a-7195-7b31-a503-e235d5da7b45",
547 "b514bbf5-a5b7-4bbe-83f9-5824ab41c35c",
548 "codex",
549 "codex-",
550 "-1",
551 "Codex-1",
552 "codex-1a",
553 "../codex-1",
554 ] {
555 assert!(!is_handle(not_handle), "{not_handle}");
556 }
557 let store = store(8, DAY, None);
558 let cwd = Path::new("/w");
559 let first = store.begin("codex", None, cwd, false).unwrap();
560 assert_eq!((first.handle.as_str(), first.turn), ("codex-1", 1));
561 assert_eq!(first.vendor, None);
562 assert_eq!(
563 first.finish(Some("t-1".into()), true).as_deref(),
564 Some("codex-1")
565 );
566 let claude = store.begin("claude", None, cwd, true).unwrap();
567 assert_eq!(claude.handle, "claude-1");
568 assert!(claude.vendor.is_some());
569 claude.finish(None, true);
570 let second = store.begin("codex", None, cwd, false).unwrap();
571 assert_eq!(second.handle, "codex-2");
572 }
573
574 #[test]
575 fn continuing_pins_agent_and_cwd_and_counts_turns() {
576 let store = store(8, DAY, None);
577 let cwd = Path::new("/w/scv");
578 let turn = store.begin("codex", None, cwd, false).unwrap();
579 turn.finish(Some("thread-a".into()), true);
580 let next = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
581 assert_eq!(next.turn, 2);
582 assert_eq!(next.vendor.as_deref(), Some("thread-a"));
583 let busy = store
585 .begin("codex", Some("codex-1"), cwd, false)
586 .unwrap_err();
587 assert!(busy.0.starts_with("session busy"), "{}", busy.0);
588 next.finish(Some("thread-a".into()), true);
589 let moved = store
590 .begin("codex", Some("codex-1"), Path::new("/w/other"), false)
591 .unwrap_err();
592 assert!(moved.0.contains("runs in"), "{}", moved.0);
593 let other_agent = store
594 .begin("claude", Some("codex-1"), cwd, false)
595 .unwrap_err();
596 assert!(
597 other_agent.0.contains("belongs to agent_codex"),
598 "{}",
599 other_agent.0
600 );
601 let third = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
602 assert_eq!(third.turn, 3);
603 }
604
605 #[test]
606 fn vendor_ids_and_unknown_handles_are_rejected() {
607 let store = store(8, DAY, None);
608 let cwd = Path::new("/w");
609 store
610 .begin("codex", None, cwd, false)
611 .unwrap()
612 .finish(Some("01a0cd5a-7195-7b31".into()), true);
613 let vendor = store
614 .begin("codex", Some("01a0cd5a-7195-7b31"), cwd, false)
615 .unwrap_err();
616 assert!(
617 vendor.0.contains("not a conversation handle"),
618 "{}",
619 vendor.0
620 );
621 let unknown = store
622 .begin("codex", Some("codex-9"), cwd, false)
623 .unwrap_err();
624 assert!(
625 unknown.0.contains("unknown in this session"),
626 "{}",
627 unknown.0
628 );
629 let other = super::tests::store(8, DAY, None);
631 assert!(other.begin("codex", Some("codex-1"), cwd, false).is_err());
632 }
633
634 #[test]
635 fn timed_out_turns_stay_resumable_but_failed_first_turns_are_forgotten() {
636 let store = store(8, DAY, None);
637 let cwd = Path::new("/w");
638 let turn = store.begin("codex", None, cwd, false).unwrap();
640 assert_eq!(
641 turn.finish(Some("t".into()), false).as_deref(),
642 Some("codex-1")
643 );
644 assert_eq!(
645 store
646 .begin("codex", Some("codex-1"), cwd, false)
647 .unwrap()
648 .turn,
649 2
650 );
651 let failed = store.begin("codex", None, cwd, false).unwrap();
653 assert_eq!(failed.finish(None, false), None);
654 assert!(!store.handles().contains(&"codex-2".to_owned()));
655 drop(store.begin("codex", None, cwd, false).unwrap());
657 assert_eq!(store.handles(), vec!["codex-1".to_owned()]);
658 }
659
660 #[test]
661 fn limits_forget_the_least_recently_used_and_idle_conversations() {
662 let store = store(2, DAY, None);
663 let cwd = Path::new("/w");
664 for id in ["a", "b"] {
665 store
666 .begin("codex", None, cwd, false)
667 .unwrap()
668 .finish(Some(id.into()), true);
669 }
670 store
672 .begin("codex", Some("codex-1"), cwd, false)
673 .unwrap()
674 .finish(Some("a".into()), true);
675 store
676 .begin("codex", None, cwd, false)
677 .unwrap()
678 .finish(Some("c".into()), true);
679 assert_eq!(
680 store.handles(),
681 vec!["codex-1".to_owned(), "codex-3".to_owned()]
682 );
683 let busy_a = store.begin("codex", Some("codex-1"), cwd, false).unwrap();
685 let busy_b = store.begin("codex", Some("codex-3"), cwd, false).unwrap();
686 assert!(store.begin("codex", None, cwd, false).is_err());
687 drop((busy_a, busy_b));
688
689 let idle = super::tests::store(8, Duration::ZERO, None);
690 idle.begin("pi", None, cwd, true)
691 .unwrap()
692 .finish(None, true);
693 let expired = idle.begin("pi", Some("pi-1"), cwd, true).unwrap_err();
694 assert!(
695 expired.0.contains("forgotten after 0 seconds idle"),
696 "{}",
697 expired.0
698 );
699 }
700
701 #[test]
702 fn markers_follow_the_conversation_and_gc_keeps_live_transcripts() {
703 let home = tempfile::tempdir().unwrap();
704 let markers = home.path().join("run/conversations");
705 let adapter = home.path().join("adapters/codex");
706 let day = adapter.join("sessions/2026/01/02");
707 std::fs::create_dir_all(&day).unwrap();
708 let old = SystemTime::now() - Duration::from_secs(10 * 86400);
709 let transcript = |id: &str| {
710 let path = day.join(format!("rollout-2026-01-02T00-00-00-{id}.jsonl"));
711 std::fs::write(&path, "{}\n").unwrap();
712 std::fs::File::options()
713 .write(true)
714 .open(&path)
715 .unwrap()
716 .set_modified(old)
717 .unwrap();
718 path
719 };
720 let live = transcript("live-id");
721 let stale = transcript("stale-id");
722 let recent = day.join("rollout-recent-id.jsonl");
723 std::fs::write(&recent, "{}\n").unwrap();
724 let outside = tempfile::NamedTempFile::new().unwrap();
726 std::os::unix::fs::symlink(outside.path(), day.join("link.jsonl")).unwrap();
727
728 let store = store(8, DAY, Some(&markers));
729 store
730 .begin("codex", None, Path::new("/w"), false)
731 .unwrap()
732 .finish(Some("live-id".into()), true);
733 assert!(markers.join("live-id.json").is_file());
734 write_private_json(
736 &markers,
737 "stale-id.json",
738 &Marker {
739 owner: ProcessIdentity {
740 pid: u32::MAX - 1,
741 start_time: 1,
742 },
743 agent: "codex".into(),
744 handle: "codex-9".into(),
745 },
746 )
747 .unwrap();
748 let files = ConversationFiles {
749 dir: "sessions",
750 extension: "jsonl",
751 };
752 let dry = collect_garbage(&adapter, files, &markers, DAY, true).unwrap();
753 assert_eq!(dry.removed, vec![stale.clone()]);
754 assert_eq!(dry.kept_live, 1);
755 assert!(stale.exists() && markers.join("stale-id.json").exists());
756 let report = collect_garbage(&adapter, files, &markers, DAY, false).unwrap();
757 assert_eq!(report.removed, vec![stale.clone()]);
758 assert!(!stale.exists() && live.exists() && recent.exists());
759 assert!(outside.path().exists());
760 assert!(!markers.join("stale-id.json").exists());
761 drop(store);
763 assert!(!markers.join("live-id.json").exists());
764 let report = collect_garbage(&adapter, files, &markers, Duration::ZERO, false).unwrap();
766 assert_eq!(report.removed, vec![live]);
767 assert!(recent.exists());
768 }
769
770 #[test]
771 fn ages_parse_with_units() {
772 assert_eq!(parse_age("30d"), Ok(Duration::from_secs(30 * 86400)));
773 assert_eq!(parse_age("12h"), Ok(Duration::from_secs(12 * 3600)));
774 assert_eq!(parse_age("90m"), Ok(Duration::from_secs(5400)));
775 assert_eq!(parse_age("45"), Ok(Duration::from_secs(45)));
776 assert!(parse_age("3w").is_err());
777 assert!(parse_age("d").is_err());
778 assert!(parse_age("-1d").is_err());
779 }
780}