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
32pub const SYNTHETIC_EVENT_THREAD_IDS: &[&str] = &["app-server", "runtime", "thread-workflow"];
34
35pub 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 #[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 #[serde(default, skip_serializing_if = "Vec::is_empty")]
111 pub tool_allowlist: Vec<String>,
112 #[serde(default, skip_serializing_if = "Option::is_none")]
114 pub developer_instructions: Option<String>,
115 #[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 #[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 #[serde(default, skip_serializing_if = "Option::is_none")]
143 pub parent_thread_id: Option<ThreadId>,
144 #[serde(default, skip_serializing_if = "Option::is_none")]
146 pub forked_from_turn_id: Option<TurnId>,
147 #[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 #[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 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 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 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 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}