1use 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
10pub 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
106pub 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#[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
125pub 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#[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
233pub 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 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}