Skip to main content

ai_agents_runtime/spawner/
storage.rs

1//! Storage adapter that isolates agent data in a shared backend.
2
3use std::collections::BTreeSet;
4use std::sync::Arc;
5
6use async_trait::async_trait;
7
8use ai_agents_core::traits::storage::StorageCapability;
9use ai_agents_core::{
10    AgentError, AgentSnapshot, AgentStorage, FactFilter, KeyFact, Result, SessionFilter,
11    SessionMetadata, SessionSummary,
12};
13
14const FORWARDED_CAPABILITIES: [StorageCapability; 6] = [
15    StorageCapability::Snapshot,
16    StorageCapability::SessionMetadata,
17    StorageCapability::SessionFiltering,
18    StorageCapability::ActorFacts,
19    StorageCapability::ActorRelationships,
20    StorageCapability::ActorDataDeletion,
21];
22
23/// Wraps shared storage with reversible flat keys scoped to one agent namespace.
24pub struct NamespacedStorage {
25    inner: Arc<dyn AgentStorage>,
26    namespace: String,
27    encoded_namespace: String,
28}
29
30impl NamespacedStorage {
31    pub fn new(inner: Arc<dyn AgentStorage>, namespace: impl Into<String>) -> Self {
32        let namespace = namespace.into();
33        Self {
34            inner,
35            encoded_namespace: encode_component(&namespace),
36            namespace,
37        }
38    }
39
40    fn require(&self, capability: StorageCapability) -> Result<()> {
41        if self.supports(capability) {
42            Ok(())
43        } else {
44            Err(AgentError::UnsupportedStorageCapability(capability))
45        }
46    }
47
48    fn encode_key(&self, kind: &str, value: &str) -> String {
49        // Hex encoding keeps every generated key flat ASCII without exposing path separators.
50        format!(
51            "a7ns1_{}_{}_{}",
52            kind,
53            self.encoded_namespace,
54            encode_component(value)
55        )
56    }
57
58    fn decode_key(&self, kind: &str, value: &str) -> Option<String> {
59        let marker = format!("a7ns1_{}_{}_", kind, self.encoded_namespace);
60        value.strip_prefix(&marker).and_then(decode_component)
61    }
62
63    fn session_key(&self, session_id: &str) -> String {
64        self.encode_key("session", session_id)
65    }
66
67    fn legacy_session_key(&self, session_id: &str) -> String {
68        format!("{}/{}", self.namespace, session_id)
69    }
70
71    fn decode_legacy_session_key(&self, value: &str) -> Option<String> {
72        value
73            .strip_prefix(&format!("{}/", self.namespace))
74            .map(str::to_owned)
75    }
76
77    fn decode_session_key(&self, value: &str) -> Result<Option<String>> {
78        let marker = format!("a7ns1_session_{}_", self.encoded_namespace);
79        if let Some(encoded) = value.strip_prefix(&marker) {
80            return decode_component(encoded).map(Some).ok_or_else(|| {
81                AgentError::Persistence(
82                    "Storage returned a malformed session key for this namespace".into(),
83                )
84            });
85        }
86        Ok(self.decode_legacy_session_key(value))
87    }
88
89    fn agent_key(&self, agent_id: &str) -> String {
90        self.encode_key("agent", agent_id)
91    }
92
93    fn actor_key(&self, actor_id: &str) -> String {
94        self.encode_key("actor", actor_id)
95    }
96
97    fn decode_agent(&self, agent_id: &str) -> Result<String> {
98        self.decode_key("agent", agent_id).ok_or_else(|| {
99            AgentError::Persistence("Storage returned an agent outside this namespace".into())
100        })
101    }
102
103    fn decode_actor(&self, actor_id: &str) -> Result<String> {
104        self.decode_key("actor", actor_id).ok_or_else(|| {
105            AgentError::Persistence("Storage returned an actor outside this namespace".into())
106        })
107    }
108
109    fn encode_snapshot(&self, snapshot: &AgentSnapshot) -> AgentSnapshot {
110        let mut snapshot = snapshot.clone();
111        snapshot.agent_id = self.agent_key(&snapshot.agent_id);
112        snapshot
113    }
114
115    fn decode_snapshot(&self, mut snapshot: AgentSnapshot) -> Result<AgentSnapshot> {
116        snapshot.agent_id = self.decode_agent(&snapshot.agent_id)?;
117        Ok(snapshot)
118    }
119
120    fn encode_metadata(&self, metadata: &SessionMetadata) -> SessionMetadata {
121        let mut metadata = metadata.clone();
122        metadata.actor_id = metadata.actor_id.map(|actor| self.actor_key(&actor));
123        metadata.actors = metadata
124            .actors
125            .into_iter()
126            .map(|actor| self.actor_key(&actor))
127            .collect();
128        metadata
129    }
130
131    fn decode_metadata(&self, mut metadata: SessionMetadata) -> Result<SessionMetadata> {
132        metadata.actor_id = metadata
133            .actor_id
134            .map(|actor| self.decode_actor(&actor))
135            .transpose()?;
136        metadata.actors = metadata
137            .actors
138            .into_iter()
139            .map(|actor| self.decode_actor(&actor))
140            .collect::<Result<Vec<_>>>()?;
141        Ok(metadata)
142    }
143
144    fn encode_fact(&self, fact: &KeyFact) -> KeyFact {
145        let mut fact = fact.clone();
146        fact.actor_id = fact.actor_id.map(|actor| self.actor_key(&actor));
147        fact
148    }
149
150    fn decode_fact(&self, mut fact: KeyFact) -> Option<KeyFact> {
151        if let Some(actor_id) = fact.actor_id.take() {
152            fact.actor_id = Some(self.decode_key("actor", &actor_id)?);
153        }
154        Some(fact)
155    }
156
157    fn backend_session_filter(&self, filter: &SessionFilter) -> SessionFilter {
158        let mut filter = filter.clone();
159        filter.actor_id = None;
160        filter.agent_id = None;
161        filter.limit = None;
162        filter
163    }
164
165    fn decode_session_summary(
166        &self,
167        mut summary: SessionSummary,
168    ) -> Result<Option<SessionSummary>> {
169        let raw_session = summary.session_id.clone();
170        let Some(session_id) = self.decode_session_key(&raw_session)? else {
171            return Ok(None);
172        };
173        summary.session_id = session_id;
174
175        let current_marker = format!("a7ns1_session_{}_", self.encoded_namespace);
176        if raw_session.starts_with(&current_marker) {
177            summary.agent_id = self.decode_agent(&summary.agent_id)?;
178            if let Some(actor_id) = summary.actor_id.take() {
179                summary.actor_id = Some(self.decode_actor(&actor_id)?);
180            }
181        }
182        Ok(Some(summary))
183    }
184
185    fn summary_matches_identity_filter(summary: &SessionSummary, filter: &SessionFilter) -> bool {
186        filter
187            .agent_id
188            .as_ref()
189            .is_none_or(|agent| &summary.agent_id == agent)
190            && filter
191                .actor_id
192                .as_ref()
193                .is_none_or(|actor| summary.actor_id.as_ref() == Some(actor))
194    }
195}
196
197fn encode_component(value: &str) -> String {
198    const HEX: &[u8; 16] = b"0123456789abcdef";
199    let mut encoded = String::with_capacity(value.len() * 2);
200    for byte in value.as_bytes() {
201        encoded.push(HEX[(byte >> 4) as usize] as char);
202        encoded.push(HEX[(byte & 0x0f) as usize] as char);
203    }
204    encoded
205}
206
207fn decode_component(value: &str) -> Option<String> {
208    if !value.len().is_multiple_of(2) {
209        return None;
210    }
211    let mut decoded = Vec::with_capacity(value.len() / 2);
212    for pair in value.as_bytes().chunks_exact(2) {
213        decoded.push((decode_nibble(pair[0])? << 4) | decode_nibble(pair[1])?);
214    }
215    String::from_utf8(decoded).ok()
216}
217
218fn decode_nibble(value: u8) -> Option<u8> {
219    match value {
220        b'0'..=b'9' => Some(value - b'0'),
221        b'a'..=b'f' => Some(value - b'a' + 10),
222        _ => None,
223    }
224}
225
226#[async_trait]
227impl AgentStorage for NamespacedStorage {
228    fn supports(&self, capability: StorageCapability) -> bool {
229        FORWARDED_CAPABILITIES.contains(&capability) && self.inner.supports(capability)
230    }
231
232    async fn save(&self, session_id: &str, snapshot: &AgentSnapshot) -> Result<()> {
233        self.require(StorageCapability::Snapshot)?;
234        self.inner
235            .save(
236                &self.session_key(session_id),
237                &self.encode_snapshot(snapshot),
238            )
239            .await?;
240        self.inner
241            .delete(&self.legacy_session_key(session_id))
242            .await
243    }
244
245    async fn load(&self, session_id: &str) -> Result<Option<AgentSnapshot>> {
246        self.require(StorageCapability::Snapshot)?;
247        if let Some(snapshot) = self.inner.load(&self.session_key(session_id)).await? {
248            return self.decode_snapshot(snapshot).map(Some);
249        }
250
251        //
252        // Legacy reads remain side-effect free. A later save writes the new key and removes the legacy key.
253        //
254        self.inner.load(&self.legacy_session_key(session_id)).await
255    }
256
257    async fn delete(&self, session_id: &str) -> Result<()> {
258        self.require(StorageCapability::Snapshot)?;
259        self.inner.delete(&self.session_key(session_id)).await?;
260        self.inner
261            .delete(&self.legacy_session_key(session_id))
262            .await
263    }
264
265    async fn list_sessions(&self) -> Result<Vec<String>> {
266        self.require(StorageCapability::Snapshot)?;
267        let mut sessions = BTreeSet::new();
268        for session in self.inner.list_sessions().await? {
269            if let Some(session) = self.decode_session_key(&session)? {
270                sessions.insert(session);
271            }
272        }
273        Ok(sessions.into_iter().collect())
274    }
275
276    async fn save_snapshot_with_metadata(
277        &self,
278        session_id: &str,
279        snapshot: &AgentSnapshot,
280        metadata: &SessionMetadata,
281    ) -> Result<()> {
282        self.require(StorageCapability::SessionMetadata)?;
283        self.inner
284            .save_snapshot_with_metadata(
285                &self.session_key(session_id),
286                &self.encode_snapshot(snapshot),
287                &self.encode_metadata(metadata),
288            )
289            .await?;
290        self.inner
291            .delete(&self.legacy_session_key(session_id))
292            .await
293    }
294
295    async fn save_metadata(&self, session_id: &str, metadata: &SessionMetadata) -> Result<()> {
296        self.require(StorageCapability::SessionMetadata)?;
297        self.inner
298            .save_metadata(
299                &self.session_key(session_id),
300                &self.encode_metadata(metadata),
301            )
302            .await
303    }
304
305    async fn load_metadata(&self, session_id: &str) -> Result<Option<SessionMetadata>> {
306        self.require(StorageCapability::SessionMetadata)?;
307        self.inner
308            .load_metadata(&self.session_key(session_id))
309            .await?
310            .map(|metadata| self.decode_metadata(metadata))
311            .transpose()
312    }
313
314    async fn list_sessions_filtered(&self, filter: &SessionFilter) -> Result<Vec<SessionSummary>> {
315        self.require(StorageCapability::SessionFiltering)?;
316        let mut summaries = Vec::new();
317        for summary in self
318            .inner
319            .list_sessions_filtered(&self.backend_session_filter(filter))
320            .await?
321        {
322            if let Some(summary) = self.decode_session_summary(summary)?
323                && Self::summary_matches_identity_filter(&summary, filter)
324            {
325                summaries.push(summary);
326            }
327        }
328        if let Some(limit) = filter.limit {
329            summaries.truncate(limit);
330        }
331        Ok(summaries)
332    }
333
334    async fn cleanup_expired(&self) -> Result<usize> {
335        Err(AgentError::UnsupportedStorageCapability(
336            StorageCapability::ExpiryCleanup,
337        ))
338    }
339
340    async fn save_facts(&self, agent_id: &str, actor_id: &str, facts: &[KeyFact]) -> Result<()> {
341        self.require(StorageCapability::ActorFacts)?;
342        let facts = facts
343            .iter()
344            .map(|fact| self.encode_fact(fact))
345            .collect::<Vec<_>>();
346        self.inner
347            .save_facts(&self.agent_key(agent_id), &self.actor_key(actor_id), &facts)
348            .await
349    }
350
351    async fn load_facts(&self, agent_id: &str, actor_id: &str) -> Result<Vec<KeyFact>> {
352        self.require(StorageCapability::ActorFacts)?;
353        Ok(self
354            .inner
355            .load_facts(&self.agent_key(agent_id), &self.actor_key(actor_id))
356            .await?
357            .into_iter()
358            .filter_map(|fact| self.decode_fact(fact))
359            .collect())
360    }
361
362    async fn query_facts(&self, agent_id: &str, filter: &FactFilter) -> Result<Vec<KeyFact>> {
363        self.require(StorageCapability::ActorFacts)?;
364        let mut filter = filter.clone();
365        filter.actor_id = filter.actor_id.map(|actor| self.actor_key(&actor));
366        Ok(self
367            .inner
368            .query_facts(&self.agent_key(agent_id), &filter)
369            .await?
370            .into_iter()
371            .filter_map(|fact| self.decode_fact(fact))
372            .collect())
373    }
374
375    async fn delete_fact(&self, agent_id: &str, actor_id: &str, fact_id: &str) -> Result<()> {
376        self.require(StorageCapability::ActorFacts)?;
377        self.inner
378            .delete_fact(
379                &self.agent_key(agent_id),
380                &self.actor_key(actor_id),
381                fact_id,
382            )
383            .await
384    }
385
386    async fn delete_actor_data(&self, agent_id: &str, actor_id: &str) -> Result<()> {
387        self.require(StorageCapability::ActorDataDeletion)?;
388        self.inner
389            .delete_actor_data(&self.agent_key(agent_id), &self.actor_key(actor_id))
390            .await
391    }
392
393    async fn save_relationship(
394        &self,
395        agent_id: &str,
396        actor_id: &str,
397        relationship: &serde_json::Value,
398    ) -> Result<()> {
399        self.require(StorageCapability::ActorRelationships)?;
400        self.inner
401            .save_relationship(
402                &self.agent_key(agent_id),
403                &self.actor_key(actor_id),
404                relationship,
405            )
406            .await
407    }
408
409    async fn load_relationship(
410        &self,
411        agent_id: &str,
412        actor_id: &str,
413    ) -> Result<Option<serde_json::Value>> {
414        self.require(StorageCapability::ActorRelationships)?;
415        self.inner
416            .load_relationship(&self.agent_key(agent_id), &self.actor_key(actor_id))
417            .await
418    }
419
420    async fn list_relationship_actors(&self, agent_id: &str) -> Result<Vec<String>> {
421        self.require(StorageCapability::ActorRelationships)?;
422        Ok(self
423            .inner
424            .list_relationship_actors(&self.agent_key(agent_id))
425            .await?
426            .into_iter()
427            .filter_map(|actor| self.decode_key("actor", &actor))
428            .collect())
429    }
430
431    async fn delete_relationship(&self, agent_id: &str, actor_id: &str) -> Result<()> {
432        self.require(StorageCapability::ActorRelationships)?;
433        self.inner
434            .delete_relationship(&self.agent_key(agent_id), &self.actor_key(actor_id))
435            .await
436    }
437}
438
439#[cfg(test)]
440mod tests {
441    use std::collections::HashMap;
442    use std::sync::atomic::{AtomicUsize, Ordering};
443
444    use ai_agents_core::FactCategory;
445    use chrono::Utc;
446    use parking_lot::RwLock;
447
448    use super::*;
449
450    const ALL_CAPABILITIES: &[StorageCapability] = &[
451        StorageCapability::Snapshot,
452        StorageCapability::SessionMetadata,
453        StorageCapability::SessionFiltering,
454        StorageCapability::ExpiryCleanup,
455        StorageCapability::ActorFacts,
456        StorageCapability::ActorRelationships,
457        StorageCapability::ActorDataDeletion,
458    ];
459
460    #[derive(Default)]
461    struct MemStorage {
462        data: RwLock<HashMap<String, AgentSnapshot>>,
463        metadata: RwLock<HashMap<String, SessionMetadata>>,
464        facts: RwLock<HashMap<(String, String), Vec<KeyFact>>>,
465        relationships: RwLock<HashMap<(String, String), serde_json::Value>>,
466        last_session_filter: RwLock<Option<SessionFilter>>,
467        expiry_calls: AtomicUsize,
468        full_capabilities: bool,
469    }
470
471    impl MemStorage {
472        fn full() -> Self {
473            Self {
474                full_capabilities: true,
475                ..Self::default()
476            }
477        }
478    }
479
480    #[async_trait]
481    impl AgentStorage for MemStorage {
482        fn supports(&self, capability: StorageCapability) -> bool {
483            if self.full_capabilities {
484                ALL_CAPABILITIES.contains(&capability)
485            } else {
486                capability == StorageCapability::Snapshot
487            }
488        }
489
490        async fn save(&self, session_id: &str, snapshot: &AgentSnapshot) -> Result<()> {
491            self.data
492                .write()
493                .insert(session_id.to_string(), snapshot.clone());
494            Ok(())
495        }
496
497        async fn load(&self, session_id: &str) -> Result<Option<AgentSnapshot>> {
498            Ok(self.data.read().get(session_id).cloned())
499        }
500
501        async fn delete(&self, session_id: &str) -> Result<()> {
502            self.data.write().remove(session_id);
503            Ok(())
504        }
505
506        async fn list_sessions(&self) -> Result<Vec<String>> {
507            Ok(self.data.read().keys().cloned().collect())
508        }
509
510        async fn save_snapshot_with_metadata(
511            &self,
512            session_id: &str,
513            snapshot: &AgentSnapshot,
514            metadata: &SessionMetadata,
515        ) -> Result<()> {
516            self.save(session_id, snapshot).await?;
517            self.save_metadata(session_id, metadata).await
518        }
519
520        async fn save_metadata(&self, session_id: &str, metadata: &SessionMetadata) -> Result<()> {
521            self.metadata
522                .write()
523                .insert(session_id.to_string(), metadata.clone());
524            Ok(())
525        }
526
527        async fn load_metadata(&self, session_id: &str) -> Result<Option<SessionMetadata>> {
528            Ok(self.metadata.read().get(session_id).cloned())
529        }
530
531        async fn list_sessions_filtered(
532            &self,
533            filter: &SessionFilter,
534        ) -> Result<Vec<SessionSummary>> {
535            *self.last_session_filter.write() = Some(filter.clone());
536            let metadata = self.metadata.read();
537            let mut summaries = self
538                .data
539                .read()
540                .iter()
541                .filter_map(|(session_id, snapshot)| {
542                    let meta = metadata.get(session_id)?;
543                    if filter
544                        .actor_id
545                        .as_ref()
546                        .is_some_and(|actor| meta.actor_id.as_ref() != Some(actor))
547                        || filter
548                            .agent_id
549                            .as_ref()
550                            .is_some_and(|agent| &snapshot.agent_id != agent)
551                    {
552                        return None;
553                    }
554                    Some(SessionSummary {
555                        session_id: session_id.clone(),
556                        agent_id: snapshot.agent_id.clone(),
557                        actor_id: meta.actor_id.clone(),
558                        tags: meta.tags.clone(),
559                        created_at: meta.created_at,
560                        last_active: meta.last_active,
561                        message_count: meta.message_count,
562                    })
563                })
564                .collect::<Vec<_>>();
565            summaries.sort_by_key(|summary| std::cmp::Reverse(summary.last_active));
566            if let Some(limit) = filter.limit {
567                summaries.truncate(limit);
568            }
569            Ok(summaries)
570        }
571
572        async fn cleanup_expired(&self) -> Result<usize> {
573            self.expiry_calls.fetch_add(1, Ordering::SeqCst);
574            Ok(0)
575        }
576
577        async fn save_facts(
578            &self,
579            agent_id: &str,
580            actor_id: &str,
581            facts: &[KeyFact],
582        ) -> Result<()> {
583            self.facts
584                .write()
585                .insert((agent_id.to_string(), actor_id.to_string()), facts.to_vec());
586            Ok(())
587        }
588
589        async fn load_facts(&self, agent_id: &str, actor_id: &str) -> Result<Vec<KeyFact>> {
590            Ok(self
591                .facts
592                .read()
593                .get(&(agent_id.to_string(), actor_id.to_string()))
594                .cloned()
595                .unwrap_or_default())
596        }
597
598        async fn query_facts(&self, agent_id: &str, filter: &FactFilter) -> Result<Vec<KeyFact>> {
599            Ok(self
600                .facts
601                .read()
602                .iter()
603                .filter(|((stored_agent, stored_actor), _)| {
604                    stored_agent == agent_id
605                        && filter
606                            .actor_id
607                            .as_ref()
608                            .is_none_or(|actor| stored_actor == actor)
609                })
610                .flat_map(|(_, facts)| facts.clone())
611                .collect())
612        }
613
614        async fn delete_fact(&self, agent_id: &str, actor_id: &str, fact_id: &str) -> Result<()> {
615            if let Some(facts) = self
616                .facts
617                .write()
618                .get_mut(&(agent_id.to_string(), actor_id.to_string()))
619            {
620                facts.retain(|fact| fact.id != fact_id);
621            }
622            Ok(())
623        }
624
625        async fn delete_actor_data(&self, agent_id: &str, actor_id: &str) -> Result<()> {
626            self.facts
627                .write()
628                .remove(&(agent_id.to_string(), actor_id.to_string()));
629            Ok(())
630        }
631
632        async fn save_relationship(
633            &self,
634            agent_id: &str,
635            actor_id: &str,
636            relationship: &serde_json::Value,
637        ) -> Result<()> {
638            self.relationships.write().insert(
639                (agent_id.to_string(), actor_id.to_string()),
640                relationship.clone(),
641            );
642            Ok(())
643        }
644
645        async fn load_relationship(
646            &self,
647            agent_id: &str,
648            actor_id: &str,
649        ) -> Result<Option<serde_json::Value>> {
650            Ok(self
651                .relationships
652                .read()
653                .get(&(agent_id.to_string(), actor_id.to_string()))
654                .cloned())
655        }
656
657        async fn list_relationship_actors(&self, agent_id: &str) -> Result<Vec<String>> {
658            Ok(self
659                .relationships
660                .read()
661                .keys()
662                .filter(|(stored_agent, _)| stored_agent == agent_id)
663                .map(|(_, actor)| actor.clone())
664                .collect())
665        }
666
667        async fn delete_relationship(&self, agent_id: &str, actor_id: &str) -> Result<()> {
668            self.relationships
669                .write()
670                .remove(&(agent_id.to_string(), actor_id.to_string()));
671            Ok(())
672        }
673    }
674
675    fn fact(actor_id: &str) -> KeyFact {
676        KeyFact {
677            id: "fact-1".into(),
678            actor_id: Some(actor_id.into()),
679            category: FactCategory::UserContext,
680            content: "context".into(),
681            confidence: 1.0,
682            salience: 1.0,
683            extracted_at: Utc::now(),
684            last_accessed: None,
685            source_message_id: None,
686            source_language: None,
687        }
688    }
689
690    #[tokio::test]
691    async fn parent_and_sibling_keys_are_isolated_and_arbitrary_ids_round_trip() {
692        let inner = Arc::new(MemStorage::full());
693        let first = NamespacedStorage::new(inner.clone(), "parent/child");
694        let second = NamespacedStorage::new(inner.clone(), "parent");
695        let first_session = "sibling/../session\0πŸ˜€";
696        let second_session = "child/sibling/../session\0πŸ˜€";
697
698        first
699            .save(first_session, &AgentSnapshot::new("agent/Ξ±".into()))
700            .await
701            .unwrap();
702        second
703            .save(second_session, &AgentSnapshot::new("agent/Ξ²".into()))
704            .await
705            .unwrap();
706
707        let keys = inner.data.read().keys().cloned().collect::<Vec<_>>();
708        assert_eq!(keys.len(), 2);
709        assert_ne!(keys[0], keys[1]);
710        assert!(keys.iter().all(|key| key.is_ascii() && !key.contains('/')));
711        assert_eq!(first.list_sessions().await.unwrap(), vec![first_session]);
712        assert_eq!(second.list_sessions().await.unwrap(), vec![second_session]);
713        assert!(first.load(second_session).await.unwrap().is_none());
714        assert_eq!(
715            first.load(first_session).await.unwrap().unwrap().agent_id,
716            "agent/Ξ±"
717        );
718    }
719
720    #[tokio::test]
721    async fn capabilities_are_derived_without_expiry_forwarding() {
722        let inner = Arc::new(MemStorage::full());
723        let storage = NamespacedStorage::new(inner.clone(), "agent");
724
725        for capability in FORWARDED_CAPABILITIES {
726            assert!(storage.supports(capability));
727        }
728        assert!(!storage.supports(StorageCapability::ExpiryCleanup));
729        assert!(matches!(
730            storage.cleanup_expired().await,
731            Err(AgentError::UnsupportedStorageCapability(
732                StorageCapability::ExpiryCleanup
733            ))
734        ));
735        assert_eq!(inner.expiry_calls.load(Ordering::SeqCst), 0);
736
737        let snapshot_only = NamespacedStorage::new(Arc::new(MemStorage::default()), "agent");
738        assert!(snapshot_only.supports(StorageCapability::Snapshot));
739        assert!(!snapshot_only.supports(StorageCapability::SessionMetadata));
740        assert!(matches!(
741            snapshot_only
742                .save_metadata("session", &SessionMetadata::default())
743                .await,
744            Err(AgentError::UnsupportedStorageCapability(
745                StorageCapability::SessionMetadata
746            ))
747        ));
748    }
749
750    #[tokio::test]
751    async fn session_agent_and_actor_keys_are_transformed_and_stripped() {
752        let inner = Arc::new(MemStorage::full());
753        let storage = NamespacedStorage::new(inner.clone(), "spawned/πŸ˜€");
754        let sibling = NamespacedStorage::new(inner.clone(), "sibling");
755        let session_id = "session/ι›Ά";
756        let agent_id = "agent/ι›Ά";
757        let actor_id = "actor/ι›Ά";
758        let metadata = SessionMetadata {
759            actor_id: Some(actor_id.into()),
760            actors: vec![actor_id.into(), "other/πŸ˜€".into()],
761            ..SessionMetadata::default()
762        };
763
764        storage
765            .save(session_id, &AgentSnapshot::new(agent_id.into()))
766            .await
767            .unwrap();
768        storage.save_metadata(session_id, &metadata).await.unwrap();
769        sibling
770            .save(
771                "session/sibling",
772                &AgentSnapshot::new("agent/sibling".into()),
773            )
774            .await
775            .unwrap();
776        sibling
777            .save_metadata("session/sibling", &SessionMetadata::default())
778            .await
779            .unwrap();
780
781        let raw_session = storage.session_key(session_id);
782        let raw_snapshot = inner.data.read().get(&raw_session).cloned().unwrap();
783        let raw_metadata = inner.metadata.read().get(&raw_session).cloned().unwrap();
784        assert!(raw_snapshot.agent_id.is_ascii());
785        assert_ne!(raw_snapshot.agent_id, agent_id);
786        assert!(
787            raw_metadata
788                .actor_id
789                .as_ref()
790                .is_some_and(|actor| actor.is_ascii() && actor != actor_id)
791        );
792        let loaded_metadata = storage.load_metadata(session_id).await.unwrap().unwrap();
793        assert_eq!(loaded_metadata.actor_id, metadata.actor_id);
794        assert_eq!(loaded_metadata.actors, metadata.actors);
795        assert_eq!(loaded_metadata.tags, metadata.tags);
796
797        storage
798            .save_facts(agent_id, actor_id, &[fact(actor_id)])
799            .await
800            .unwrap();
801        let loaded_facts = storage.load_facts(agent_id, actor_id).await.unwrap();
802        assert_eq!(loaded_facts[0].actor_id.as_deref(), Some(actor_id));
803        let queried_facts = storage
804            .query_facts(
805                agent_id,
806                &FactFilter {
807                    actor_id: Some(actor_id.into()),
808                    ..FactFilter::default()
809                },
810            )
811            .await
812            .unwrap();
813        assert_eq!(queried_facts[0].actor_id.as_deref(), Some(actor_id));
814        let raw_fact_key = (storage.agent_key(agent_id), storage.actor_key(actor_id));
815        let raw_actor = storage.actor_key(actor_id);
816        assert_eq!(
817            inner.facts.read()[&raw_fact_key][0].actor_id.as_deref(),
818            Some(raw_actor.as_str())
819        );
820
821        let relationship = serde_json::json!({"actor_id": actor_id, "score": 3});
822        storage
823            .save_relationship(agent_id, actor_id, &relationship)
824            .await
825            .unwrap();
826        assert_eq!(
827            storage.load_relationship(agent_id, actor_id).await.unwrap(),
828            Some(relationship)
829        );
830        assert_eq!(
831            storage.list_relationship_actors(agent_id).await.unwrap(),
832            vec![actor_id]
833        );
834
835        let filter = SessionFilter {
836            actor_id: Some(actor_id.into()),
837            agent_id: Some(agent_id.into()),
838            ..SessionFilter::default()
839        };
840        let unfiltered = storage
841            .list_sessions_filtered(&SessionFilter::default())
842            .await
843            .unwrap();
844        assert_eq!(unfiltered.len(), 1);
845        assert_eq!(unfiltered[0].session_id, session_id);
846
847        let summaries = storage.list_sessions_filtered(&filter).await.unwrap();
848        assert_eq!(summaries.len(), 1);
849        assert_eq!(summaries[0].session_id, session_id);
850        assert_eq!(summaries[0].agent_id, agent_id);
851        assert_eq!(summaries[0].actor_id.as_deref(), Some(actor_id));
852        let raw_filter = inner.last_session_filter.read().clone().unwrap();
853        assert_eq!(raw_filter.agent_id, None);
854        assert_eq!(raw_filter.actor_id, None);
855        assert_eq!(raw_filter.limit, None);
856    }
857
858    #[tokio::test]
859    async fn filtered_limit_is_applied_after_namespace_isolation() {
860        let inner = Arc::new(MemStorage::full());
861        let storage = NamespacedStorage::new(inner.clone(), "owned");
862        let sibling = NamespacedStorage::new(inner.clone(), "sibling");
863        let older = Utc::now() - chrono::Duration::minutes(1);
864        let newer = Utc::now();
865        let owned_metadata = SessionMetadata {
866            last_active: older,
867            ..SessionMetadata::default()
868        };
869        let sibling_metadata = SessionMetadata {
870            last_active: newer,
871            ..SessionMetadata::default()
872        };
873
874        storage
875            .save("owned-session", &AgentSnapshot::new("agent".into()))
876            .await
877            .unwrap();
878        storage
879            .save_metadata("owned-session", &owned_metadata)
880            .await
881            .unwrap();
882        sibling
883            .save("sibling-session", &AgentSnapshot::new("agent".into()))
884            .await
885            .unwrap();
886        sibling
887            .save_metadata("sibling-session", &sibling_metadata)
888            .await
889            .unwrap();
890
891        let summaries = storage
892            .list_sessions_filtered(&SessionFilter {
893                limit: Some(1),
894                ..SessionFilter::default()
895            })
896            .await
897            .unwrap();
898
899        assert_eq!(summaries.len(), 1);
900        assert_eq!(summaries[0].session_id, "owned-session");
901        assert_eq!(
902            inner.last_session_filter.read().as_ref().unwrap().limit,
903            None
904        );
905    }
906
907    #[tokio::test]
908    async fn legacy_session_keys_are_read_and_migrated_on_save() {
909        let inner = Arc::new(MemStorage::full());
910        let storage = NamespacedStorage::new(inner.clone(), "child");
911        let legacy_key = "child/legacy";
912        inner
913            .save(legacy_key, &AgentSnapshot::new("legacy-agent".into()))
914            .await
915            .unwrap();
916
917        assert_eq!(
918            storage.load("legacy").await.unwrap().unwrap().agent_id,
919            "legacy-agent"
920        );
921        assert_eq!(storage.list_sessions().await.unwrap(), vec!["legacy"]);
922        assert!(inner.data.read().contains_key(legacy_key));
923
924        storage
925            .save("legacy", &AgentSnapshot::new("new-agent".into()))
926            .await
927            .unwrap();
928        assert!(!inner.data.read().contains_key(legacy_key));
929        assert_eq!(
930            storage.load("legacy").await.unwrap().unwrap().agent_id,
931            "new-agent"
932        );
933    }
934
935    #[cfg(feature = "sqlite")]
936    #[tokio::test]
937    async fn sqlite_legacy_session_keys_are_read_migrated_preferred_and_deleted() {
938        let directory = std::env::temp_dir().join(format!(
939            "ai-agents-namespaced-sqlite-{}",
940            uuid::Uuid::new_v4()
941        ));
942        let path = directory.join("sessions.sqlite");
943        let path = path.to_string_lossy().into_owned();
944
945        {
946            let inner = Arc::new(ai_agents_storage::SqliteStorage::new(&path).await.unwrap());
947            let storage = NamespacedStorage::new(inner.clone(), "child");
948            let legacy_key = "child/legacy";
949            inner
950                .save(legacy_key, &AgentSnapshot::new("legacy-agent".into()))
951                .await
952                .unwrap();
953
954            assert_eq!(
955                storage.load("legacy").await.unwrap().unwrap().agent_id,
956                "legacy-agent"
957            );
958            assert_eq!(storage.list_sessions().await.unwrap(), vec!["legacy"]);
959            assert!(inner.load(legacy_key).await.unwrap().is_some());
960
961            storage
962                .save("legacy", &AgentSnapshot::new("current-agent".into()))
963                .await
964                .unwrap();
965            assert!(inner.load(legacy_key).await.unwrap().is_none());
966            inner
967                .save(legacy_key, &AgentSnapshot::new("stale-agent".into()))
968                .await
969                .unwrap();
970            assert_eq!(
971                storage.load("legacy").await.unwrap().unwrap().agent_id,
972                "current-agent"
973            );
974            assert_eq!(storage.list_sessions().await.unwrap(), vec!["legacy"]);
975
976            storage.delete("legacy").await.unwrap();
977            assert!(storage.load("legacy").await.unwrap().is_none());
978            assert!(inner.load(legacy_key).await.unwrap().is_none());
979
980            drop(storage);
981            inner.close().await;
982        }
983
984        crate::remove_sqlite_test_directory(&directory)
985            .await
986            .unwrap();
987    }
988
989    #[cfg(feature = "redis-storage")]
990    #[tokio::test]
991    #[ignore = "requires a Redis service"]
992    async fn redis_legacy_session_keys_are_read_migrated_preferred_and_deleted() {
993        let url = std::env::var("REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1/".into());
994        let prefix = format!("ai-agents-namespace-test:{}:", uuid::Uuid::new_v4());
995        let inner = Arc::new(
996            ai_agents_storage::RedisStorage::new(&url)
997                .unwrap()
998                .with_prefix(prefix),
999        );
1000        let storage = NamespacedStorage::new(inner.clone(), "child");
1001        let legacy_key = "child/legacy";
1002        inner
1003            .save(legacy_key, &AgentSnapshot::new("legacy-agent".into()))
1004            .await
1005            .unwrap();
1006
1007        assert_eq!(
1008            storage.load("legacy").await.unwrap().unwrap().agent_id,
1009            "legacy-agent"
1010        );
1011        assert_eq!(storage.list_sessions().await.unwrap(), vec!["legacy"]);
1012
1013        storage
1014            .save("legacy", &AgentSnapshot::new("current-agent".into()))
1015            .await
1016            .unwrap();
1017        assert!(inner.load(legacy_key).await.unwrap().is_none());
1018        inner
1019            .save(legacy_key, &AgentSnapshot::new("stale-agent".into()))
1020            .await
1021            .unwrap();
1022        assert_eq!(
1023            storage.load("legacy").await.unwrap().unwrap().agent_id,
1024            "current-agent"
1025        );
1026
1027        storage.delete("legacy").await.unwrap();
1028        assert!(storage.load("legacy").await.unwrap().is_none());
1029        assert!(inner.load(legacy_key).await.unwrap().is_none());
1030    }
1031
1032    #[tokio::test]
1033    async fn malformed_current_namespace_key_fails_closed() {
1034        let inner = Arc::new(MemStorage::full());
1035        let storage = NamespacedStorage::new(inner.clone(), "child");
1036        inner.data.write().insert(
1037            format!("a7ns1_session_{}_not-hex", storage.encoded_namespace),
1038            AgentSnapshot::new("agent".into()),
1039        );
1040
1041        assert!(matches!(
1042            storage.list_sessions().await,
1043            Err(AgentError::Persistence(message)) if message.contains("malformed session key")
1044        ));
1045    }
1046
1047    #[tokio::test]
1048    async fn file_storage_composition_supports_long_flat_namespace_keys() {
1049        let directory = std::env::temp_dir().join(format!(
1050            "ai-agents-namespaced-storage-{}",
1051            uuid::Uuid::new_v4()
1052        ));
1053        let inner = Arc::new(ai_agents_storage::FileStorage::new(&directory));
1054        let storage = NamespacedStorage::new(inner, "n".repeat(128));
1055        let session_id = "session-segment/".repeat(256);
1056        let snapshot = AgentSnapshot::new("agent".repeat(128));
1057
1058        storage.save(&session_id, &snapshot).await.unwrap();
1059        assert_eq!(
1060            storage.list_sessions().await.unwrap(),
1061            vec![session_id.clone()]
1062        );
1063        assert_eq!(
1064            storage.load(&session_id).await.unwrap().unwrap().agent_id,
1065            snapshot.agent_id
1066        );
1067        storage.delete(&session_id).await.unwrap();
1068        assert!(storage.load(&session_id).await.unwrap().is_none());
1069
1070        std::fs::remove_dir_all(directory).unwrap();
1071    }
1072
1073    #[tokio::test]
1074    async fn relationship_payload_is_opaque() {
1075        let inner = Arc::new(MemStorage::full());
1076        let storage = NamespacedStorage::new(inner.clone(), "child");
1077        let payload = serde_json::json!({
1078            "kind": "custom",
1079            "nested": { "actor_id": "payload-owned-value" }
1080        });
1081
1082        storage
1083            .save_relationship("agent", "actor", &payload)
1084            .await
1085            .unwrap();
1086
1087        let raw_key = (storage.agent_key("agent"), storage.actor_key("actor"));
1088        assert_eq!(inner.relationships.read().get(&raw_key), Some(&payload));
1089        assert_eq!(
1090            storage.load_relationship("agent", "actor").await.unwrap(),
1091            Some(payload)
1092        );
1093    }
1094
1095    #[test]
1096    fn component_codec_round_trips_empty_ascii_and_unicode_values() {
1097        for value in ["", "plain", "a/b:c_0", "ι›ΆπŸ˜€\0"] {
1098            let encoded = encode_component(value);
1099            assert!(encoded.is_ascii());
1100            assert_eq!(decode_component(&encoded).as_deref(), Some(value));
1101        }
1102        assert!(decode_component("0").is_none());
1103        assert!(decode_component("zz").is_none());
1104    }
1105}