Skip to main content

agentic_core/storage/
response.rs

1//! Response storage operations and queries.
2
3use 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/// Response storage operations.
13#[derive(Clone, Debug)]
14pub struct ResponseStore {
15    pool: Option<Arc<DbPool>>,
16}
17
18impl ResponseStore {
19    /// Creates a disabled response store (no persistence).
20    ///
21    /// Useful for testing or when response storage is not configured.
22    #[must_use]
23    pub fn disabled() -> Self {
24        Self { pool: None }
25    }
26
27    /// Creates a new response store with database pool.
28    ///
29    /// # Arguments
30    ///
31    /// * `pool` - Connection pool for database access
32    #[must_use]
33    pub fn new(pool: Arc<DbPool>) -> Self {
34        Self { pool: Some(pool) }
35    }
36
37    /// Returns a reference to the database pool.
38    ///
39    /// # Errors
40    ///
41    /// Returns [`StorageError::NotConfigured`] if store is disabled (no pool configured).
42    fn pool(&self) -> StoreResult<&DbPool> {
43        self.pool.as_deref().ok_or(StorageError::NotConfigured)
44    }
45
46    /// Retrieves a response by ID.
47    ///
48    /// # Errors
49    ///
50    /// Returns error if response not found, database query fails, or store is disabled.
51    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    /// Rehydrates a response with full history.
60    ///
61    /// Fetches all history items referenced by a response.
62    ///
63    /// # Errors
64    ///
65    /// Returns error if database query fails or store is disabled.
66    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    /// Persists a response with its items and metadata.
88    ///
89    /// Creates items and stores the associated response record.
90    ///
91    /// # Errors
92    ///
93    /// Returns [`StorageError`] if database operation fails or store is disabled.
94    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    /// Persists a response while retaining its inherited conversation ID.
106    ///
107    /// # Errors
108    ///
109    /// Returns [`StorageError`] if database operation fails or store is disabled.
110    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}