Skip to main content

pardarsh_core/
projection.rs

1//! Deriving current state from history, and explaining history to humans.
2
3use serde::{Deserialize, Serialize};
4
5use crate::error::Error;
6use crate::event::{Change, Event, EventKind, redact_record};
7use crate::id::NamespacedName;
8use crate::model::{Attribution, Lifecycle, Record, Timestamp, Visibility};
9
10/// Apply one event to the current state.
11pub fn apply(state: Option<Record>, event: &Event) -> Result<Record, Error> {
12    let fail = |m: String| Error::Projection(format!("event {} (seq {}): {m}", event.id, event.sequence));
13    let mut record = match (&event.kind, state) {
14        (EventKind::Created { record }, None) => {
15            if event.sequence != 1 {
16                return Err(fail("created event must have sequence 1".into()));
17            }
18            if record.id != event.record_id {
19                return Err(fail("created snapshot id does not match event record id".into()));
20            }
21            (**record).clone()
22        }
23        (EventKind::Created { .. }, Some(_)) => return Err(fail("record already exists".into())),
24        (_, None) => return Err(fail("first event must be created".into())),
25        (kind, Some(mut record)) => {
26            if record.id != event.record_id {
27                return Err(fail("event belongs to another record".into()));
28            }
29            if event.sequence != record.version + 1 {
30                return Err(fail(format!("expected sequence {}", record.version + 1)));
31            }
32            match kind {
33                EventKind::Created { .. } => unreachable!(),
34                EventKind::Amended { changes } | EventKind::Corrected { changes, .. } => {
35                    for change in changes {
36                        apply_change(&mut record, change).map_err(fail)?;
37                    }
38                }
39                EventKind::RelationAdded { relation } => {
40                    if record.relations.iter().any(|r| r.id == relation.id) {
41                        return Err(fail(format!("relation {} already exists", relation.id)));
42                    }
43                    record.relations.push(relation.clone());
44                }
45                EventKind::RelationRetracted { relation_id } => {
46                    let rel = record
47                        .relations
48                        .iter_mut()
49                        .find(|r| r.id == *relation_id && !r.retracted)
50                        .ok_or_else(|| fail(format!("no active relation {relation_id}")))?;
51                    rel.retracted = true;
52                }
53                EventKind::Retracted => record.lifecycle = Lifecycle::Retracted,
54                EventKind::Superseded { by } => record.lifecycle = Lifecycle::Superseded { by: by.clone() },
55                EventKind::Redacted { fields } => {
56                    redact_record(&mut record, fields);
57                    for f in fields {
58                        if !record.redactions.contains(f) {
59                            record.redactions.push(f.clone());
60                        }
61                    }
62                }
63                EventKind::Annotated { .. } => {}
64            }
65            record
66        }
67    };
68    record.version = event.sequence;
69    record.updated_at = event.recorded_at;
70    Ok(record)
71}
72
73fn apply_change(record: &mut Record, change: &Change) -> Result<(), String> {
74    match change {
75        Change::Title(t) => record.title = t.clone(),
76        Change::Summary(s) => record.summary = s.clone(),
77        Change::Body(b) => record.body = b.clone(),
78        Change::PublishedAt(t) => record.times.published_at = *t,
79        Change::EffectiveFrom(t) => record.times.effective_from = *t,
80        Change::EffectiveUntil(t) => record.times.effective_until = *t,
81        Change::ObservedAt(t) => record.times.observed_at = *t,
82        Change::Visibility(v) => record.visibility = v.clone(),
83        Change::AddActor(a) => record.actors.push(a.clone()),
84        Change::AddSource(s) => {
85            if record.sources.iter().any(|x| x.id == s.id) {
86                return Err(format!("source {} already attached", s.id));
87            }
88            record.sources.push(s.clone())
89        }
90        Change::SetExtensionField { namespace, field, value } => {
91            let ext =
92                record.extensions.get_mut(namespace).ok_or_else(|| format!("record has no extension {namespace}"))?;
93            match value {
94                Some(v) => {
95                    ext.fields.insert(field.clone(), v.clone());
96                }
97                None => {
98                    ext.fields.remove(field);
99                }
100            }
101        }
102    }
103    Ok(())
104}
105
106/// Replay a complete history. Returns `None` for an empty history.
107pub fn project(events: &[Event]) -> Result<Option<Record>, Error> {
108    let mut state = None;
109    for e in events {
110        state = Some(apply(state, e)?);
111    }
112    Ok(state)
113}
114
115/// One human-readable field difference.
116#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
117pub struct FieldDiff {
118    pub field: String,
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub before: Option<String>,
121    #[serde(default, skip_serializing_if = "Option::is_none")]
122    pub after: Option<String>,
123}
124
125/// Field-level differences between two states of a record.
126pub fn diff(before: Option<&Record>, after: &Record) -> Vec<FieldDiff> {
127    let mut out = Vec::new();
128    let mut push = |field: String, b: Option<String>, a: Option<String>| {
129        if b != a {
130            out.push(FieldDiff { field, before: b, after: a });
131        }
132    };
133    let ts = |t: Option<Timestamp>| t.map(|t| t.to_rfc3339());
134    push("title".into(), before.map(|r| r.title.clone()), Some(after.title.clone()));
135    push("summary".into(), before.and_then(|r| r.summary.clone()), after.summary.clone());
136    push("body".into(), before.and_then(|r| r.body.clone()), after.body.clone());
137    push("published_at".into(), before.and_then(|r| ts(r.times.published_at)), ts(after.times.published_at));
138    push("effective_from".into(), before.and_then(|r| ts(r.times.effective_from)), ts(after.times.effective_from));
139    push("effective_until".into(), before.and_then(|r| ts(r.times.effective_until)), ts(after.times.effective_until));
140    push("observed_at".into(), before.and_then(|r| ts(r.times.observed_at)), ts(after.times.observed_at));
141    push(
142        "visibility".into(),
143        before.map(|r| visibility_label(&r.visibility)),
144        Some(visibility_label(&after.visibility)),
145    );
146    if before.is_some() || after.lifecycle != Lifecycle::Active {
147        push(
148            "lifecycle".into(),
149            before.map(|r| lifecycle_label(&r.lifecycle)),
150            Some(lifecycle_label(&after.lifecycle)),
151        );
152    }
153
154    let mut namespaces: Vec<&String> = after.extensions.keys().collect();
155    if let Some(b) = before {
156        namespaces.extend(b.extensions.keys());
157    }
158    namespaces.sort();
159    namespaces.dedup();
160    for ns in namespaces {
161        let b = before.and_then(|r| r.extensions.get(ns));
162        let a = after.extensions.get(ns);
163        let mut fields: Vec<&String> = a.map(|a| a.fields.keys().collect()).unwrap_or_default();
164        if let Some(b) = b {
165            fields.extend(b.fields.keys());
166        }
167        fields.sort();
168        fields.dedup();
169        for f in fields {
170            push(
171                format!("{ns}.{f}"),
172                b.and_then(|e| e.fields.get(f)).map(|v| v.display()),
173                a.and_then(|e| e.fields.get(f)).map(|v| v.display()),
174            );
175        }
176    }
177
178    for s in &after.sources {
179        if !before.is_some_and(|b| b.sources.iter().any(|x| x.id == s.id)) {
180            push("sources".into(), None, Some(s.locator.clone()));
181        }
182    }
183    for r in &after.relations {
184        let prev = before.and_then(|b| b.relations.iter().find(|x| x.id == r.id));
185        match prev {
186            None => push("relations".into(), None, Some(format!("{} {}", r.kind, r.target))),
187            Some(p) if !p.retracted && r.retracted => {
188                push("relations".into(), Some(format!("{} {}", r.kind, r.target)), Some("retracted".into()))
189            }
190            _ => {}
191        }
192    }
193    let before_actors = before.map_or(0, |b| b.actors.len());
194    for a in after.actors.iter().skip(before_actors) {
195        push("actors".into(), None, Some(a.role.to_string()));
196    }
197    out
198}
199
200fn visibility_label(v: &Visibility) -> String {
201    match v {
202        Visibility::Public => "public".into(),
203        Visibility::Restricted { audiences } if audiences.is_empty() => "private".into(),
204        Visibility::Restricted { audiences } => {
205            format!("restricted({})", audiences.iter().cloned().collect::<Vec<_>>().join(","))
206        }
207    }
208}
209
210fn lifecycle_label(l: &Lifecycle) -> String {
211    match l {
212        Lifecycle::Active => "active".into(),
213        Lifecycle::Retracted => "retracted".into(),
214        Lifecycle::Superseded { by } => format!("superseded by {by}"),
215    }
216}
217
218/// A readable history entry.
219#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
220pub struct TimelineEntry {
221    pub sequence: u64,
222    pub event: String,
223    #[serde(default, skip_serializing_if = "Option::is_none")]
224    pub tag: Option<NamespacedName>,
225    pub attribution: Attribution,
226    pub occurred_at: Timestamp,
227    pub recorded_at: Timestamp,
228    #[serde(default, skip_serializing_if = "Option::is_none")]
229    pub reason: Option<String>,
230    pub changes: Vec<FieldDiff>,
231}
232
233/// Replay a history and describe each step. Works on filtered histories
234/// (gaps are tolerated) by diffing only the events present.
235pub fn timeline(events: &[Event]) -> Vec<TimelineEntry> {
236    let mut state: Option<Record> = None;
237    let mut out = Vec::with_capacity(events.len());
238    for e in events {
239        let next = match (&e.kind, &state) {
240            (EventKind::Created { record }, _) => Some((**record).clone()),
241            (_, Some(s)) => {
242                // Tolerate gaps from filtered histories by forcing the sequence.
243                let mut s = s.clone();
244                s.version = e.sequence - 1;
245                apply(Some(s), e).ok()
246            }
247            (_, None) => None,
248        };
249        let changes = match (&next, &e.kind) {
250            (_, EventKind::Annotated { data, .. }) => data
251                .iter()
252                .map(|(k, v)| FieldDiff { field: k.clone(), before: None, after: Some(v.display()) })
253                .collect(),
254            (Some(n), _) => diff(state.as_ref(), n),
255            (None, _) => Vec::new(),
256        };
257        out.push(TimelineEntry {
258            sequence: e.sequence,
259            event: e.kind.name().to_string(),
260            tag: e.tag.clone(),
261            attribution: e.attribution.clone(),
262            occurred_at: e.occurred_at,
263            recorded_at: e.recorded_at,
264            reason: e.reason.clone(),
265            changes,
266        });
267        if next.is_some() {
268            state = next;
269        }
270    }
271    out
272}