Skip to main content

pardarsh_core/storage/
mod.rs

1//! Storage interface and reference implementations.
2//!
3//! Stores persist events and the projected snapshot of each record. They are
4//! deliberately dumb: validation, authorization and projection happen in the
5//! [`crate::repository::Repository`]. A store's only semantic duties are
6//! atomicity, optimistic concurrency (gap-free sequences), idempotency-key
7//! uniqueness, and scrubbing history when a redaction is committed.
8
9mod memory;
10#[cfg(feature = "sqlite")]
11mod sqlite;
12
13pub use memory::MemoryStore;
14#[cfg(feature = "sqlite")]
15pub use sqlite::SqliteStore;
16
17use crate::error::Result;
18use crate::event::Event;
19use crate::id::{Origin, RecordId};
20use crate::model::{Lifecycle, Record, RecordKind};
21use crate::relation::Relation;
22
23/// Filter for listing records. Text search is an application concern.
24#[derive(Clone, Debug, Default)]
25pub struct RecordQuery {
26    /// Only these kinds (empty = all kinds).
27    pub kinds: Vec<RecordKind>,
28    /// Only records carrying this extension namespace.
29    pub extension: Option<String>,
30    /// Include retracted and superseded records.
31    pub include_inactive: bool,
32    pub limit: Option<usize>,
33    pub offset: usize,
34}
35
36impl RecordQuery {
37    pub fn kind(kind: RecordKind) -> Self {
38        Self { kinds: vec![kind], ..Self::default() }
39    }
40
41    pub(crate) fn matches(&self, r: &Record) -> bool {
42        (self.kinds.is_empty() || self.kinds.contains(&r.kind))
43            && self.extension.as_ref().is_none_or(|ns| r.extensions.contains_key(ns))
44            && (self.include_inactive || r.lifecycle == Lifecycle::Active)
45    }
46
47    pub(crate) fn page<I: Iterator<Item = Record>>(&self, it: I) -> Vec<Record> {
48        let it = it.filter(|r| self.matches(r)).skip(self.offset);
49        match self.limit {
50            Some(n) => it.take(n).collect(),
51            None => it.collect(),
52        }
53    }
54}
55
56/// An edge pointing *at* a record, as seen from the target.
57#[derive(Clone, Debug, PartialEq)]
58pub struct IncomingEdge {
59    pub source: RecordId,
60    pub source_kind: RecordKind,
61    pub relation: Relation,
62}
63
64pub trait Store: Send {
65    /// The authority this store mints identifiers for.
66    fn origin(&self) -> &Origin;
67
68    /// Current snapshot of a record.
69    fn get(&self, id: &RecordId) -> Result<Option<Record>>;
70
71    /// Full history of a record, ordered by sequence.
72    fn history(&self, id: &RecordId) -> Result<Vec<Event>>;
73
74    fn event_by_idempotency_key(&self, key: &str) -> Result<Option<Event>>;
75
76    /// Atomically append `event` and replace the snapshot with `snapshot`.
77    ///
78    /// Must fail with [`crate::Error::Conflict`] unless `event.sequence` is
79    /// exactly one more than the stored version, and with
80    /// [`crate::Error::IdempotencyConflict`] if the event's idempotency key is
81    /// already used. When `event` is a redaction, earlier events of the record
82    /// must be scrubbed with [`crate::event::scrub_event`] in the same
83    /// transaction.
84    fn commit(&mut self, event: Event, snapshot: Record) -> Result<()>;
85
86    /// Records matching the query, newest first.
87    fn list(&self, query: &RecordQuery) -> Result<Vec<Record>>;
88
89    /// Relations (including retracted ones) whose target is `target`.
90    fn incoming(&self, target: &RecordId) -> Result<Vec<IncomingEdge>>;
91}
92
93pub(crate) fn check_sequence(id: &RecordId, current: u64, event: &Event, snapshot: &Record) -> Result<()> {
94    if event.sequence != current + 1 || snapshot.version != event.sequence || &snapshot.id != id {
95        return Err(crate::Error::Conflict {
96            id: id.clone(),
97            expected: event.sequence.saturating_sub(1),
98            actual: current,
99        });
100    }
101    Ok(())
102}