Skip to main content

aep_agent/
store.rs

1use std::{
2    collections::BTreeMap,
3    sync::{Arc, Mutex},
4};
5
6use async_trait::async_trait;
7use time::OffsetDateTime;
8use url::Url;
9use uuid::Uuid;
10
11use crate::{
12    AgentError, AgentIdentity, Clock, CredentialRecord, CredentialStore, IdempotencyKeyProvider,
13    IdentityStore, InspectCache, InspectCacheEntry, OperationKey,
14};
15
16#[derive(Default)]
17pub struct SystemClock;
18
19impl Clock for SystemClock {
20    fn now(&self) -> OffsetDateTime {
21        OffsetDateTime::now_utc()
22    }
23}
24
25#[derive(Default)]
26pub struct TimerDelay;
27
28#[async_trait]
29impl crate::Delay for TimerDelay {
30    async fn sleep(&self, duration: std::time::Duration) {
31        futures_timer::Delay::new(duration).await;
32    }
33}
34
35#[derive(Default)]
36pub struct RandomIdempotencyKeyProvider;
37
38#[async_trait]
39impl IdempotencyKeyProvider for RandomIdempotencyKeyProvider {
40    async fn create_key(&self, _operation: &OperationKey) -> Result<String, AgentError> {
41        Ok(Uuid::new_v4().to_string())
42    }
43}
44
45#[derive(Default)]
46pub struct MemoryIdentityStore {
47    entries: Mutex<BTreeMap<String, AgentIdentity>>,
48}
49
50#[async_trait]
51impl IdentityStore for MemoryIdentityStore {
52    async fn find(&self, service_did: &str) -> Result<Option<AgentIdentity>, AgentError> {
53        Ok(self
54            .entries
55            .lock()
56            .map_err(lock_error)?
57            .get(service_did)
58            .cloned())
59    }
60    async fn save(&self, identity: AgentIdentity) -> Result<(), AgentError> {
61        self.entries
62            .lock()
63            .map_err(lock_error)?
64            .insert(identity.service_did.clone(), identity);
65        Ok(())
66    }
67}
68
69pub struct MemoryCredentialStore {
70    clock: Arc<dyn Clock>,
71    entries: Mutex<BTreeMap<(String, String), CredentialRecord>>,
72}
73
74impl MemoryCredentialStore {
75    pub fn new(clock: Arc<dyn Clock>) -> Self {
76        Self {
77            clock,
78            entries: Mutex::new(BTreeMap::new()),
79        }
80    }
81}
82
83#[async_trait]
84impl CredentialStore for MemoryCredentialStore {
85    async fn delete(&self, service_did: &str, credential_id: &str) -> Result<(), AgentError> {
86        self.entries
87            .lock()
88            .map_err(lock_error)?
89            .remove(&(service_did.to_owned(), credential_id.to_owned()));
90        Ok(())
91    }
92    async fn find(
93        &self,
94        service_did: &str,
95        credential_id: &str,
96    ) -> Result<Option<CredentialRecord>, AgentError> {
97        let key = (service_did.to_owned(), credential_id.to_owned());
98        let mut entries = self.entries.lock().map_err(lock_error)?;
99        if entries
100            .get(&key)
101            .is_some_and(|record| record.expires_at <= self.clock.now())
102        {
103            entries.remove(&key);
104        }
105        Ok(entries.get(&key).cloned())
106    }
107    async fn list(&self, service_did: &str) -> Result<Vec<CredentialRecord>, AgentError> {
108        let now = self.clock.now();
109        let mut entries = self.entries.lock().map_err(lock_error)?;
110        entries.retain(|_, record| record.expires_at > now);
111        let mut records = entries
112            .values()
113            .filter(|record| record.service_did == service_did)
114            .cloned()
115            .collect::<Vec<_>>();
116        records.sort_by(|left, right| {
117            right
118                .issued_at
119                .cmp(&left.issued_at)
120                .then_with(|| left.credential_id.cmp(&right.credential_id))
121        });
122        Ok(records)
123    }
124    async fn save(&self, credential: CredentialRecord) -> Result<(), AgentError> {
125        crate::authentication::validate_record(
126            &credential,
127            &credential.service_did,
128            self.clock.now(),
129        )?;
130        let key = (
131            credential.service_did.clone(),
132            credential.credential_id.clone(),
133        );
134        self.entries
135            .lock()
136            .map_err(lock_error)?
137            .insert(key, credential);
138        Ok(())
139    }
140}
141
142#[derive(Default)]
143pub struct MemoryInspectCache {
144    entries: Mutex<BTreeMap<String, InspectCacheEntry>>,
145}
146
147#[async_trait]
148impl InspectCache for MemoryInspectCache {
149    async fn delete(&self, inspect_url: &Url) -> Result<(), AgentError> {
150        self.entries
151            .lock()
152            .map_err(lock_error)?
153            .remove(inspect_url.as_str());
154        Ok(())
155    }
156    async fn find(&self, inspect_url: &Url) -> Result<Option<InspectCacheEntry>, AgentError> {
157        Ok(self
158            .entries
159            .lock()
160            .map_err(lock_error)?
161            .get(inspect_url.as_str())
162            .cloned())
163    }
164    async fn save(&self, inspect_url: &Url, entry: InspectCacheEntry) -> Result<(), AgentError> {
165        self.entries
166            .lock()
167            .map_err(lock_error)?
168            .insert(inspect_url.to_string(), entry);
169        Ok(())
170    }
171}
172
173fn lock_error<T>(_error: std::sync::PoisonError<T>) -> AgentError {
174    AgentError::Store("AEP Agent memory store lock is poisoned".to_owned())
175}