agentic_core/storage/
response.rs1use std::collections::HashMap;
4use std::convert::TryFrom;
5use std::sync::Arc;
6
7use super::models::{item, response};
8use super::pool::DbPool;
9use super::types::{InOutItem, ResponseData, ResponseMetadata, StorageError, StoreResult};
10use crate::utils::common::{serialize_to_string, uuid7_str};
11
12#[derive(Clone, Debug)]
14pub struct ResponseStore {
15 pool: Option<Arc<DbPool>>,
16}
17
18impl ResponseStore {
19 #[must_use]
23 pub fn disabled() -> Self {
24 Self { pool: None }
25 }
26
27 #[must_use]
33 pub fn new(pool: Arc<DbPool>) -> Self {
34 Self { pool: Some(pool) }
35 }
36
37 fn pool(&self) -> StoreResult<&DbPool> {
43 self.pool.as_deref().ok_or(StorageError::NotConfigured)
44 }
45
46 pub async fn get(&self, response_id: &str) -> StoreResult<ResponseData> {
52 let pool = self.pool()?;
53 let row = response::get(pool, response_id)
54 .await?
55 .ok_or_else(|| StorageError::not_found("Response", response_id))?;
56 Ok(row.into())
57 }
58
59 pub async fn rehydrate(&self, response_id: &str) -> StoreResult<Vec<InOutItem>> {
67 let pool = self.pool()?;
68 let response = self.get(response_id).await?;
69 let rows = item::get_items(pool, &response.history_item_ids).await?;
70 let mut items_by_id: HashMap<String, InOutItem> = rows
71 .into_iter()
72 .filter_map(|row| {
73 let id = row.id.clone();
74 row.as_inout().map(|item| (id, item))
75 })
76 .collect();
77
78 let ordered_items = response
79 .history_item_ids
80 .iter()
81 .filter_map(|id| items_by_id.remove(id))
82 .collect();
83
84 Ok(ordered_items)
85 }
86
87 pub async fn persist(
95 &self,
96 response_id: &str,
97 previous_response_id: Option<&str>,
98 new_items: Vec<InOutItem>,
99 metadata: &ResponseMetadata,
100 ) -> StoreResult<()> {
101 self.persist_with_conversation_id(response_id, None, previous_response_id, new_items, metadata)
102 .await
103 }
104
105 pub(crate) async fn persist_with_conversation_id(
111 &self,
112 response_id: &str,
113 conversation_id: Option<&str>,
114 previous_response_id: Option<&str>,
115 new_items: Vec<InOutItem>,
116 metadata: &ResponseMetadata,
117 ) -> StoreResult<()> {
118 let pool = self.pool()?;
119
120 let mut item_ids: Vec<String> = match previous_response_id {
121 Some(prev_id) => self.get(prev_id).await?.history_item_ids,
122 None => Vec::new(),
123 };
124 let mut items_: Vec<(String, String)> = Vec::new();
125 for any_item in new_items {
126 let item_id = uuid7_str("item_");
127 item_ids.push(item_id.clone());
128 let data_str = String::try_from(&any_item)?;
129 items_.push((item_id, data_str));
130 }
131 let history_item_ids_json = serialize_to_string(&item_ids)?;
132 let metadata_json = String::try_from(metadata)?;
133
134 let mut tx = pool.begin().await?;
135
136 item::create_in_tx(&mut tx, items_, None).await?;
137
138 response::create_in_tx(
139 &mut tx,
140 response_id,
141 conversation_id,
142 previous_response_id,
143 Some(&history_item_ids_json),
144 Some(&metadata_json),
145 )
146 .await?;
147 tx.commit().await?;
148
149 Ok(())
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use super::super::types::ResponseMetadata;
156 use super::*;
157
158 #[test]
159 fn test_response_store_disabled() {
160 let store = ResponseStore::disabled();
161 assert!(store.pool().is_err());
162 }
163
164 #[test]
165 fn test_response_metadata_default() {
166 let meta = ResponseMetadata::default();
167 assert!(meta.model.is_empty());
168 assert!(meta.previous_response_id.is_none());
169 }
170}