Skip to main content

vtcode_core/core/
threads.rs

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/// Unique identifier for a thread in the runtime.
20#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, StringNewtype)]
21pub struct ThreadId(String);
22
23/// Unique identifier for a submission within a thread turn.
24#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize, StringNewtype)]
25pub struct SubmissionId(String);
26
27impl SubmissionId {
28    /// Generate a new submission identifier with a `sub-` prefix.
29    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    // Preserve the source session's primary agent ("mode") so a resume/fork keeps
361    // running the same agent instead of falling back to the config default.
362    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        // Resume preserves the source session's primary agent ("mode").
634        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}