Skip to main content

roder_api/
thread.rs

1use std::path::PathBuf;
2use std::sync::Arc;
3
4use serde::{Deserialize, Deserializer, Serialize};
5use time::OffsetDateTime;
6
7use crate::artifacts::ContextArtifactStore;
8use crate::events::{EventEnvelope, ThreadId, TurnId};
9pub use crate::extension::{CheckpointStoreId, ThreadStoreId};
10use crate::extension_state::ExtensionStateRecord;
11use crate::inference::{TokenUsage, cache_hit_rate};
12use crate::inference_routing::ModelSelectionMode;
13use crate::remote_runner::{RunnerDestination, RunnerSessionState, ThreadRunnerBinding};
14use crate::transcript::{InputImage, TranscriptItem};
15
16mod projection;
17pub use projection::{project_thread_item_events, project_turns_from_events};
18
19#[derive(Debug, Clone, Default, PartialEq, Eq)]
20pub struct ThreadListOptions {
21    pub limit: Option<usize>,
22    pub cursor: Option<String>,
23}
24
25#[derive(Debug, Clone, Default)]
26pub struct ThreadListPage {
27    pub threads: Vec<ThreadMetadata>,
28    pub next_cursor: Option<String>,
29    pub backwards_cursor: Option<String>,
30}
31
32/// Thread IDs reserved for system event streams that are not user-visible conversations.
33pub const SYNTHETIC_EVENT_THREAD_IDS: &[&str] = &["app-server", "runtime", "thread-workflow"];
34
35/// Returns true for production-emitted synthetic event streams.
36pub fn is_synthetic_event_thread_id(thread_id: &str) -> bool {
37    SYNTHETIC_EVENT_THREAD_IDS.contains(&thread_id)
38}
39
40#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
41pub struct ThreadUsageMetadata {
42    #[serde(default)]
43    pub prompt_tokens: u64,
44    #[serde(default)]
45    pub completion_tokens: u64,
46    #[serde(default)]
47    pub total_tokens: u64,
48    #[serde(default)]
49    pub cached_prompt_tokens: u64,
50    /// Subset of `prompt_tokens` written to the provider prompt cache.
51    #[serde(default)]
52    pub cache_creation_prompt_tokens: u64,
53    #[serde(default, skip_serializing_if = "Option::is_none")]
54    pub cache_hit_rate: Option<f64>,
55}
56
57impl ThreadUsageMetadata {
58    pub fn add_token_usage(&mut self, usage: &TokenUsage) {
59        self.prompt_tokens = self
60            .prompt_tokens
61            .saturating_add(u64::from(usage.prompt_tokens));
62        self.completion_tokens = self
63            .completion_tokens
64            .saturating_add(u64::from(usage.completion_tokens));
65        self.total_tokens = self
66            .total_tokens
67            .saturating_add(u64::from(usage.total_tokens));
68        self.cached_prompt_tokens = self
69            .cached_prompt_tokens
70            .saturating_add(u64::from(usage.cached_prompt_tokens));
71        self.cache_creation_prompt_tokens = self
72            .cache_creation_prompt_tokens
73            .saturating_add(u64::from(usage.cache_creation_prompt_tokens));
74        self.cache_hit_rate = if self.prompt_tokens == 0 {
75            None
76        } else if self.prompt_tokens > u64::from(u32::MAX) {
77            Some(
78                (self.cached_prompt_tokens.min(self.prompt_tokens) as f64)
79                    / (self.prompt_tokens as f64),
80            )
81        } else {
82            cache_hit_rate(self.prompt_tokens as u32, self.cached_prompt_tokens as u32)
83        };
84    }
85
86    pub fn is_empty(&self) -> bool {
87        self.prompt_tokens == 0
88            && self.completion_tokens == 0
89            && self.total_tokens == 0
90            && self.cached_prompt_tokens == 0
91            && self.cache_creation_prompt_tokens == 0
92    }
93}
94
95#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
96pub struct ThreadMetadata {
97    pub thread_id: ThreadId,
98    pub title: Option<String>,
99    #[serde(deserialize_with = "deserialize_thread_workspace")]
100    pub workspace: String,
101    #[serde(default, skip_serializing_if = "Option::is_none")]
102    pub workspace_id: Option<String>,
103    #[serde(default, skip_serializing_if = "Option::is_none")]
104    pub root_id: Option<String>,
105    pub provider: Option<String>,
106    pub model: Option<String>,
107    #[serde(default, skip_serializing_if = "Option::is_none")]
108    pub selection_mode: Option<ModelSelectionMode>,
109    /// Per-thread tool filter applied on top of the runtime allowlist. Empty = no filtering.
110    #[serde(default, skip_serializing_if = "Vec::is_empty")]
111    pub tool_allowlist: Vec<String>,
112    /// Host-supplied instructions added to the developer slot of every turn's inference request.
113    #[serde(default, skip_serializing_if = "Option::is_none")]
114    pub developer_instructions: Option<String>,
115    /**
116     * Host-executed tool specs advertised to the model on every turn of this thread. Calls to
117     * these tools pause on a `thread/toolExecutionRequested` notification until the host client
118     * answers with `tools/resolve`.
119     */
120    #[serde(default, skip_serializing_if = "Vec::is_empty")]
121    pub external_tools: Vec<crate::tools::ToolSpec>,
122    #[serde(default, skip_serializing_if = "Option::is_none")]
123    pub runner_destination: Option<RunnerDestination>,
124    #[serde(default, skip_serializing_if = "Option::is_none")]
125    pub runner_state: Option<RunnerSessionState>,
126    /**
127     * Set when the thread explicitly selected a remote runner at creation.
128     * Native coding tools for the thread route through this runner workspace;
129     * absent = local tool execution even when a runtime-level runner
130     * destination is configured.
131     */
132    #[serde(default, skip_serializing_if = "Option::is_none")]
133    pub runner_binding: Option<ThreadRunnerBinding>,
134    #[serde(with = "time::serde::rfc3339")]
135    pub created_at: OffsetDateTime,
136    #[serde(with = "time::serde::rfc3339")]
137    pub updated_at: OffsetDateTime,
138    pub message_count: u32,
139    #[serde(default, skip_serializing_if = "Option::is_none")]
140    pub usage: Option<ThreadUsageMetadata>,
141    /// Parent thread this conversation was forked from. Absent for normal threads.
142    #[serde(default, skip_serializing_if = "Option::is_none")]
143    pub parent_thread_id: Option<ThreadId>,
144    /// Parent turn the fork branched at, when the fork targeted a specific turn.
145    #[serde(default, skip_serializing_if = "Option::is_none")]
146    pub forked_from_turn_id: Option<TurnId>,
147    /**
148     * Provider-neutral workspace-fork provenance for threads whose
149     * workspace is a fork of another workspace (roadmap phase 81; the
150     * phase-90 Git-worktree MVP shape was folded into this canonical type).
151     */
152    #[serde(default, skip_serializing_if = "Option::is_none")]
153    pub workspace_fork: Option<crate::forks::WorkspaceFork>,
154}
155
156pub fn validate_thread_workspace(workspace: &str) -> anyhow::Result<String> {
157    let workspace = workspace.trim();
158    anyhow::ensure!(!workspace.is_empty(), "thread workspace is required");
159    anyhow::ensure!(
160        std::path::Path::new(workspace).is_absolute(),
161        "thread workspace must be an absolute path: {workspace}"
162    );
163    Ok(workspace.to_string())
164}
165
166fn deserialize_thread_workspace<'de, D>(deserializer: D) -> Result<String, D::Error>
167where
168    D: Deserializer<'de>,
169{
170    let workspace = String::deserialize(deserializer)?;
171    validate_thread_workspace(&workspace).map_err(serde::de::Error::custom)
172}
173
174#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
175pub struct TurnRecord {
176    pub thread_id: ThreadId,
177    pub turn_id: TurnId,
178    pub items: Vec<TranscriptItem>,
179    #[serde(with = "time::serde::rfc3339")]
180    pub created_at: OffsetDateTime,
181    #[serde(with = "time::serde::rfc3339::option")]
182    pub completed_at: Option<OffsetDateTime>,
183    #[serde(default, skip_serializing_if = "Option::is_none")]
184    pub usage: Option<TokenUsage>,
185    /// Normalized stop reason from `TurnCompleted`; `None` for failed or interrupted turns.
186    #[serde(default, skip_serializing_if = "Option::is_none")]
187    pub finish_reason: Option<String>,
188}
189
190#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
191#[serde(rename_all = "camelCase")]
192pub enum ThreadItemStatus {
193    InProgress,
194    Completed,
195    Failed,
196}
197
198#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
199#[serde(tag = "type", rename_all = "camelCase")]
200pub enum ThreadItem {
201    UserMessage {
202        id: String,
203        text: String,
204        #[serde(default, skip_serializing_if = "Vec::is_empty")]
205        images: Vec<InputImage>,
206        #[serde(default, skip_serializing_if = "Option::is_none")]
207        status: Option<ThreadItemStatus>,
208    },
209    AgentMessage {
210        id: String,
211        text: String,
212        #[serde(default, skip_serializing_if = "Option::is_none")]
213        phase: Option<String>,
214        #[serde(default, skip_serializing_if = "Option::is_none")]
215        status: Option<ThreadItemStatus>,
216    },
217    Reasoning {
218        id: String,
219        #[serde(default, skip_serializing_if = "Vec::is_empty")]
220        summary: Vec<String>,
221        #[serde(default, skip_serializing_if = "Vec::is_empty")]
222        content: Vec<String>,
223        #[serde(default, skip_serializing_if = "Option::is_none")]
224        status: Option<ThreadItemStatus>,
225    },
226    ToolExecution {
227        id: String,
228        #[serde(rename = "toolCallId")]
229        tool_call_id: String,
230        #[serde(rename = "toolName")]
231        tool_name: String,
232        status: ThreadItemStatus,
233        #[serde(default, skip_serializing_if = "Option::is_none")]
234        input: Option<serde_json::Value>,
235        #[serde(default, skip_serializing_if = "Option::is_none")]
236        output: Option<String>,
237        #[serde(default, skip_serializing_if = "Option::is_none")]
238        error: Option<String>,
239    },
240    RoutingDecision {
241        id: String,
242        decision: crate::events::InferenceRoutingDecisionEvent,
243        #[serde(default, skip_serializing_if = "Option::is_none")]
244        status: Option<ThreadItemStatus>,
245    },
246    Compaction {
247        id: String,
248        summary: String,
249        #[serde(default, skip_serializing_if = "Option::is_none")]
250        status: Option<ThreadItemStatus>,
251    },
252    Error {
253        id: String,
254        message: String,
255        #[serde(default, skip_serializing_if = "Option::is_none")]
256        status: Option<ThreadItemStatus>,
257    },
258    Raw {
259        id: String,
260        payload: serde_json::Value,
261        #[serde(default, skip_serializing_if = "Option::is_none")]
262        status: Option<ThreadItemStatus>,
263    },
264}
265
266#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
267pub struct ThreadItemTurnRecord {
268    pub thread_id: ThreadId,
269    pub turn_id: TurnId,
270    #[serde(with = "time::serde::rfc3339")]
271    pub created_at: OffsetDateTime,
272    pub items: Vec<ThreadItem>,
273}
274
275#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
276#[serde(tag = "type", rename_all = "camelCase")]
277pub enum ThreadItemDelta {
278    AgentMessageText {
279        delta: String,
280        #[serde(default, skip_serializing_if = "Option::is_none")]
281        phase: Option<String>,
282    },
283    ReasoningText {
284        delta: String,
285        #[serde(rename = "contentIndex")]
286        content_index: usize,
287    },
288    ReasoningSummaryPartAdded {
289        #[serde(rename = "summaryIndex")]
290        summary_index: usize,
291    },
292    ReasoningSummaryText {
293        delta: String,
294        #[serde(rename = "summaryIndex")]
295        summary_index: usize,
296    },
297}
298
299#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
300#[serde(tag = "type", rename_all = "camelCase")]
301pub enum ThreadItemEventKind {
302    ItemStarted {
303        item: ThreadItem,
304    },
305    ItemDelta {
306        #[serde(rename = "itemId")]
307        item_id: String,
308        delta: ThreadItemDelta,
309    },
310    ItemCompleted {
311        item: ThreadItem,
312    },
313}
314
315#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
316#[serde(rename_all = "camelCase")]
317pub struct ThreadItemEvent {
318    pub seq: u64,
319    #[serde(rename = "eventId")]
320    pub event_id: String,
321    #[serde(rename = "threadId")]
322    pub thread_id: ThreadId,
323    #[serde(rename = "turnId")]
324    pub turn_id: TurnId,
325    #[serde(with = "time::serde::rfc3339")]
326    pub timestamp: OffsetDateTime,
327    pub event: ThreadItemEventKind,
328}
329
330#[derive(Debug, Clone, Default, Serialize, Deserialize)]
331pub struct ThreadSnapshot {
332    pub metadata: Option<ThreadMetadata>,
333    pub events: Vec<EventEnvelope>,
334    pub turns: Vec<TurnRecord>,
335    #[serde(default)]
336    pub item_events: Vec<ThreadItemEvent>,
337    pub extension_states: Vec<ExtensionStateRecord>,
338}
339
340impl ThreadItem {
341    pub fn id(&self) -> &str {
342        match self {
343            ThreadItem::UserMessage { id, .. }
344            | ThreadItem::AgentMessage { id, .. }
345            | ThreadItem::Reasoning { id, .. }
346            | ThreadItem::ToolExecution { id, .. }
347            | ThreadItem::RoutingDecision { id, .. }
348            | ThreadItem::Compaction { id, .. }
349            | ThreadItem::Error { id, .. }
350            | ThreadItem::Raw { id, .. } => id,
351        }
352    }
353}
354
355#[async_trait::async_trait]
356pub trait ThreadStore: Send + Sync {
357    fn id(&self) -> ThreadStoreId;
358
359    fn local_thread_root(&self) -> Option<PathBuf> {
360        None
361    }
362
363    fn context_artifact_store(&self) -> Option<ContextArtifactStore> {
364        None
365    }
366
367    async fn create_thread(&self, metadata: ThreadMetadata) -> anyhow::Result<ThreadMetadata>;
368    async fn update_thread_metadata(
369        &self,
370        metadata: ThreadMetadata,
371    ) -> anyhow::Result<ThreadMetadata> {
372        Ok(metadata)
373    }
374    async fn list_threads(&self) -> anyhow::Result<Vec<ThreadMetadata>>;
375    async fn list_threads_page(
376        &self,
377        options: ThreadListOptions,
378    ) -> anyhow::Result<ThreadListPage> {
379        let mut threads = self.list_threads().await?;
380        threads.sort_by_key(|thread| std::cmp::Reverse(thread.updated_at));
381        let offset = options
382            .cursor
383            .as_deref()
384            .and_then(|cursor| cursor.parse::<usize>().ok())
385            .unwrap_or(0)
386            .min(threads.len());
387        let limit = options
388            .limit
389            .unwrap_or(threads.len().saturating_sub(offset));
390        let next_offset = offset.saturating_add(limit).min(threads.len());
391        let total = threads.len();
392        let page_threads = threads
393            .into_iter()
394            .skip(offset)
395            .take(limit)
396            .collect::<Vec<_>>();
397        Ok(ThreadListPage {
398            threads: page_threads,
399            next_cursor: (next_offset < total).then(|| next_offset.to_string()),
400            backwards_cursor: (offset > 0).then(|| offset.saturating_sub(limit).to_string()),
401        })
402    }
403    async fn load_thread_metadata(
404        &self,
405        thread_id: &ThreadId,
406    ) -> anyhow::Result<Option<ThreadMetadata>> {
407        Ok(self
408            .load_thread(thread_id)
409            .await?
410            .and_then(|snapshot| snapshot.metadata))
411    }
412    async fn load_thread(&self, thread_id: &ThreadId) -> anyhow::Result<Option<ThreadSnapshot>>;
413    /// Loads extension state without requiring callers to project a full thread snapshot.
414    ///
415    /// Stores with dedicated extension-state backing should override this method. The default
416    /// preserves compatibility for stores that only expose extension state through `load_thread`.
417    async fn load_extension_states(
418        &self,
419        thread_id: &ThreadId,
420    ) -> anyhow::Result<Vec<ExtensionStateRecord>> {
421        Ok(self
422            .load_thread(thread_id)
423            .await?
424            .map_or_else(Vec::new, |snapshot| snapshot.extension_states))
425    }
426    async fn archive_thread(&self, thread_id: &ThreadId) -> anyhow::Result<bool> {
427        let _ = thread_id;
428        anyhow::bail!("thread store {} does not support archive", self.id())
429    }
430    async fn append_event(
431        &self,
432        thread_id: &ThreadId,
433        envelope: &EventEnvelope,
434    ) -> anyhow::Result<()>;
435    async fn append_item_event(
436        &self,
437        thread_id: &ThreadId,
438        item_event: &ThreadItemEvent,
439    ) -> anyhow::Result<()> {
440        let _ = (thread_id, item_event);
441        Ok(())
442    }
443    async fn append_extension_state(
444        &self,
445        thread_id: &ThreadId,
446        record: &ExtensionStateRecord,
447    ) -> anyhow::Result<()> {
448        let _ = (thread_id, record);
449        anyhow::bail!(
450            "thread store {} does not support extension state",
451            self.id()
452        )
453    }
454}
455
456pub trait ThreadStoreFactory: Send + Sync + 'static {
457    fn id(&self) -> ThreadStoreId;
458    fn create(&self) -> Arc<dyn ThreadStore>;
459}
460
461#[async_trait::async_trait]
462pub trait CheckpointStore: Send + Sync {
463    fn id(&self) -> CheckpointStoreId;
464    async fn save_snapshot(&self, snapshot: ThreadSnapshot) -> anyhow::Result<()>;
465    async fn load_snapshot(&self, thread_id: &ThreadId) -> anyhow::Result<Option<ThreadSnapshot>>;
466}
467
468pub trait CheckpointStoreFactory: Send + Sync + 'static {
469    fn id(&self) -> CheckpointStoreId;
470    fn create(&self) -> Arc<dyn CheckpointStore>;
471}
472
473#[cfg(test)]
474mod tests {
475    use super::*;
476    use crate::inference::ModelSelection;
477
478    #[test]
479    fn synthetic_event_thread_ids_are_reserved_production_ids() {
480        assert!(is_synthetic_event_thread_id("app-server"));
481        assert!(is_synthetic_event_thread_id("runtime"));
482        assert!(is_synthetic_event_thread_id("thread-workflow"));
483
484        assert!(!is_synthetic_event_thread_id("thread-discovery"));
485        assert!(!is_synthetic_event_thread_id("thread-plan"));
486        assert!(!is_synthetic_event_thread_id("thread-process"));
487        assert!(!is_synthetic_event_thread_id("thread-1"));
488    }
489
490    #[test]
491    fn thread_fork_metadata_is_additive_for_legacy_records() {
492        // A legacy thread record without any fork fields still deserializes.
493        let legacy = serde_json::json!({
494            "thread_id": "thread-old",
495            "title": null,
496            "workspace": "/workspace",
497            "provider": null,
498            "model": null,
499            "created_at": "1970-01-01T00:00:00Z",
500            "updated_at": "1970-01-01T00:00:00Z",
501            "message_count": 3
502        });
503        let metadata: ThreadMetadata = serde_json::from_value(legacy).unwrap();
504        assert!(metadata.parent_thread_id.is_none());
505        assert!(metadata.forked_from_turn_id.is_none());
506        assert!(metadata.workspace_fork.is_none());
507
508        // Non-forked threads serialize without any fork keys.
509        let value = serde_json::to_value(&metadata).unwrap();
510        assert!(value.get("parentThreadId").is_none());
511        assert!(value.get("parent_thread_id").is_none());
512        assert!(value.get("workspace_fork").is_none());
513    }
514
515    #[test]
516    fn thread_fork_metadata_round_trips_workspace_fork_provenance() {
517        let fork = crate::forks::WorkspaceFork {
518            id: "/repo/.roder/worktrees/parser-experiment".to_string(),
519            provider_id: "git-worktree".to_string(),
520            source_workspace: std::path::PathBuf::from("/repo"),
521            workspace: std::path::PathBuf::from("/repo/.roder/worktrees/parser-experiment"),
522            status: crate::forks::ForkStatus::Active,
523            provenance: crate::forks::ForkProvenance {
524                branch: Some("roder/fork/parser-experiment".to_string()),
525                source_branch: Some("main".to_string()),
526                source_commit: Some("abc123".to_string()),
527                snapshot_id: None,
528                session_id: None,
529                created_at: OffsetDateTime::UNIX_EPOCH,
530            },
531            cleanup: crate::forks::ForkCleanupPolicy::Explicit,
532            metadata: serde_json::json!({}),
533        };
534        let value = serde_json::to_value(&fork).unwrap();
535        assert_eq!(value["providerId"], "git-worktree");
536        assert_eq!(value["status"], "active");
537        assert_eq!(value["cleanup"], "explicit");
538        assert_eq!(value["provenance"]["sourceCommit"], "abc123");
539
540        let round_trip: crate::forks::WorkspaceFork = serde_json::from_value(value).unwrap();
541        assert_eq!(round_trip, fork);
542
543        // A detached-HEAD fork keeps clear provenance without a branch name.
544        let detached = crate::forks::WorkspaceFork {
545            provenance: crate::forks::ForkProvenance {
546                source_branch: None,
547                ..fork.provenance.clone()
548            },
549            ..fork
550        };
551        let value = serde_json::to_value(&detached).unwrap();
552        assert!(value["provenance"].get("sourceBranch").is_none());
553    }
554
555    #[test]
556    fn thread_metadata_timestamps_serialize_as_rfc3339_strings() {
557        let value = serde_json::to_value(ThreadMetadata {
558            thread_id: "thread-a".to_string(),
559            title: None,
560            workspace: "/workspace".to_string(),
561            workspace_id: None,
562            root_id: None,
563            provider: None,
564            model: None,
565            selection_mode: None,
566            tool_allowlist: Vec::new(),
567            developer_instructions: None,
568            external_tools: Vec::new(),
569            runner_destination: None,
570            runner_state: None,
571            runner_binding: None,
572            parent_thread_id: None,
573            forked_from_turn_id: None,
574            workspace_fork: None,
575            created_at: OffsetDateTime::UNIX_EPOCH,
576            updated_at: OffsetDateTime::UNIX_EPOCH,
577            message_count: 0,
578            usage: None,
579        })
580        .unwrap();
581
582        assert_eq!(value["created_at"], "1970-01-01T00:00:00Z");
583        assert_eq!(value["updated_at"], "1970-01-01T00:00:00Z");
584        assert_eq!(value["workspace"], "/workspace");
585    }
586
587    #[test]
588    fn thread_metadata_deserializes_without_selection_mode() {
589        let value = serde_json::json!({
590            "thread_id": "thread-a",
591            "title": null,
592            "workspace": "/workspace",
593            "provider": "codex",
594            "model": "gpt-5.5",
595            "created_at": "1970-01-01T00:00:00Z",
596            "updated_at": "1970-01-01T00:00:00Z",
597            "message_count": 0
598        });
599
600        let metadata = serde_json::from_value::<ThreadMetadata>(value).unwrap();
601
602        assert_eq!(metadata.provider.as_deref(), Some("codex"));
603        assert_eq!(metadata.model.as_deref(), Some("gpt-5.5"));
604        assert_eq!(metadata.selection_mode, None);
605    }
606
607    #[test]
608    fn thread_metadata_round_trips_auto_selection_mode() {
609        let metadata = ThreadMetadata {
610            thread_id: "thread-a".to_string(),
611            title: None,
612            workspace: "/workspace".to_string(),
613            workspace_id: None,
614            root_id: None,
615            provider: Some("codex".to_string()),
616            model: Some("gpt-5.5".to_string()),
617            selection_mode: Some(ModelSelectionMode::auto(
618                "local-router:coding",
619                "local-router",
620                "Auto: Coding",
621                ModelSelection {
622                    provider: "codex".to_string(),
623                    model: "gpt-5.5".to_string(),
624                },
625                Some("coding".to_string()),
626                Some("low".to_string()),
627            )),
628            tool_allowlist: Vec::new(),
629            developer_instructions: None,
630            external_tools: Vec::new(),
631            runner_destination: None,
632            runner_state: None,
633            runner_binding: None,
634            parent_thread_id: None,
635            forked_from_turn_id: None,
636            workspace_fork: None,
637            created_at: OffsetDateTime::UNIX_EPOCH,
638            updated_at: OffsetDateTime::UNIX_EPOCH,
639            message_count: 0,
640            usage: None,
641        };
642
643        let value = serde_json::to_value(&metadata).unwrap();
644        let round_trip = serde_json::from_value::<ThreadMetadata>(value).unwrap();
645
646        assert_eq!(round_trip, metadata);
647    }
648
649    #[test]
650    fn thread_metadata_requires_workspace_when_deserializing() {
651        let value = serde_json::json!({
652            "thread_id": "thread-a",
653            "title": null,
654            "provider": null,
655            "model": null,
656            "created_at": "1970-01-01T00:00:00Z",
657            "updated_at": "1970-01-01T00:00:00Z",
658            "message_count": 0
659        });
660
661        let result = serde_json::from_value::<ThreadMetadata>(value);
662
663        assert!(result.is_err());
664    }
665
666    #[test]
667    fn thread_metadata_rejects_blank_or_relative_workspace_when_deserializing() {
668        for workspace in ["", "project"] {
669            let value = serde_json::json!({
670                "thread_id": "thread-a",
671                "title": null,
672                "workspace": workspace,
673                "provider": null,
674                "model": null,
675                "created_at": "1970-01-01T00:00:00Z",
676                "updated_at": "1970-01-01T00:00:00Z",
677                "message_count": 0
678            });
679
680            let result = serde_json::from_value::<ThreadMetadata>(value);
681
682            assert!(result.is_err(), "workspace {workspace:?} should fail");
683        }
684    }
685
686    #[test]
687    fn thread_usage_metadata_accumulates_cache_hit_rate() {
688        let mut usage = ThreadUsageMetadata::default();
689
690        usage.add_token_usage(
691            &TokenUsage::new(100, 10, 110)
692                .with_cached_prompt_tokens(92)
693                .with_cache_creation_prompt_tokens(5),
694        );
695        usage.add_token_usage(
696            &TokenUsage::new(50, 5, 55)
697                .with_cached_prompt_tokens(43)
698                .with_cache_creation_prompt_tokens(3),
699        );
700
701        assert_eq!(usage.prompt_tokens, 150);
702        assert_eq!(usage.cached_prompt_tokens, 135);
703        assert_eq!(usage.cache_creation_prompt_tokens, 8);
704        assert!((usage.cache_hit_rate.unwrap() - 0.9).abs() < f64::EPSILON);
705    }
706
707    #[test]
708    fn thread_item_events_replay_reasoning_and_final_answer_into_stable_items() {
709        let timestamp = OffsetDateTime::UNIX_EPOCH;
710        let events = vec![
711            ThreadItemEvent {
712                seq: 1,
713                event_id: "event-1".to_string(),
714                thread_id: "thread-1".to_string(),
715                turn_id: "turn-1".to_string(),
716                timestamp,
717                event: ThreadItemEventKind::ItemStarted {
718                    item: ThreadItem::Reasoning {
719                        id: "turn-1-agent-reasoning".to_string(),
720                        summary: Vec::new(),
721                        content: vec![String::new()],
722                        status: Some(ThreadItemStatus::InProgress),
723                    },
724                },
725            },
726            ThreadItemEvent {
727                seq: 2,
728                event_id: "event-2".to_string(),
729                thread_id: "thread-1".to_string(),
730                turn_id: "turn-1".to_string(),
731                timestamp,
732                event: ThreadItemEventKind::ItemDelta {
733                    item_id: "turn-1-agent-reasoning".to_string(),
734                    delta: ThreadItemDelta::ReasoningText {
735                        delta: "Inspecting".to_string(),
736                        content_index: 0,
737                    },
738                },
739            },
740            ThreadItemEvent {
741                seq: 3,
742                event_id: "event-3".to_string(),
743                thread_id: "thread-1".to_string(),
744                turn_id: "turn-1".to_string(),
745                timestamp,
746                event: ThreadItemEventKind::ItemDelta {
747                    item_id: "turn-1-agent-final_answer".to_string(),
748                    delta: ThreadItemDelta::AgentMessageText {
749                        delta: "Done".to_string(),
750                        phase: Some("final_answer".to_string()),
751                    },
752                },
753            },
754            ThreadItemEvent {
755                seq: 4,
756                event_id: "event-4".to_string(),
757                thread_id: "thread-1".to_string(),
758                turn_id: "turn-1".to_string(),
759                timestamp,
760                event: ThreadItemEventKind::ItemCompleted {
761                    item: ThreadItem::AgentMessage {
762                        id: "turn-1-agent-final_answer".to_string(),
763                        text: "Done.".to_string(),
764                        phase: Some("final_answer".to_string()),
765                        status: Some(ThreadItemStatus::Completed),
766                    },
767                },
768            },
769        ];
770
771        let turns = project_thread_item_events(&events);
772
773        assert_eq!(turns.len(), 1);
774        assert_eq!(turns[0].turn_id, "turn-1");
775        assert_eq!(
776            turns[0].items,
777            vec![
778                ThreadItem::Reasoning {
779                    id: "turn-1-agent-reasoning".to_string(),
780                    summary: Vec::new(),
781                    content: vec!["Inspecting".to_string()],
782                    status: Some(ThreadItemStatus::InProgress),
783                },
784                ThreadItem::AgentMessage {
785                    id: "turn-1-agent-final_answer".to_string(),
786                    text: "Done.".to_string(),
787                    phase: Some("final_answer".to_string()),
788                    status: Some(ThreadItemStatus::Completed),
789                }
790            ]
791        );
792    }
793}