1use 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
23pub 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 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(¤t_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 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}