1use crate::exec::events::ThreadEvent;
2use crate::llm::provider::Message;
3use crate::utils::session_archive::{
4 SessionArchive, SessionArchiveMetadata, SessionForkMode, SessionListing, find_session_by_identifier,
5 list_recent_sessions, reserve_session_archive_identifier, session_listing_matches_workspace,
6};
7use crate::utils::session_debug::runtime_archive_session_id;
8use anyhow::{Result, anyhow};
9use parking_lot::Mutex;
10use serde::{Deserialize, Serialize};
11use std::collections::VecDeque;
12use std::path::{Path, PathBuf};
13use std::sync::Arc;
14use uuid::Uuid;
15use vtcode_macros::StringNewtype;
16
17const DEFAULT_EVENT_BUFFER_CAPACITY: usize = 512;
18
19#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, StringNewtype)]
21pub struct ThreadId(String);
22
23#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, StringNewtype)]
25pub struct SubmissionId(String);
26
27impl SubmissionId {
28 pub fn generate() -> Self {
30 Self(format!("sub-{}", Uuid::new_v4()))
31 }
32}
33
34impl Default for SubmissionId {
35 fn default() -> Self {
36 Self::generate()
37 }
38}
39
40#[derive(Debug, Clone)]
41pub struct ThreadEventRecord {
42 pub sequence: u64,
43 pub thread_id: ThreadId,
44 pub submission_id: Option<SubmissionId>,
45 pub turn_id: Option<String>,
46 pub event: ThreadEvent,
47}
48
49#[derive(Debug, Clone)]
50pub struct ThreadSnapshot {
51 pub thread_id: ThreadId,
52 pub metadata: Option<SessionArchiveMetadata>,
53 pub archive_listing: Option<SessionListing>,
54 pub messages: Vec<Message>,
55 pub loaded_skills: Vec<String>,
56 pub turn_in_flight: bool,
57}
58
59#[derive(Debug, Clone)]
60pub struct ThreadBootstrap {
61 pub metadata: Option<SessionArchiveMetadata>,
62 pub archive_listing: Option<SessionListing>,
63 pub messages: Vec<Message>,
64 pub loaded_skills: Vec<String>,
65}
66
67impl ThreadBootstrap {
68 pub fn new(metadata: Option<SessionArchiveMetadata>) -> Self {
69 Self {
70 metadata,
71 archive_listing: None,
72 messages: Vec::new(),
73 loaded_skills: Vec::new(),
74 }
75 }
76
77 pub fn from_listing(listing: SessionListing) -> Self {
78 Self {
79 metadata: Some(listing.snapshot.metadata.clone()),
80 messages: messages_from_session_listing(&listing),
81 loaded_skills: loaded_skills_from_session_listing(&listing),
82 archive_listing: Some(listing),
83 }
84 }
85
86 pub fn from_snapshot(snapshot: ThreadSnapshot) -> Self {
87 Self {
88 metadata: snapshot.metadata,
89 archive_listing: snapshot.archive_listing,
90 messages: snapshot.messages,
91 loaded_skills: snapshot.loaded_skills,
92 }
93 }
94
95 pub fn with_messages(mut self, messages: Vec<Message>) -> Self {
96 self.messages = messages;
97 self
98 }
99
100 pub fn with_loaded_skills(mut self, loaded_skills: Vec<String>) -> Self {
101 self.loaded_skills = loaded_skills;
102 self
103 }
104}
105
106#[derive(Debug, Clone, PartialEq, Eq)]
107pub enum SessionQueryScope {
108 CurrentWorkspace(PathBuf),
109 All,
110}
111
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub enum ArchivedSessionIntent {
114 ResumeInPlace,
115 ForkNewArchive {
116 custom_suffix: Option<String>,
117 summarize: bool,
118 },
119}
120
121#[derive(Debug, Clone)]
122pub struct PreparedArchivedSession {
123 pub source: SessionListing,
124 pub workspace: PathBuf,
125 pub bootstrap: ThreadBootstrap,
126 pub thread_id: String,
127 pub archive: SessionArchive,
128}
129
130#[derive(Default)]
131struct ThreadEventStore {
132 capacity: usize,
133 next_sequence: u64,
134 events: VecDeque<ThreadEventRecord>,
135}
136
137impl ThreadEventStore {
138 fn with_capacity(capacity: usize) -> Self {
139 Self { capacity: capacity.max(1), ..Self::default() }
140 }
141
142 fn push(
143 &mut self,
144 thread_id: &ThreadId,
145 submission_id: Option<SubmissionId>,
146 turn_id: Option<String>,
147 event: ThreadEvent,
148 ) {
149 let record = ThreadEventRecord {
150 sequence: self.next_sequence,
151 thread_id: thread_id.clone(),
152 submission_id,
153 turn_id,
154 event,
155 };
156 self.next_sequence = self.next_sequence.saturating_add(1);
157
158 if self.events.len() >= self.capacity {
159 self.events.pop_front();
160 }
161 self.events.push_back(record);
162 }
163
164 fn snapshot(&self) -> Vec<ThreadEventRecord> {
165 self.events.iter().cloned().collect()
166 }
167}
168
169struct ThreadSessionState {
170 thread_id: ThreadId,
171 metadata: Option<SessionArchiveMetadata>,
172 archive_listing: Option<SessionListing>,
173 messages: Vec<Message>,
174 loaded_skills: Vec<String>,
175 turn_in_flight: bool,
176}
177
178impl ThreadSessionState {
179 fn snapshot(&self) -> ThreadSnapshot {
180 ThreadSnapshot {
181 thread_id: self.thread_id.clone(),
182 metadata: self.metadata.clone(),
183 archive_listing: self.archive_listing.clone(),
184 messages: self.messages.clone(),
185 loaded_skills: self.loaded_skills.clone(),
186 turn_in_flight: self.turn_in_flight,
187 }
188 }
189}
190
191#[derive(Clone)]
192pub struct ThreadRuntimeHandle {
193 inner: Arc<ThreadRuntimeInner>,
194}
195
196struct ThreadRuntimeInner {
197 session: Mutex<ThreadSessionState>,
198 event_store: Mutex<ThreadEventStore>,
199}
200
201impl ThreadRuntimeHandle {
202 fn new(thread_id: ThreadId, bootstrap: ThreadBootstrap, event_capacity: usize) -> Self {
203 let session = ThreadSessionState {
204 thread_id,
205 metadata: bootstrap.metadata,
206 archive_listing: bootstrap.archive_listing,
207 messages: bootstrap.messages,
208 loaded_skills: bootstrap.loaded_skills,
209 turn_in_flight: false,
210 };
211
212 Self {
213 inner: Arc::new(ThreadRuntimeInner {
214 session: Mutex::new(session),
215 event_store: Mutex::new(ThreadEventStore::with_capacity(event_capacity)),
216 }),
217 }
218 }
219
220 pub fn thread_id(&self) -> ThreadId {
221 self.inner.session.lock().thread_id.clone()
222 }
223
224 pub fn snapshot(&self) -> ThreadSnapshot {
225 self.inner.session.lock().snapshot()
226 }
227
228 pub fn metadata(&self) -> Option<SessionArchiveMetadata> {
229 self.inner.session.lock().metadata.clone()
230 }
231
232 pub fn replace_metadata(&self, metadata: Option<SessionArchiveMetadata>) {
233 self.inner.session.lock().metadata = metadata;
234 }
235
236 pub fn archive_listing(&self) -> Option<SessionListing> {
237 self.inner.session.lock().archive_listing.clone()
238 }
239
240 pub fn messages(&self) -> Vec<Message> {
241 self.inner.session.lock().messages.clone()
242 }
243
244 pub fn replace_messages(&self, messages: Vec<Message>) {
245 self.inner.session.lock().messages = messages;
246 }
247
248 pub fn append_message(&self, message: Message) {
249 self.inner.session.lock().messages.push(message);
250 }
251
252 pub fn begin_turn(&self) -> Result<SubmissionId> {
253 let mut session = self.inner.session.lock();
254 if session.turn_in_flight {
255 return Err(anyhow!("thread '{}' already has an in-flight turn", session.thread_id));
256 }
257
258 session.turn_in_flight = true;
259 Ok(SubmissionId::generate())
260 }
261
262 pub fn finish_turn(&self) {
263 self.inner.session.lock().turn_in_flight = false;
264 }
265
266 pub fn record_event(&self, submission_id: Option<SubmissionId>, turn_id: Option<String>, event: ThreadEvent) {
267 let thread_id = self.thread_id();
268 self.inner.event_store.lock().push(&thread_id, submission_id, turn_id, event);
269 }
270
271 pub fn replay_recent(&self) -> Vec<ThreadEventRecord> {
272 self.inner.event_store.lock().snapshot()
273 }
274
275 pub fn recent_events(&self) -> Vec<ThreadEvent> {
276 self.replay_recent().into_iter().map(|record| record.event).collect()
277 }
278}
279
280#[derive(Clone)]
281pub struct ThreadManager {
282 event_buffer_capacity: usize,
283}
284
285impl Default for ThreadManager {
286 fn default() -> Self {
287 Self::new()
288 }
289}
290
291impl ThreadManager {
292 pub fn new() -> Self {
293 Self {
294 event_buffer_capacity: DEFAULT_EVENT_BUFFER_CAPACITY,
295 }
296 }
297
298 pub fn with_event_buffer_capacity(event_buffer_capacity: usize) -> Self {
299 Self {
300 event_buffer_capacity: event_buffer_capacity.max(1),
301 }
302 }
303
304 pub fn start_thread_with_identifier(
305 &self,
306 identifier: impl Into<String>,
307 bootstrap: ThreadBootstrap,
308 ) -> ThreadRuntimeHandle {
309 ThreadRuntimeHandle::new(ThreadId::new(identifier.into()), bootstrap, self.event_buffer_capacity)
310 }
311
312 pub async fn start_thread(
313 &self,
314 workspace_label: &str,
315 custom_suffix: Option<String>,
316 bootstrap: ThreadBootstrap,
317 ) -> Result<ThreadRuntimeHandle> {
318 let identifier = reserve_session_archive_identifier(workspace_label, custom_suffix).await?;
319 Ok(self.start_thread_with_identifier(identifier, bootstrap))
320 }
321
322 pub async fn resume_thread(&self, identifier: &str) -> Result<Option<ThreadRuntimeHandle>> {
323 let listing = find_session_by_identifier(identifier).await?;
324 Ok(listing.map(|listing| {
325 self.start_thread_with_identifier(listing.identifier(), ThreadBootstrap::from_listing(listing))
326 }))
327 }
328}
329
330pub async fn list_recent_sessions_in_scope(limit: usize, scope: &SessionQueryScope) -> Result<Vec<SessionListing>> {
331 let mut listings = list_recent_sessions(limit.saturating_mul(4).max(limit)).await?;
332 if let SessionQueryScope::CurrentWorkspace(workspace) = scope {
333 listings.retain(|listing| session_listing_matches_workspace(listing, workspace));
334 }
335 let current_id = runtime_archive_session_id();
336 listings.retain(|listing| {
337 if listing.snapshot.total_messages == 0 {
338 return false;
339 }
340 if let Some(ref current) = current_id
341 && listing.identifier() == *current
342 {
343 return false;
344 }
345 true
346 });
347 listings.truncate(limit);
348 Ok(listings)
349}
350
351pub async fn prepare_archived_session(
352 source: SessionListing,
353 workspace: PathBuf,
354 metadata: SessionArchiveMetadata,
355 intent: ArchivedSessionIntent,
356 reserved_identifier: Option<String>,
357) -> Result<PreparedArchivedSession> {
358 let mut metadata = preserve_prompt_cache_lineage_if_compatible(metadata, &source.snapshot.metadata);
359 metadata.continuation_metadata = source.snapshot.metadata.continuation_metadata.clone();
360 if metadata.primary_agent.is_none() {
363 metadata.primary_agent = source.snapshot.metadata.primary_agent.clone();
364 }
365 let mut bootstrap = ThreadBootstrap::from_listing(source.clone());
366 bootstrap.metadata = Some(metadata.clone());
367
368 let thread_id = match &intent {
369 ArchivedSessionIntent::ResumeInPlace => source.identifier(),
370 ArchivedSessionIntent::ForkNewArchive { custom_suffix, .. } => {
371 if let Some(identifier) = reserved_identifier {
372 identifier
373 } else {
374 reserve_session_archive_identifier(&metadata.workspace_label, custom_suffix.clone()).await?
375 }
376 }
377 };
378
379 if let ArchivedSessionIntent::ForkNewArchive { summarize, .. } = &intent {
380 metadata.parent_session_id = Some(source.identifier());
381 metadata.fork_mode = Some(if *summarize {
382 SessionForkMode::Summarized
383 } else {
384 SessionForkMode::FullCopy
385 });
386 bootstrap.metadata = Some(metadata.clone());
387 }
388
389 let archive = match intent {
390 ArchivedSessionIntent::ResumeInPlace => SessionArchive::resume_from_listing(&source, metadata),
391 ArchivedSessionIntent::ForkNewArchive { .. } => {
392 SessionArchive::new_with_identifier(metadata, thread_id.clone()).await?
393 }
394 };
395
396 Ok(PreparedArchivedSession { source, workspace, bootstrap, thread_id, archive })
397}
398
399fn preserve_prompt_cache_lineage_if_compatible(
400 mut metadata: SessionArchiveMetadata,
401 source: &SessionArchiveMetadata,
402) -> SessionArchiveMetadata {
403 let is_compatible = metadata.workspace_path == source.workspace_path
404 && metadata.provider == source.provider
405 && metadata.model == source.model;
406 if is_compatible && let Some(lineage_id) = source.prompt_cache_lineage_id.as_ref() {
407 metadata.prompt_cache_lineage_id = Some(lineage_id.clone());
408 }
409 metadata
410}
411
412pub fn messages_from_session_listing(listing: &SessionListing) -> Vec<Message> {
413 if !listing.snapshot.messages.is_empty() {
414 listing.snapshot.messages.iter().map(Message::from).collect()
415 } else if let Some(progress) = &listing.snapshot.progress
416 && !progress.recent_messages.is_empty()
417 {
418 progress.recent_messages.iter().map(Message::from).collect()
419 } else {
420 Vec::new()
421 }
422}
423
424pub fn loaded_skills_from_session_listing(listing: &SessionListing) -> Vec<String> {
425 listing
426 .snapshot
427 .progress
428 .as_ref()
429 .map(|progress| progress.loaded_skills.clone())
430 .filter(|skills| !skills.is_empty())
431 .unwrap_or_else(|| listing.snapshot.metadata.loaded_skills.clone())
432}
433
434pub fn build_thread_archive_metadata(
435 workspace: &Path,
436 model: &str,
437 provider: &str,
438 theme: &str,
439 reasoning_effort: &str,
440) -> SessionArchiveMetadata {
441 let workspace_label = workspace.file_name().and_then(|value| value.to_str()).unwrap_or("workspace");
442
443 SessionArchiveMetadata::new(
444 workspace_label,
445 workspace.to_string_lossy().to_string(),
446 model,
447 provider,
448 theme,
449 reasoning_effort,
450 )
451 .ensure_prompt_cache_lineage_id()
452}
453
454#[cfg(test)]
455mod tests {
456 use super::*;
457 use crate::exec::events::{ThreadEvent, ThreadStartedEvent};
458 use crate::llm::provider::MessageRole;
459 use crate::utils::session_archive::{
460 SessionArchiveMetadata, SessionMessage, SessionProgress, SessionSnapshot,
461 clear_sessions_dir_override_for_tests, override_sessions_dir_for_tests,
462 };
463 use chrono::Utc;
464 use std::sync::{LazyLock, Mutex};
465 use tempfile::TempDir;
466
467 static SESSION_DIR_TEST_GUARD: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::new(()));
468
469 #[test]
470 fn event_store_evicts_old_records() {
471 let manager = ThreadManager::with_event_buffer_capacity(2);
472 let handle = manager.start_thread_with_identifier("thread-1", ThreadBootstrap::new(None));
473
474 handle.record_event(
475 None,
476 None,
477 ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread-1".to_string() }),
478 );
479 handle.record_event(
480 None,
481 Some("turn-1".to_string()),
482 ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread-1-turn-1".to_string() }),
483 );
484 handle.record_event(
485 None,
486 Some("turn-2".to_string()),
487 ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread-1-turn-2".to_string() }),
488 );
489
490 let records = handle.replay_recent();
491 assert_eq!(records.len(), 2);
492 assert_eq!(records[0].sequence, 1);
493 assert_eq!(records[1].sequence, 2);
494 }
495
496 #[test]
497 fn start_thread_with_identifier_preserves_message_history() {
498 let manager = ThreadManager::new();
499 let bootstrap = ThreadBootstrap::new(None)
500 .with_messages(vec![Message::user("hello".to_string())])
501 .with_loaded_skills(vec!["repo-skill".to_string()]);
502 let handle = manager.start_thread_with_identifier("thread-123", bootstrap);
503
504 assert_eq!(handle.thread_id().as_str(), "thread-123");
505 let snapshot = handle.snapshot();
506 assert_eq!(snapshot.messages.len(), 1);
507 assert_eq!(snapshot.loaded_skills, vec!["repo-skill".to_string()]);
508 }
509
510 #[test]
511 fn submit_enforces_single_in_flight_turn() {
512 let manager = ThreadManager::new();
513 let handle = manager.start_thread_with_identifier("thread-123", ThreadBootstrap::new(None));
514
515 let _first = handle.begin_turn().expect("first turn");
516 let err = handle.begin_turn().expect_err("second turn should fail");
517 assert!(err.to_string().contains("in-flight turn"));
518 handle.finish_turn();
519 handle.begin_turn().expect("turn after finish");
520 }
521
522 #[test]
523 #[serial_test::serial(session_dir_override)]
524 fn list_recent_sessions_in_scope_filters_by_workspace() {
525 let _guard = SESSION_DIR_TEST_GUARD.lock().expect("session dir test guard");
526 let tmp = TempDir::new().expect("temp dir");
527 override_sessions_dir_for_tests(tmp.path());
528
529 let listing = SessionListing {
530 path: tmp.path().join("session-alpha.json"),
531 snapshot: SessionSnapshot {
532 metadata: SessionArchiveMetadata::new(
533 "ws",
534 tmp.path().join("workspace").display().to_string(),
535 "model",
536 "provider",
537 "theme",
538 "medium",
539 ),
540 started_at: Utc::now(),
541 ended_at: Utc::now(),
542 total_messages: 1,
543 distinct_tools: Vec::new(),
544 transcript: Vec::new(),
545 messages: vec![SessionMessage::new(MessageRole::User, "hello")],
546 progress: None,
547 error_logs: Vec::new(),
548 },
549 };
550 std::fs::write(&listing.path, serde_json::to_string(&listing.snapshot).expect("serialize snapshot"))
551 .expect("write listing");
552
553 let runtime = tokio::runtime::Runtime::new().expect("runtime");
554 let filtered = runtime
555 .block_on(list_recent_sessions_in_scope(
556 5,
557 &SessionQueryScope::CurrentWorkspace(tmp.path().join("workspace")),
558 ))
559 .expect("filter by workspace");
560 let all = runtime
561 .block_on(list_recent_sessions_in_scope(5, &SessionQueryScope::All))
562 .expect("list all");
563
564 clear_sessions_dir_override_for_tests();
565
566 assert_eq!(filtered.len(), 1);
567 assert_eq!(all.len(), 1);
568 }
569
570 #[test]
571 #[serial_test::serial(session_dir_override)]
572 fn prepare_archived_session_resume_reuses_source_identifier_and_archive() {
573 let _guard = SESSION_DIR_TEST_GUARD.lock().expect("session dir test guard");
574 let tmp = TempDir::new().expect("temp dir");
575 override_sessions_dir_for_tests(tmp.path());
576
577 let listing = SessionListing {
578 path: tmp.path().join("session-source.json"),
579 snapshot: SessionSnapshot {
580 metadata: SessionArchiveMetadata::new(
581 "ws",
582 tmp.path().join("workspace").display().to_string(),
583 "old-model",
584 "old-provider",
585 "old-theme",
586 "medium",
587 )
588 .with_primary_agent("plan"),
589 started_at: Utc::now(),
590 ended_at: Utc::now(),
591 total_messages: 2,
592 distinct_tools: vec!["tool_a".to_string()],
593 transcript: Vec::new(),
594 messages: vec![SessionMessage::new(MessageRole::User, "hello")],
595 progress: Some(Box::new(SessionProgress {
596 turn_number: 1,
597 recent_messages: vec![SessionMessage::new(MessageRole::Assistant, "recent")],
598 tool_summaries: Vec::new(),
599 token_usage: None,
600 max_context_tokens: None,
601 loaded_skills: vec!["skill_a".to_string()],
602 turn_diagnostics: None,
603 })),
604 error_logs: Vec::new(),
605 },
606 };
607
608 let runtime = tokio::runtime::Runtime::new().expect("runtime");
609 let prepared = runtime
610 .block_on(prepare_archived_session(
611 listing.clone(),
612 tmp.path().join("workspace"),
613 SessionArchiveMetadata::new(
614 "ws",
615 tmp.path().join("workspace").display().to_string(),
616 "new-model",
617 "new-provider",
618 "new-theme",
619 "high",
620 ),
621 ArchivedSessionIntent::ResumeInPlace,
622 Some("should-not-be-used".to_string()),
623 ))
624 .expect("prepare resume");
625
626 clear_sessions_dir_override_for_tests();
627
628 assert_eq!(prepared.thread_id, listing.identifier());
629 assert_eq!(prepared.archive.path(), listing.path.as_path());
630 assert_eq!(prepared.bootstrap.messages[0].content.as_text(), "hello");
631 assert_eq!(prepared.bootstrap.loaded_skills, vec!["skill_a".to_string()]);
632 assert_eq!(prepared.bootstrap.metadata.as_ref().expect("metadata").model, "new-model");
633 assert_eq!(prepared.bootstrap.metadata.as_ref().expect("metadata").primary_agent.as_deref(), Some("plan"));
635 }
636
637 #[test]
638 #[serial_test::serial(session_dir_override)]
639 fn prepare_archived_session_fork_uses_new_identifier_and_preserves_history() {
640 let _guard = SESSION_DIR_TEST_GUARD.lock().expect("session dir test guard");
641 let tmp = TempDir::new().expect("temp dir");
642 override_sessions_dir_for_tests(tmp.path());
643
644 let listing = SessionListing {
645 path: tmp.path().join("session-source.json"),
646 snapshot: SessionSnapshot {
647 metadata: SessionArchiveMetadata::new(
648 "ws",
649 tmp.path().join("workspace").display().to_string(),
650 "model",
651 "provider",
652 "theme",
653 "medium",
654 ),
655 started_at: Utc::now(),
656 ended_at: Utc::now(),
657 total_messages: 1,
658 distinct_tools: Vec::new(),
659 transcript: Vec::new(),
660 messages: vec![SessionMessage::new(MessageRole::User, "hello")],
661 progress: None,
662 error_logs: Vec::new(),
663 },
664 };
665
666 let runtime = tokio::runtime::Runtime::new().expect("runtime");
667 let prepared = runtime
668 .block_on(prepare_archived_session(
669 listing.clone(),
670 tmp.path().join("workspace"),
671 SessionArchiveMetadata::new(
672 "ws",
673 tmp.path().join("workspace").display().to_string(),
674 "model",
675 "provider",
676 "theme",
677 "medium",
678 ),
679 ArchivedSessionIntent::ForkNewArchive {
680 custom_suffix: Some("branch".to_string()),
681 summarize: false,
682 },
683 Some("session-forked".to_string()),
684 ))
685 .expect("prepare fork");
686
687 clear_sessions_dir_override_for_tests();
688
689 assert_eq!(prepared.thread_id, "session-forked");
690 assert_ne!(prepared.archive.path(), listing.path.as_path());
691 assert!(prepared.archive.path().ends_with(Path::new("session-forked.json")));
692 assert_eq!(prepared.bootstrap.messages[0].content.as_text(), "hello");
693 assert_eq!(
694 prepared
695 .bootstrap
696 .metadata
697 .as_ref()
698 .and_then(|metadata| metadata.parent_session_id.as_deref()),
699 Some("session-source")
700 );
701 assert_eq!(
702 prepared.bootstrap.metadata.as_ref().and_then(|metadata| metadata.fork_mode),
703 Some(SessionForkMode::FullCopy)
704 );
705 }
706
707 #[test]
708 #[serial_test::serial(session_dir_override)]
709 fn prepare_archived_session_preserves_prompt_cache_lineage_when_compatible() {
710 let _guard = SESSION_DIR_TEST_GUARD.lock().expect("session dir test guard");
711 let tmp = TempDir::new().expect("temp dir");
712 override_sessions_dir_for_tests(tmp.path());
713
714 let listing = SessionListing {
715 path: tmp.path().join("session-source.json"),
716 snapshot: SessionSnapshot {
717 metadata: SessionArchiveMetadata::new(
718 "ws",
719 tmp.path().join("workspace").display().to_string(),
720 "model",
721 "provider",
722 "theme",
723 "medium",
724 )
725 .with_prompt_cache_lineage_id("lineage-source"),
726 started_at: Utc::now(),
727 ended_at: Utc::now(),
728 total_messages: 1,
729 distinct_tools: Vec::new(),
730 transcript: Vec::new(),
731 messages: vec![SessionMessage::new(MessageRole::User, "hello")],
732 progress: None,
733 error_logs: Vec::new(),
734 },
735 };
736
737 let runtime = tokio::runtime::Runtime::new().expect("runtime");
738 let prepared = runtime
739 .block_on(prepare_archived_session(
740 listing,
741 tmp.path().join("workspace"),
742 SessionArchiveMetadata::new(
743 "ws",
744 tmp.path().join("workspace").display().to_string(),
745 "model",
746 "provider",
747 "theme",
748 "medium",
749 )
750 .with_prompt_cache_lineage_id("lineage-new"),
751 ArchivedSessionIntent::ResumeInPlace,
752 None,
753 ))
754 .expect("prepare resume");
755
756 clear_sessions_dir_override_for_tests();
757
758 assert_eq!(
759 prepared
760 .bootstrap
761 .metadata
762 .as_ref()
763 .and_then(|metadata| metadata.prompt_cache_lineage_id.as_deref()),
764 Some("lineage-source")
765 );
766 }
767
768 #[test]
769 #[serial_test::serial(session_dir_override)]
770 fn prepare_archived_session_resets_prompt_cache_lineage_on_model_change() {
771 let _guard = SESSION_DIR_TEST_GUARD.lock().expect("session dir test guard");
772 let tmp = TempDir::new().expect("temp dir");
773 override_sessions_dir_for_tests(tmp.path());
774
775 let listing = SessionListing {
776 path: tmp.path().join("session-source.json"),
777 snapshot: SessionSnapshot {
778 metadata: SessionArchiveMetadata::new(
779 "ws",
780 tmp.path().join("workspace").display().to_string(),
781 "model-a",
782 "provider",
783 "theme",
784 "medium",
785 )
786 .with_prompt_cache_lineage_id("lineage-source"),
787 started_at: Utc::now(),
788 ended_at: Utc::now(),
789 total_messages: 1,
790 distinct_tools: Vec::new(),
791 transcript: Vec::new(),
792 messages: vec![SessionMessage::new(MessageRole::User, "hello")],
793 progress: None,
794 error_logs: Vec::new(),
795 },
796 };
797
798 let runtime = tokio::runtime::Runtime::new().expect("runtime");
799 let prepared = runtime
800 .block_on(prepare_archived_session(
801 listing,
802 tmp.path().join("workspace"),
803 SessionArchiveMetadata::new(
804 "ws",
805 tmp.path().join("workspace").display().to_string(),
806 "model-b",
807 "provider",
808 "theme",
809 "medium",
810 )
811 .with_prompt_cache_lineage_id("lineage-new"),
812 ArchivedSessionIntent::ResumeInPlace,
813 None,
814 ))
815 .expect("prepare resume");
816
817 clear_sessions_dir_override_for_tests();
818
819 assert_eq!(
820 prepared
821 .bootstrap
822 .metadata
823 .as_ref()
824 .and_then(|metadata| metadata.prompt_cache_lineage_id.as_deref()),
825 Some("lineage-new")
826 );
827 }
828
829 #[test]
830 fn messages_from_session_listing_preserves_assistant_phases_from_progress() {
831 let listing = SessionListing {
832 path: PathBuf::from("session.json"),
833 snapshot: SessionSnapshot {
834 metadata: SessionArchiveMetadata::new("ws", "/tmp/ws", "gpt-5.6-sol", "openai", "theme", "medium"),
835 started_at: Utc::now(),
836 ended_at: Utc::now(),
837 total_messages: 2,
838 distinct_tools: Vec::new(),
839 transcript: Vec::new(),
840 messages: Vec::new(),
841 progress: Some(Box::new(SessionProgress {
842 turn_number: 2,
843 recent_messages: vec![
844 SessionMessage::from(
845 &Message::assistant("Working".to_string())
846 .with_phase(Some(crate::llm::provider::AssistantPhase::Commentary)),
847 ),
848 SessionMessage::from(
849 &Message::assistant("Done".to_string())
850 .with_phase(Some(crate::llm::provider::AssistantPhase::FinalAnswer)),
851 ),
852 ],
853 tool_summaries: Vec::new(),
854 token_usage: None,
855 max_context_tokens: None,
856 loaded_skills: Vec::new(),
857 turn_diagnostics: None,
858 })),
859 error_logs: Vec::new(),
860 },
861 };
862
863 let messages = messages_from_session_listing(&listing);
864 assert_eq!(
865 messages.iter().map(|message| message.phase).collect::<Vec<_>>(),
866 vec![
867 Some(crate::llm::provider::AssistantPhase::Commentary),
868 Some(crate::llm::provider::AssistantPhase::FinalAnswer),
869 ]
870 );
871 }
872}