Skip to main content

pardarsh_core/
repository.rs

1//! The command side: every change goes through here.
2//!
3//! The repository turns a command into an [`Event`], projects the next state,
4//! validates it (structure, extensions, relation rules, reference targets),
5//! and commits it with optimistic concurrency and idempotency. It does *not*
6//! decide who may do what: authorization is the caller's responsibility
7//! (see [`crate::access`] for the visibility primitives).
8
9use std::collections::BTreeMap;
10use std::sync::Arc;
11use std::sync::atomic::{AtomicI64, Ordering};
12
13use chrono::{DateTime, Utc};
14
15use crate::error::{Error, Result};
16use crate::event::{Change, Event, EventKind, FieldPath};
17use crate::extension::{ExtValue, extension_references};
18use crate::id::{EventId, NamespacedName, RecordId, RelationId};
19use crate::model::{Attribution, Lifecycle, Record, RecordDraft, RecordTimes, Timestamp, Visibility};
20use crate::projection::{TimelineEntry, apply, timeline};
21use crate::relation::{Relation, RelationDraft};
22use crate::schema::ENVELOPE_SCHEMA_VERSION;
23use crate::storage::{IncomingEdge, RecordQuery, Store};
24use crate::validation::{ValidationErrors, ValidationIssue, Validator};
25
26/// Source of "now". Injected so tests are deterministic.
27pub trait Clock: Send + Sync {
28    fn now(&self) -> Timestamp;
29}
30
31#[derive(Debug, Default, Clone, Copy)]
32pub struct SystemClock;
33
34impl Clock for SystemClock {
35    fn now(&self) -> Timestamp {
36        Utc::now()
37    }
38}
39
40/// A manually advanced clock for tests and fixtures.
41#[derive(Debug, Clone)]
42pub struct ManualClock(Arc<AtomicI64>);
43
44impl ManualClock {
45    pub fn new(start: Timestamp) -> Self {
46        Self(Arc::new(AtomicI64::new(start.timestamp_millis())))
47    }
48
49    pub fn advance(&self, d: chrono::Duration) {
50        self.0.fetch_add(d.num_milliseconds(), Ordering::SeqCst);
51    }
52}
53
54impl Clock for ManualClock {
55    fn now(&self) -> Timestamp {
56        DateTime::from_timestamp_millis(self.0.load(Ordering::SeqCst)).expect("valid timestamp")
57    }
58}
59
60/// Metadata common to every command.
61#[derive(Clone, Debug)]
62pub struct Meta {
63    pub attribution: Attribution,
64    /// When it happened according to the actor; defaults to now.
65    pub occurred_at: Option<Timestamp>,
66    /// Retrying a command with the same key returns the original result.
67    pub idempotency_key: Option<String>,
68    /// Reject the command unless the record is at this version.
69    pub expected_version: Option<u64>,
70    pub reason: Option<String>,
71    pub tag: Option<NamespacedName>,
72    pub visibility: Visibility,
73}
74
75impl Meta {
76    pub fn by(attribution: Attribution) -> Self {
77        Self {
78            attribution,
79            occurred_at: None,
80            idempotency_key: None,
81            expected_version: None,
82            reason: None,
83            tag: None,
84            visibility: Visibility::Public,
85        }
86    }
87
88    pub fn reason(mut self, r: impl Into<String>) -> Self {
89        self.reason = Some(r.into());
90        self
91    }
92
93    pub fn idempotency_key(mut self, k: impl Into<String>) -> Self {
94        self.idempotency_key = Some(k.into());
95        self
96    }
97
98    pub fn expect_version(mut self, v: u64) -> Self {
99        self.expected_version = Some(v);
100        self
101    }
102
103    pub fn tag(mut self, t: NamespacedName) -> Self {
104        self.tag = Some(t);
105        self
106    }
107
108    pub fn occurred_at(mut self, t: Timestamp) -> Self {
109        self.occurred_at = Some(t);
110        self
111    }
112
113    pub fn visibility(mut self, v: Visibility) -> Self {
114        self.visibility = v;
115        self
116    }
117}
118
119/// Result of a command.
120#[derive(Clone, Debug)]
121pub struct Committed {
122    pub record: Record,
123    pub event: Event,
124    /// True when the idempotency key matched an earlier identical command and
125    /// nothing new was written.
126    pub replayed: bool,
127}
128
129pub struct Repository<S: Store> {
130    store: S,
131    validator: Validator,
132    clock: Box<dyn Clock>,
133}
134
135impl<S: Store> Repository<S> {
136    pub fn new(store: S, validator: Validator) -> Self {
137        Self { store, validator, clock: Box::new(SystemClock) }
138    }
139
140    pub fn with_clock(mut self, clock: impl Clock + 'static) -> Self {
141        self.clock = Box::new(clock);
142        self
143    }
144
145    pub fn store(&self) -> &S {
146        &self.store
147    }
148
149    pub fn validator(&self) -> &Validator {
150        &self.validator
151    }
152
153    pub fn now(&self) -> Timestamp {
154        self.clock.now()
155    }
156
157    // ---- queries -------------------------------------------------------
158
159    pub fn get(&self, id: &RecordId) -> Result<Record> {
160        self.store.get(id)?.ok_or_else(|| Error::NotFound(id.clone()))
161    }
162
163    pub fn find(&self, id: &RecordId) -> Result<Option<Record>> {
164        self.store.get(id)
165    }
166
167    pub fn history(&self, id: &RecordId) -> Result<Vec<Event>> {
168        self.store.history(id)
169    }
170
171    pub fn timeline(&self, id: &RecordId) -> Result<Vec<TimelineEntry>> {
172        Ok(timeline(&self.store.history(id)?))
173    }
174
175    pub fn list(&self, query: &RecordQuery) -> Result<Vec<Record>> {
176        self.store.list(query)
177    }
178
179    pub fn incoming(&self, id: &RecordId) -> Result<Vec<IncomingEdge>> {
180        self.store.incoming(id)
181    }
182
183    // ---- commands ------------------------------------------------------
184
185    pub fn create(&mut self, draft: RecordDraft, meta: Meta) -> Result<Committed> {
186        if let Some(c) = self.replay(&meta, None)? {
187            return Ok(c);
188        }
189        let now = self.clock.now();
190        let id = RecordId::generate(self.store.origin());
191        let relations = draft.relations.into_iter().map(|r| r.into_relation(meta.attribution.clone(), now)).collect();
192        let record = Record {
193            id: id.clone(),
194            kind: draft.kind,
195            schema_version: ENVELOPE_SCHEMA_VERSION,
196            title: draft.title,
197            summary: draft.summary,
198            body: draft.body,
199            times: RecordTimes {
200                recorded_at: now,
201                published_at: draft.published_at,
202                effective_from: draft.effective_from,
203                effective_until: draft.effective_until,
204                observed_at: draft.observed_at,
205            },
206            actors: draft.actors,
207            sources: draft.sources,
208            relations,
209            visibility: draft.visibility,
210            lifecycle: Lifecycle::Active,
211            extensions: draft.extensions,
212            redactions: Vec::new(),
213            version: 1,
214            updated_at: now,
215        };
216        self.execute(id, EventKind::Created { record: Box::new(record) }, meta, None)
217    }
218
219    pub fn amend(&mut self, id: &RecordId, changes: Vec<Change>, meta: Meta) -> Result<Committed> {
220        self.run(id, EventKind::Amended { changes }, meta)
221    }
222
223    pub fn correct(
224        &mut self,
225        id: &RecordId,
226        changes: Vec<Change>,
227        corrects: Option<EventId>,
228        meta: Meta,
229    ) -> Result<Committed> {
230        if let Some(target) = corrects {
231            if !self.store.history(id)?.iter().any(|e| e.id == target) {
232                return Err(ValidationErrors::single("corrects", "event is not part of this record's history").into());
233            }
234        }
235        self.run(id, EventKind::Corrected { changes, corrects }, meta)
236    }
237
238    pub fn set_extension_field(
239        &mut self,
240        id: &RecordId,
241        namespace: &str,
242        field: &str,
243        value: Option<ExtValue>,
244        meta: Meta,
245    ) -> Result<Committed> {
246        let change = Change::SetExtensionField { namespace: namespace.into(), field: field.into(), value };
247        self.amend(id, vec![change], meta)
248    }
249
250    pub fn link(&mut self, source: &RecordId, draft: RelationDraft, meta: Meta) -> Result<Committed> {
251        let now = self.clock.now();
252        let relation = draft.into_relation(meta.attribution.clone(), now);
253        self.run(source, EventKind::RelationAdded { relation }, meta)
254    }
255
256    pub fn unlink(&mut self, source: &RecordId, relation_id: RelationId, meta: Meta) -> Result<Committed> {
257        self.run(source, EventKind::RelationRetracted { relation_id }, meta)
258    }
259
260    pub fn retract(&mut self, id: &RecordId, meta: Meta) -> Result<Committed> {
261        self.run(id, EventKind::Retracted, meta)
262    }
263
264    pub fn supersede(&mut self, id: &RecordId, by: &RecordId, meta: Meta) -> Result<Committed> {
265        let replacement = self.get(by)?;
266        let current = self.get(id)?;
267        if replacement.kind != current.kind {
268            return Err(
269                ValidationErrors::single("by", "a record can only be superseded by a record of the same kind").into()
270            );
271        }
272        self.run(id, EventKind::Superseded { by: by.clone() }, meta)
273    }
274
275    pub fn redact(&mut self, id: &RecordId, fields: Vec<FieldPath>, meta: Meta) -> Result<Committed> {
276        self.run(id, EventKind::Redacted { fields }, meta)
277    }
278
279    pub fn annotate(
280        &mut self,
281        id: &RecordId,
282        name: NamespacedName,
283        data: BTreeMap<String, ExtValue>,
284        meta: Meta,
285    ) -> Result<Committed> {
286        self.run(id, EventKind::Annotated { name, data }, meta)
287    }
288
289    /// Commit already-formed events received from elsewhere (import). The
290    /// events keep their ids and times; they are re-projected and
291    /// validated exactly like local commands.
292    pub fn ingest(&mut self, event: Event) -> Result<Record> {
293        let current = self.store.get(&event.record_id)?;
294        let next = apply(current.clone(), &event)?;
295        self.validator.validate_event(&event, event.recorded_at)?;
296        self.validator.validate_record(&next)?;
297        self.store.commit(event, next.clone())?;
298        Ok(next)
299    }
300
301    // ---- internals -----------------------------------------------------
302
303    fn run(&mut self, id: &RecordId, kind: EventKind, meta: Meta) -> Result<Committed> {
304        if let Some(c) = self.replay(&meta, Some((id, &kind)))? {
305            return Ok(c);
306        }
307        let current = self.get(id)?;
308        if let Some(expected) = meta.expected_version {
309            if expected != current.version {
310                return Err(Error::Conflict { id: id.clone(), expected, actual: current.version });
311            }
312        }
313        if current.lifecycle == Lifecycle::Retracted
314            && !matches!(kind, EventKind::Redacted { .. } | EventKind::Annotated { .. } | EventKind::Corrected { .. })
315        {
316            return Err(Error::NotPermitted(format!(
317                "{id} is retracted; only corrections, redactions and annotations are allowed"
318            )));
319        }
320        if matches!(kind, EventKind::Retracted) && current.lifecycle == Lifecycle::Retracted {
321            return Err(Error::NotPermitted(format!("{id} is already retracted")));
322        }
323        self.execute(id.clone(), kind, meta, Some(current))
324    }
325
326    /// If the idempotency key was used before, return the original result
327    /// when the command is the same, or fail when it is different.
328    fn replay(&self, meta: &Meta, target: Option<(&RecordId, &EventKind)>) -> Result<Option<Committed>> {
329        let Some(key) = &meta.idempotency_key else { return Ok(None) };
330        let Some(previous) = self.store.event_by_idempotency_key(key)? else { return Ok(None) };
331        let same = match (target, &previous.kind) {
332            (None, EventKind::Created { .. }) => true,
333            (Some((id, kind)), prev) => &previous.record_id == id && same_command(kind, prev),
334            _ => false,
335        };
336        if !same || previous.attribution != meta.attribution {
337            return Err(Error::IdempotencyConflict(key.clone()));
338        }
339        let record = self.get(&previous.record_id)?;
340        Ok(Some(Committed { record, event: previous, replayed: true }))
341    }
342
343    fn execute(&mut self, id: RecordId, kind: EventKind, meta: Meta, current: Option<Record>) -> Result<Committed> {
344        let now = match &kind {
345            // Keep the snapshot's recorded time and the event's identical.
346            EventKind::Created { record } => record.times.recorded_at,
347            _ => self.clock.now(),
348        };
349        let event = Event {
350            id: EventId::generate(),
351            record_id: id.clone(),
352            sequence: current.as_ref().map_or(1, |r| r.version + 1),
353            kind,
354            attribution: meta.attribution,
355            occurred_at: meta.occurred_at.unwrap_or(now),
356            recorded_at: now,
357            idempotency_key: meta.idempotency_key,
358            reason: meta.reason,
359            tag: meta.tag,
360            visibility: meta.visibility,
361        };
362        self.validator.validate_event(&event, now)?;
363        let next = apply(current.clone(), &event)?;
364        self.validator.validate_record(&next)?;
365        self.check_references(current.as_ref(), &next)?;
366        self.store.commit(event.clone(), next.clone())?;
367        Ok(Committed { record: next, event, replayed: false })
368    }
369
370    /// Local reference targets must exist and satisfy relation/field kind
371    /// rules. Remote targets cannot be checked here and are accepted.
372    fn check_references(&self, before: Option<&Record>, next: &Record) -> Result<()> {
373        let origin = self.store.origin();
374        let mut issues = Vec::new();
375        let is_new = |r: &Relation| before.is_none_or(|b| !b.relations.iter().any(|x| x.id == r.id));
376        for (i, rel) in next.relations.iter().enumerate().filter(|(_, r)| is_new(r)) {
377            let path = format!("relations[{i}]");
378            if !rel.target.is_from(origin) {
379                continue;
380            }
381            let Some(target) = self.store.get(&rel.target)? else {
382                issues.push(ValidationIssue::new(format!("{path}.target"), "target record does not exist"));
383                continue;
384            };
385            if let Some(kinds) = rel.kind.target_kinds() {
386                if !kinds.contains(&target.kind) {
387                    issues.push(ValidationIssue::new(
388                        format!("{path}.target"),
389                        format!(
390                            "{} requires a target of kind {:?}, got {}",
391                            rel.kind,
392                            kinds.iter().map(|k| k.to_string()).collect::<Vec<_>>(),
393                            target.kind
394                        ),
395                    ));
396                }
397            }
398            if rel.kind.requires_same_kind() && target.kind != next.kind {
399                issues.push(ValidationIssue::new(
400                    format!("{path}.target"),
401                    format!("{} requires records of the same kind", rel.kind),
402                ));
403            }
404            let duplicate = next
405                .relations
406                .iter()
407                .any(|r| r.id != rel.id && !r.retracted && r.kind == rel.kind && r.target == rel.target);
408            if duplicate {
409                issues.push(ValidationIssue::new(path, "an identical active relation already exists"));
410            }
411        }
412        for (path, target, kinds) in extension_references(&self.validator.registry, next) {
413            if !target.is_from(origin) {
414                continue;
415            }
416            match self.store.get(target)? {
417                None => issues.push(ValidationIssue::new(path, "referenced record does not exist")),
418                Some(t) if !kinds.is_empty() && !kinds.contains(&t.kind) => {
419                    issues.push(ValidationIssue::new(path, format!("referenced record has kind {}", t.kind)))
420                }
421                Some(_) => {}
422            }
423        }
424        if issues.is_empty() { Ok(()) } else { Err(ValidationErrors { issues }.into()) }
425    }
426}
427
428/// Two commands are "the same" for idempotency if their payloads match,
429/// ignoring ids and times generated while handling them.
430fn same_command(a: &EventKind, b: &EventKind) -> bool {
431    fn normalize(v: serde_json::Value) -> serde_json::Value {
432        match v {
433            serde_json::Value::Object(map) => map
434                .into_iter()
435                .filter(|(k, _)| !matches!(k.as_str(), "id" | "asserted_at"))
436                .map(|(k, v)| (k, normalize(v)))
437                .collect(),
438            serde_json::Value::Array(items) => items.into_iter().map(normalize).collect(),
439            other => other,
440        }
441    }
442    match (serde_json::to_value(a), serde_json::to_value(b)) {
443        (Ok(x), Ok(y)) => normalize(x) == normalize(y),
444        _ => false,
445    }
446}