Skip to main content

agentplane/keyring/
memory.rs

1//! Cryptographically erasable governed memory for single-node deployments.
2//!
3//! Content is sealed before it reaches `MemoryStore`; metadata remains clear so
4//! subject/purpose indexes and policy remain usable. Every **version** of every
5//! item is sealed in the crate's one envelope under its own tenant-qualified
6//! key scope, `memory-item/<id>@<version>`, so each erasure verb destroys
7//! exactly the keys of what it erased: `forget` every version of one id,
8//! `forget_cascading` and `sweep_expired` every version of each id they erased
9//! whole, `erase_subject` every version of every id the subject holds — and a
10//! cascade that trims an id's superseded versions destroys exactly those
11//! versions' keys, leaving its current one readable. Destroying a scope makes
12//! the live rows, replicas and backups of that version unreadable at once. A
13//! subject is not a key scope, so erasing one does not stop it being written
14//! to again under new ids.
15//!
16//! Subject erasure is serialized with writes and legal-hold changes by this
17//! wrapper. That mutex is process-local, so this concrete adapter is for redb or
18//! another single-writer deployment. Active-active deployments need a
19//! distributed erasure coordinator spanning their database lock and KMS call;
20//! pretending a local mutex supplies that contract would create a hold race.
21
22use std::sync::Arc;
23
24use async_trait::async_trait;
25use serde::{Deserialize, Serialize};
26
27use crate::core::{StoreError, TenantId, Timestamp};
28use crate::journal::payload;
29use crate::memory::{Cascade, MemoryItem, MemoryStore, Recall, Selected};
30
31use super::{Erasure, KeyError, KeyRing};
32
33/// What the envelope seals: the content and the lineage it was derived from.
34///
35/// Integrity is the AEAD's: the envelope authenticates these bytes under the
36/// item's identity, so no digest of the plaintext is stored beside them.
37#[derive(Debug, Serialize, Deserialize)]
38#[serde(deny_unknown_fields)]
39struct PlainMemory {
40    content: serde_json::Value,
41    derived_from: Vec<Selected>,
42}
43
44/// A memory store whose content is unreadable once its items' keys are destroyed.
45pub struct EncryptedMemoryStore {
46    inner: Arc<dyn MemoryStore>,
47    keys: Arc<dyn KeyRing>,
48    tenant: TenantId,
49    lifecycle: Arc<dyn super::ErasureCoordinator>,
50    /// Key destructions an erasure owes: rows already gone whose keys the ring
51    /// refused to destroy.
52    owed: Arc<tokio::sync::Mutex<Vec<OwedKey>>>,
53}
54
55/// One version's key an erasure still has to destroy.
56#[derive(Debug)]
57struct OwedKey {
58    id: String,
59    version: u64,
60    at: Timestamp,
61    reason: String,
62}
63
64impl std::fmt::Debug for EncryptedMemoryStore {
65    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66        f.debug_struct("EncryptedMemoryStore")
67            .field("tenant", &self.tenant)
68            .finish_non_exhaustive()
69    }
70}
71
72impl EncryptedMemoryStore {
73    /// Seal this store's content, serialised by a **process-local** lifecycle
74    /// lock.
75    ///
76    /// Single-node *by default*, not by construction: the lock is a seam. An
77    /// active-active plane calls
78    /// [`coordinated_by`](Self::coordinated_by) with a coordinator that spans
79    /// instances — and [`is_distributed`](Self::is_distributed) is how a caller
80    /// checks which it got, rather than inferring it from a constructor name.
81    ///
82    /// # Panics
83    ///
84    /// If `tenant` is not the tenant `inner` serves — see
85    /// [`SealedCases::wrap`](super::SealedCases::wrap) for why that pair is
86    /// checked rather than trusted.
87    #[must_use]
88    pub fn new(inner: Arc<dyn MemoryStore>, keys: Arc<dyn KeyRing>, tenant: TenantId) -> Self {
89        super::assert_serves(inner.tenant(), &tenant, "memory");
90        Self {
91            inner,
92            keys,
93            tenant,
94            lifecycle: Arc::new(super::LocalCoordinator::new()),
95            owed: Arc::default(),
96        }
97    }
98
99    /// Serialise this store's lifecycle operations with somebody else's lock.
100    ///
101    /// The default is [`LocalCoordinator`](super::LocalCoordinator), which is a
102    /// process-local mutex and therefore correct for a single-writer deployment
103    /// and **nothing else**. An active-active plane supplies a coordinator that
104    /// spans instances — otherwise a write on the second instance lands under a
105    /// scope the first is destroying, and the erasure reports success over a row
106    /// sealed to a key that no longer exists.
107    #[must_use]
108    pub fn coordinated_by(mut self, coordinator: Arc<dyn super::ErasureCoordinator>) -> Self {
109        self.lifecycle = coordinator;
110        self
111    }
112
113    /// Whether this store's lifecycle lock spans instances.
114    ///
115    /// Read at `build`, so a plane wiring a shared store can refuse a
116    /// single-node coordinator rather than discovering it during an erasure
117    /// that already reported success.
118    #[must_use]
119    pub fn is_distributed(&self) -> bool {
120        self.lifecycle.is_distributed()
121    }
122
123    /// The lifecycle lock's scope: one per **tenant**, not per subject.
124    ///
125    /// Per-subject would be finer and is not available: `forget`,
126    /// `forget_cascading` and `set_legal_hold` are addressed by item id, and
127    /// `sweep_expired` spans every subject at once. Looking a subject up to
128    /// decide which lock to take is a read that races the very thing the lock
129    /// protects — so the scope is the widest operation's scope, which is what
130    /// the process-local mutex this replaced was already doing.
131    ///
132    /// The cost is stated rather than implied: `remember` takes this lock on
133    /// **every write**, so all of a tenant's memory writes serialise through
134    /// it, and the coordinator's per-scope granularity buys this wrapper
135    /// nothing — one tenant is one scope. That is a throughput ceiling, not a
136    /// safety gap. What the single scope does *not* protect: nothing — it is
137    /// strictly coarser than any finer scheme; what it forgoes is concurrency
138    /// between one tenant's unrelated subjects. A deployment for which that
139    /// ceiling matters needs id-addressed operations to learn their subject
140    /// transactionally before a finer scope is sound; until then, wider and
141    /// correct beats finer and racy.
142    fn lifecycle_scope(&self) -> String {
143        super::scope(&self.tenant, "memory-lifecycle")
144    }
145
146    /// The erasure scope of one version of one memory id.
147    ///
148    /// Per version rather than per id or per subject: `forget` and the expiry
149    /// sweep erase ids, and a cascade erases ids *and* trims superseded
150    /// versions of ids that stay current. Only a key that seals exactly one
151    /// version can be destroyed by the trim without taking the current version
152    /// along — a per-id key would survive the trim, and a backup would keep
153    /// opening the version it removed. An id belongs to one subject and is never
154    /// reused once erased, and versions only grow, so erasing an id destroys
155    /// the scopes of versions `1..=` its highest.
156    fn scope(&self, id: &str, version: u64) -> String {
157        super::scope(&self.tenant, &format!("memory-item/{id}@{version}"))
158    }
159
160    /// The identity one stored version is sealed to.
161    ///
162    /// A canonical JSON array, so no field's content can spell another's: an
163    /// envelope moved to another id, version, subject, purpose or tenant fails
164    /// to authenticate there rather than opening as that row's content.
165    fn aad(&self, item: &MemoryItem, version: u64) -> Result<Vec<u8>, StoreError> {
166        crate::core::canon::to_bytes(&(
167            "memory",
168            self.tenant.as_str(),
169            item.id.as_str(),
170            version,
171            item.subject.as_str(),
172            item.purpose.as_str(),
173        ))
174        .map_err(|error| StoreError::Backend(error.to_string()))
175    }
176
177    async fn seal(&self, item: &MemoryItem, version: u64) -> Result<serde_json::Value, StoreError> {
178        let plain = crate::core::canon::to_bytes(&PlainMemory {
179            content: item.content.clone(),
180            derived_from: item.derived_from.clone(),
181        })
182        .map_err(|error| StoreError::Backend(error.to_string()))?;
183        let envelope = super::envelope::seal(
184            self.keys.as_ref(),
185            &self.scope(&item.id, version),
186            &self.aad(item, version)?,
187            &plain,
188        )
189        .await
190        .map_err(|error| match error {
191            KeyError::Destroyed { .. } => StoreError::Backend(format!(
192                "memory id '{}' was erased and cannot be reused",
193                item.id
194            )),
195            other => key_error(other),
196        })?;
197        Ok(payload::wrap(&envelope))
198    }
199
200    /// `Ok(None)` when the row's key was destroyed — a completed erasure
201    /// reporting itself, not a fault. Every other failure stays loud: a row
202    /// that is not an envelope, an envelope that does not authenticate under
203    /// this row's identity, a retired key version and an unreachable ring are
204    /// all things someone must be told about, and folding them into the skip
205    /// would make tampering read as erasure.
206    async fn open_item(&self, mut item: MemoryItem) -> Result<Option<MemoryItem>, StoreError> {
207        let envelope = payload::unwrap(&item.content).ok_or_else(|| {
208            StoreError::Backend(
209                "encrypted memory row does not contain a sealed envelope".to_owned(),
210            )
211        })?;
212        let aad = self.aad(&item, item.version)?;
213        let Some(plain) = super::envelope::open_or_erased(self.keys.as_ref(), &aad, &envelope)
214            .await
215            .map_err(key_error)?
216        else {
217            return Ok(None);
218        };
219        let plain: PlainMemory = serde_json::from_slice(&plain)
220            .map_err(|error| StoreError::Backend(format!("encrypted memory: {error}")))?;
221        item.content = plain.content;
222        item.derived_from = plain.derived_from;
223        Ok(Some(item))
224    }
225
226    async fn backing_selection(&self, source: &Selected) -> Result<Selected, StoreError> {
227        let stored = self
228            .inner
229            .version(&source.id, source.version)
230            .await?
231            .ok_or_else(|| {
232                StoreError::Backend(format!(
233                    "derived memory source '{}' version {} is absent",
234                    source.id, source.version
235                ))
236            })?;
237        // An erased source cannot anchor new lineage: deriving from a version
238        // whose key is destroyed would commit to content nobody can verify.
239        let opened = self.open_item(stored.clone()).await?.ok_or_else(|| {
240            StoreError::Backend(format!(
241                "derived memory source '{}' version {} was erased",
242                source.id, source.version
243            ))
244        })?;
245        if opened.selection_digest() != source.digest {
246            return Err(StoreError::Backend(format!(
247                "derived memory source '{}' version {} changed",
248                source.id, source.version
249            )));
250        }
251        Ok(Selected {
252            id: source.id.clone(),
253            version: source.version,
254            digest: stored.selection_digest(),
255        })
256    }
257
258    /// Destroy the keys of versions the inner store already erased — each
259    /// entry an id and the versions of it that went — after whatever earlier
260    /// erasures still owe.
261    ///
262    /// The rows are gone from the live store, so repeating the verb that
263    /// erased them would find nothing: a key the ring refuses stays owed, and
264    /// every later erasure verb destroys it first, failing while it cannot.
265    /// The debt lives in this process; a restart before it is paid loses it,
266    /// and the error names the scope to destroy by hand.
267    async fn destroy_erased(
268        &self,
269        erased: &[(String, Vec<u64>)],
270        at: Timestamp,
271        reason: &str,
272    ) -> Result<(), StoreError> {
273        let mut owed = self.owed.lock().await;
274        for (id, versions) in erased {
275            owed.extend(versions.iter().map(|version| OwedKey {
276                id: id.clone(),
277                version: *version,
278                at,
279                reason: reason.to_owned(),
280            }));
281        }
282        while let Some(key) = owed.first() {
283            let scope = self.scope(&key.id, key.version);
284            self.keys
285                .destroy(&scope, key.at, &key.reason)
286                .await
287                .map_err(|error| {
288                    StoreError::Backend(format!(
289                        "memory '{}' version {} was erased from the store, and destroying its \
290                         key failed ({error}) — its backups still open until scope '{scope}' is \
291                         destroyed, which the next erasure retries",
292                        key.id, key.version
293                    ))
294                })?;
295            owed.remove(0);
296        }
297        Ok(())
298    }
299
300    /// Every version of each id, up to the highest it held.
301    fn every_version(ids: impl IntoIterator<Item = (String, u64)>) -> Vec<(String, Vec<u64>)> {
302        ids.into_iter()
303            .map(|(id, highest)| (id, (1..=highest).collect()))
304            .collect()
305    }
306
307    /// Each of `ids` with the highest version the store holds for it, read
308    /// before an erasure removes the rows that say.
309    async fn highest_versions(&self, ids: &[String]) -> Result<Vec<(String, u64)>, StoreError> {
310        let mut out = Vec::with_capacity(ids.len());
311        for id in ids {
312            if let Some(current) = self.inner.current(id, None).await? {
313                out.push((id.clone(), current.version));
314            }
315        }
316        Ok(out)
317    }
318
319    /// Refuse when any of `ids` is under legal hold.
320    async fn refuse_held(&self, ids: &[String]) -> Result<(), StoreError> {
321        for id in ids {
322            if self.inner.legal_hold(id).await? {
323                return Err(StoreError::UnderLegalHold { id: id.clone() });
324            }
325        }
326        Ok(())
327    }
328
329    /// Destroy every item key of a subject, then clean unreadable ciphertext.
330    ///
331    /// Holds are checked first, and a held item refuses the whole erasure with
332    /// [`StoreError::UnderLegalHold`] before any key is touched. Then each of
333    /// the subject's ids has its key destroyed — the erasure, reaching every
334    /// copy — and only then are the live rows removed. A cleanup failure after
335    /// the keys are gone is reported in [`Erasure::cleanup_failed`], not as
336    /// success. The subject stays writable under new ids.
337    ///
338    /// `at` and `reason` come from the caller's audited lifecycle operation;
339    /// this adapter never reads an ambient clock.
340    ///
341    /// # Errors
342    ///
343    /// [`StoreError::UnderLegalHold`] naming a held item, or a failure to read
344    /// the subject or destroy a key.
345    pub async fn erase_subject(
346        &self,
347        subject: &str,
348        at: Timestamp,
349        reason: &str,
350    ) -> Result<Erasure, StoreError> {
351        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
352            // The subject's ids, enumerated by the dedicated erasure-path
353            // operation rather than by a recall with an enormous limit. A
354            // recall is a bounded content query, and a backend may cap or
355            // refuse extreme limits — the PostgreSQL store refuses anything
356            // past BIGINT — so an erasure riding one either failed outright or
357            // silently checked holds for a truncated page of the subject.
358            let ids = self.inner.subject_ids(subject).await?;
359            self.refuse_held(&ids).await?;
360            self.destroy_erased(&[], at, reason).await?;
361            for (id, versions) in Self::every_version(self.highest_versions(&ids).await?) {
362                for version in versions {
363                    self.keys
364                        .destroy(&self.scope(&id, version), at, reason)
365                        .await
366                        .map_err(key_error)?;
367                }
368            }
369            Ok(match self.inner.forget_subject(subject).await {
370                Ok(count) => Erasure {
371                    reached: count,
372                    cleanup_failed: None,
373                },
374                Err(error) => {
375                    tracing::warn!(%subject, %error, "memory keys were destroyed but ciphertext cleanup failed");
376                    Erasure {
377                        reached: ids.len(),
378                        cleanup_failed: Some(error.to_string()),
379                    }
380                }
381            })
382        })
383        .await
384    }
385}
386
387#[allow(clippy::needless_pass_by_value)]
388fn key_error(error: KeyError) -> StoreError {
389    StoreError::Backend(error.to_string())
390}
391
392/// The instant and reason a key destruction records for a verb that carries
393/// neither: the trait's erasure verbs are addressed by id alone.
394///
395/// A wall-clock read, because the instant is the key ring's own record of when
396/// the key went — descriptive metadata no run reads back and nothing replays.
397/// Erasures that carry their caller's instant (`erase_subject`, the sweep)
398/// record that instead.
399#[allow(clippy::disallowed_methods)]
400fn verb_erasure(verb: &str) -> (Timestamp, String) {
401    (Timestamp::now_utc(), format!("memory {verb}"))
402}
403
404#[async_trait]
405impl MemoryStore for EncryptedMemoryStore {
406    fn tenant(&self) -> &str {
407        self.tenant.as_str()
408    }
409
410    /// This store *does* have a lifecycle lock, so the answer is never `None` —
411    /// and whether it spans instances is the coordinator's to say.
412    fn erasure_is_distributed(&self) -> Option<bool> {
413        Some(self.lifecycle.is_distributed())
414    }
415
416    fn seals(&self) -> bool {
417        true
418    }
419
420    fn erasure_index(&self) -> Option<Arc<dyn crate::memory::SemanticRetriever>> {
421        self.inner.erasure_index()
422    }
423
424    async fn remember(&self, item: &MemoryItem) -> Result<u64, StoreError> {
425        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
426            // The version is part of what the envelope is sealed to, and the
427            // store assigns it: the next after the current one. Predicted here
428            // under the lifecycle lock every write takes, and checked against
429            // the store's answer, so a disagreement fails this write rather
430            // than leaving a row that will never authenticate.
431            let version = self
432                .inner
433                .current(&item.id, None)
434                .await?
435                .map_or(1, |current| current.version + 1);
436            let mut sealed = item.clone();
437            sealed.content = self.seal(item, version).await?;
438            sealed.derived_from.clear();
439            for source in &item.derived_from {
440                sealed
441                    .derived_from
442                    .push(self.backing_selection(source).await?);
443            }
444            let written = self.inner.remember(&sealed).await?;
445            if written != version {
446                return Err(StoreError::Backend(format!(
447                    "memory '{}' was sealed as version {version} and stored as {written}; \
448                     the row will not open, so the write is refused",
449                    item.id
450                )));
451            }
452            Ok(written)
453        })
454        .await
455    }
456
457    // The read paths follow the skip-sealed convention SealedJournal and
458    // SealedCases set: a row whose key was **destroyed** is a completed
459    // erasure, and a completed erasure must not turn every later query about
460    // the subject into a persistent error — which is exactly what happens when
461    // an erasure destroys the keys and the ciphertext cleanup then fails.
462    // Destroyed rows are silently absent from `recall` and `derivatives`, and
463    // `version` answers `None` as it does for any erased version. What the
464    // skip does **not** cover: rows that are not envelopes, envelopes that do
465    // not authenticate under their row's identity, and an unreachable key ring
466    // all stay loud, because those are faults to page about rather than
467    // erasures reporting themselves.
468    async fn recall(&self, query: &Recall) -> Result<Vec<MemoryItem>, StoreError> {
469        // A destroyed row is one the inner store's cleanup has not removed
470        // yet, so a page can come back shorter than `limit` while more
471        // readable rows exist past it — an erasure mid-cleanup, not a complete
472        // answer, and it ends when the cleanup is retried.
473        let items = self.inner.recall(query).await?;
474        let mut opened = Vec::with_capacity(items.len());
475        for item in items {
476            if let Some(item) = self.open_item(item).await? {
477                opened.push(item);
478            }
479        }
480        Ok(opened)
481    }
482
483    async fn subject_ids(&self, subject: &str) -> Result<Vec<String>, StoreError> {
484        // Ids are metadata and never sealed, so there is nothing to open —
485        // and erasure needs this to work *after* the keys are gone.
486        self.inner.subject_ids(subject).await
487    }
488
489    async fn version(&self, id: &str, version: u64) -> Result<Option<MemoryItem>, StoreError> {
490        match self.inner.version(id, version).await? {
491            Some(item) => self.open_item(item).await,
492            None => Ok(None),
493        }
494    }
495
496    async fn current(
497        &self,
498        id: &str,
499        as_of: Option<Timestamp>,
500    ) -> Result<Option<MemoryItem>, StoreError> {
501        match self.inner.current(id, as_of).await? {
502            Some(item) => self.open_item(item).await,
503            None => Ok(None),
504        }
505    }
506
507    /// Destroys the id's key, then removes its rows.
508    ///
509    /// The hold is checked first, so a held id loses nothing. A failure to
510    /// remove the rows after the key is gone is an error naming it; retrying
511    /// `forget` completes it.
512    async fn forget(&self, id: &str) -> Result<(), StoreError> {
513        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
514            self.refuse_held(&[id.to_owned()]).await?;
515            let (at, reason) = verb_erasure("forget");
516            self.destroy_erased(&[], at, &reason).await?;
517            for (id, versions) in
518                Self::every_version(self.highest_versions(&[id.to_owned()]).await?)
519            {
520                for version in versions {
521                    self.keys
522                        .destroy(&self.scope(&id, version), at, &reason)
523                        .await
524                        .map_err(key_error)?;
525                }
526            }
527            self.inner.forget(id).await.map_err(|error| {
528                StoreError::Backend(format!(
529                    "memory '{id}': its key is destroyed, so no copy opens, and removing its \
530                     rows failed ({error}) — retry forget to remove them"
531                ))
532            })
533        })
534        .await
535    }
536
537    async fn forget_subject(&self, subject: &str) -> Result<usize, StoreError> {
538        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
539            let ids = self.inner.subject_ids(subject).await?;
540            let highest = self.highest_versions(&ids).await?;
541            let count = self.inner.forget_subject(subject).await?;
542            let (at, reason) = verb_erasure("forget_subject");
543            self.destroy_erased(&Self::every_version(highest), at, &reason)
544                .await?;
545            Ok(count)
546        })
547        .await
548    }
549
550    async fn derivatives(&self, id: &str) -> Result<Vec<MemoryItem>, StoreError> {
551        let items = self.inner.derivatives(id).await?;
552        let mut opened = Vec::with_capacity(items.len());
553        for item in items {
554            if let Some(item) = self.open_item(item).await? {
555                opened.push(item);
556            }
557        }
558        Ok(opened)
559    }
560
561    /// Destroys the key of every version of every id the cascade erased
562    /// whole, and of exactly the versions it trimmed from ids that stay
563    /// current — whose current version keeps its own key and stays readable,
564    /// while a backup taken before the cascade no longer opens what it
565    /// removed.
566    async fn forget_cascading(&self, id: &str) -> Result<Cascade, StoreError> {
567        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
568            let cascade = self.inner.forget_cascading(id).await?;
569            let (at, reason) = verb_erasure("forget_cascading");
570            self.destroy_erased(&Self::every_version(cascade.erased.clone()), at, &reason)
571                .await?;
572            self.destroy_erased(&cascade.trimmed, at, &reason).await?;
573            Ok(cascade)
574        })
575        .await
576    }
577
578    async fn set_legal_hold(&self, id: &str, held: bool) -> Result<(), StoreError> {
579        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
580            self.inner.set_legal_hold(id, held).await
581        })
582        .await
583    }
584
585    async fn legal_holds(
586        &self,
587        after: Option<&str>,
588        limit: usize,
589    ) -> Result<Vec<String>, StoreError> {
590        // Ids are not sealed — only content is — so the listing passes straight
591        // through. Sealing an id would make the hold register unreadable
592        // without the very key an erasure destroys.
593        self.inner.legal_holds(after, limit).await
594    }
595
596    async fn legal_hold(&self, id: &str) -> Result<bool, StoreError> {
597        self.inner.legal_hold(id).await
598    }
599
600    /// Destroys the key of every version of every id the sweep erased.
601    async fn sweep_expired(&self, at: Timestamp) -> Result<Vec<(String, u64)>, StoreError> {
602        super::under_lock(self.lifecycle.as_ref(), &self.lifecycle_scope(), || async {
603            let swept = self.inner.sweep_expired(at).await?;
604            self.destroy_erased(
605                &Self::every_version(swept.clone()),
606                at,
607                "memory retention expired",
608            )
609            .await?;
610            Ok(swept)
611        })
612        .await
613    }
614
615    async fn touch(&self, ids: &[String], at: Timestamp) -> Result<(), StoreError> {
616        self.inner.touch(ids, at).await
617    }
618}