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}