1use 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
26pub 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#[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#[derive(Clone, Debug)]
62pub struct Meta {
63 pub attribution: Attribution,
64 pub occurred_at: Option<Timestamp>,
66 pub idempotency_key: Option<String>,
68 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#[derive(Clone, Debug)]
121pub struct Committed {
122 pub record: Record,
123 pub event: Event,
124 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 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 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 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 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 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 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 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
428fn 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}