Skip to main content

made_core/entities/
ceremony_event_reader.rs

1//! [`CeremonyEventReader`] — where a stored payload becomes a typed
2//! event, and where old shapes will be brought forward.
3
4use crate::error::DomainError;
5use crate::value_objects::{AuditEventType, EventSchemaVersion};
6
7use super::CeremonyEvent;
8
9/// Reads a sealed payload back as the event its record names.
10///
11/// The one seam through which raw event JSON enters the domain. A
12/// payload is read under the schema version its record was sealed
13/// with; when a payload shape changes, its new version deserializes
14/// directly and the old one gets an upcaster here, so records written
15/// under earlier shapes keep reading. State-repeat coordinates are the
16/// first versioned evolution: affected payloads use version 2 when the
17/// coordinate is explicitly present, while their version-1 shape remains
18/// readable without it.
19#[derive(Debug)]
20pub struct CeremonyEventReader;
21
22impl CeremonyEventReader {
23    /// Read `raw` as the event `event_type` at payload version `version`.
24    ///
25    /// Refuses a version no reader exists for, a payload that does not
26    /// deserialize as an event, and a payload whose tag names a
27    /// different type than the record does — each naming the type and
28    /// version so an operator can tell which record and which reader
29    /// disagree.
30    pub fn read(
31        event_type: AuditEventType,
32        version: EventSchemaVersion,
33        raw: serde_json::Value,
34    ) -> Result<CeremonyEvent, DomainError> {
35        let supported = match version {
36            EventSchemaVersion::V1 => true,
37            EventSchemaVersion::V2 => matches!(
38                event_type,
39                AuditEventType::CeremonyInstanceStarted
40                    | AuditEventType::StepLeaseRenewed
41                    | AuditEventType::StepStarted
42                    | AuditEventType::StepCompleted
43                    | AuditEventType::StepFailed
44                    | AuditEventType::TransitionApplied
45                    | AuditEventType::ContextWritten
46                    | AuditEventType::StateIterationStarted
47            ),
48            EventSchemaVersion::V3 => matches!(
49                event_type,
50                AuditEventType::CeremonyInstanceStarted
51                    | AuditEventType::StepStarted
52                    | AuditEventType::StepCompleted
53                    | AuditEventType::StepFailed
54                    | AuditEventType::TransitionApplied
55            ),
56            EventSchemaVersion::V4 => matches!(
57                event_type,
58                AuditEventType::CeremonyInstanceStarted
59                    | AuditEventType::StepStarted
60                    | AuditEventType::StepFailed
61            ),
62            EventSchemaVersion::V5 | EventSchemaVersion::V6 => {
63                event_type == AuditEventType::StepStarted
64            }
65            _ => false,
66        };
67        if !supported {
68            return Err(unreadable(
69                event_type,
70                version,
71                "no reader exists for this schema version",
72            ));
73        }
74        let event = serde_json::from_value::<CeremonyEvent>(raw).map_err(|_| {
75            unreadable(
76                event_type,
77                version,
78                "the payload does not deserialize as a ceremony event",
79            )
80        })?;
81        if event.event_type() != event_type {
82            return Err(unreadable(
83                event_type,
84                version,
85                "the payload's tag names a different event type",
86            ));
87        }
88        validate_seals(&event, event_type, version)?;
89        let coordinates_valid = match &event {
90            CeremonyEvent::StepStarted(e) => e.state_visit.is_none() || e.state_iteration.is_some(),
91            CeremonyEvent::StepCompleted(e) => {
92                e.state_visit.is_none() || e.state_iteration.is_some()
93            }
94            CeremonyEvent::StepFailed(e) => e.state_visit.is_none() || e.state_iteration.is_some(),
95            CeremonyEvent::TransitionApplied(e) => match &e.destination {
96                Some(destination) => {
97                    e.transition.has_explicit_state_visit()
98                        && e.transition.has_explicit_state_iteration()
99                        && e.transition.state_visit().next().ok() == Some(destination.state_visit)
100                        && destination
101                            .step_ids
102                            .iter()
103                            .collect::<std::collections::BTreeSet<_>>()
104                            .len()
105                            == destination.step_ids.len()
106                }
107                None => !e.transition.has_explicit_state_visit(),
108            },
109            _ => true,
110        };
111        if !coordinates_valid {
112            return Err(unreadable(
113                event_type,
114                version,
115                "inconsistent state visit coordinates or destination reset",
116            ));
117        }
118        if event.schema_version() != version {
119            return Err(unreadable(
120                event_type,
121                version,
122                "the payload shape does not match its schema version",
123            ));
124        }
125        Ok(event)
126    }
127}
128
129fn unreadable(
130    event_type: AuditEventType,
131    version: EventSchemaVersion,
132    reason: &'static str,
133) -> DomainError {
134    DomainError::UnreadableCeremonyEvent {
135        event_type: event_type.as_str(),
136        version: version.get(),
137        reason,
138    }
139}
140
141fn validate_seals(
142    event: &CeremonyEvent,
143    event_type: AuditEventType,
144    version: EventSchemaVersion,
145) -> Result<(), DomainError> {
146    if let CeremonyEvent::StepCompleted(completed) = event {
147        if completed.result.failure_kind().is_some() {
148            return Err(unreadable(
149                event_type,
150                version,
151                "step completion cannot carry a failure kind",
152            ));
153        }
154    }
155    if let CeremonyEvent::StepFailed(failed) = event {
156        if failed.result.failure_kind().is_some()
157            && failed.result.status() != crate::value_objects::StepStatus::Failed
158        {
159            return Err(unreadable(
160                event_type,
161                version,
162                "only failed results carry a failure kind",
163            ));
164        }
165    }
166    if let CeremonyEvent::StepStarted(started) = event {
167        if started.role_from.is_some() && started.sealed_role.is_some() {
168            return Err(unreadable(
169                event_type,
170                version,
171                "a step start cannot carry both dynamic and static role seals",
172            ));
173        }
174        if started
175            .sealed_role
176            .as_ref()
177            .is_some_and(|sealed| sealed != &started.started_by)
178        {
179            return Err(unreadable(
180                event_type,
181                version,
182                "the sealed static role differs from started_by",
183            ));
184        }
185    }
186    if let CeremonyEvent::ExecutionReceiptLinked(linked) = event {
187        linked
188            .link
189            .validate()
190            .map_err(|_| unreadable(event_type, version, "execution receipt link is invalid"))?;
191    }
192    Ok(())
193}
194
195#[cfg(test)]
196mod tests {
197    use super::*;
198    use crate::entities::ceremony_events::CeremonyCompleted;
199    use crate::value_objects::StateId;
200    use serde_json::json;
201    use time::macros::datetime;
202
203    fn completed_json() -> serde_json::Value {
204        json!({
205            "type": "ceremony_completed",
206            "final_state": "DONE",
207            "completed_at": "2026-07-29T09:00:00Z",
208        })
209    }
210
211    fn step_started_json() -> serde_json::Value {
212        json!({
213            "type": "step_started",
214            "step_id": "draft",
215            "iteration": 1,
216            "attempt": 1,
217            "lease": {
218                "owner_id": "host-1",
219                "idempotency_key": "ceremony-1:draft:1",
220                "acquired_at": "2026-07-29T09:00:00Z",
221                "expires_at": "2026-07-29T09:01:00Z"
222            },
223            "started_by": "writer",
224            "started_at": "2026-07-29T09:00:00Z"
225        })
226    }
227
228    #[test]
229    fn version_one_reads_directly() {
230        let event = CeremonyEventReader::read(
231            AuditEventType::CeremonyCompleted,
232            EventSchemaVersion::V1,
233            completed_json(),
234        )
235        .unwrap();
236
237        assert_eq!(
238            event,
239            CeremonyEvent::CeremonyCompleted(CeremonyCompleted {
240                final_state: StateId::new("DONE").unwrap(),
241                completed_at: datetime!(2026-07-29 09:00:00 UTC),
242            })
243        );
244    }
245
246    #[test]
247    fn completed_event_rejects_injected_failure_classification_at_every_readable_version() {
248        let legacy: serde_json::Value = serde_json::from_str(include_str!(
249            "../../tests/fixtures/ceremony_events/v1/step_completed.json"
250        ))
251        .unwrap();
252        for version in [
253            EventSchemaVersion::V1,
254            EventSchemaVersion::V2,
255            EventSchemaVersion::V3,
256        ] {
257            let mut raw = legacy.clone();
258            if version != EventSchemaVersion::V1 {
259                raw["state_iteration"] = json!(1);
260            }
261            if version == EventSchemaVersion::V3 {
262                raw["state_visit"] = json!(1);
263            }
264            CeremonyEventReader::read(AuditEventType::StepCompleted, version, raw.clone())
265                .expect("unclassified completion remains readable");
266            raw["result"]["failure_kind"] = json!("no_valid_proposal");
267            assert!(matches!(
268                CeremonyEventReader::read(AuditEventType::StepCompleted, version, raw),
269                Err(DomainError::UnreadableCeremonyEvent {
270                    reason: "step completion cannot carry a failure kind",
271                    ..
272                })
273            ));
274        }
275    }
276
277    #[test]
278    fn an_unknown_schema_version_is_refused_by_name() {
279        let error = CeremonyEventReader::read(
280            AuditEventType::CeremonyCompleted,
281            EventSchemaVersion::new(2).unwrap(),
282            completed_json(),
283        )
284        .unwrap_err();
285
286        assert!(matches!(
287            error,
288            DomainError::UnreadableCeremonyEvent {
289                event_type: "ceremony_completed",
290                version: 2,
291                ..
292            }
293        ));
294        assert!(error.to_string().contains("ceremony_completed"));
295        assert!(error.to_string().contains("version 2"));
296    }
297
298    #[test]
299    fn a_tag_that_disagrees_with_the_record_is_refused() {
300        let error = CeremonyEventReader::read(
301            AuditEventType::StepCompleted,
302            EventSchemaVersion::V1,
303            completed_json(),
304        )
305        .unwrap_err();
306
307        assert!(matches!(
308            error,
309            DomainError::UnreadableCeremonyEvent {
310                event_type: "step_completed",
311                version: 1,
312                ..
313            }
314        ));
315    }
316
317    #[test]
318    fn a_payload_that_is_not_an_event_is_refused() {
319        let error = CeremonyEventReader::read(
320            AuditEventType::CeremonyCompleted,
321            EventSchemaVersion::V1,
322            json!({ "type": "ceremony_completed" }),
323        )
324        .unwrap_err();
325
326        assert!(matches!(
327            error,
328            DomainError::UnreadableCeremonyEvent {
329                event_type: "ceremony_completed",
330                ..
331            }
332        ));
333    }
334
335    #[test]
336    fn version_one_refuses_a_version_two_coordinate() {
337        let mut raw = step_started_json();
338        raw["state_iteration"] = json!(1);
339
340        let error =
341            CeremonyEventReader::read(AuditEventType::StepStarted, EventSchemaVersion::V1, raw)
342                .unwrap_err();
343
344        assert!(matches!(
345            error,
346            DomainError::UnreadableCeremonyEvent {
347                event_type: "step_started",
348                version: 1,
349                ..
350            }
351        ));
352    }
353
354    #[test]
355    fn version_two_requires_an_explicit_coordinate() {
356        let error = CeremonyEventReader::read(
357            AuditEventType::StepStarted,
358            EventSchemaVersion::V2,
359            step_started_json(),
360        )
361        .unwrap_err();
362
363        assert!(matches!(
364            error,
365            DomainError::UnreadableCeremonyEvent {
366                event_type: "step_started",
367                version: 2,
368                ..
369            }
370        ));
371    }
372
373    #[test]
374    fn version_three_refuses_contradictory_static_role_seals() {
375        let mut raw = step_started_json();
376        raw["state_iteration"] = json!(1);
377        raw["sealed_role"] = json!("reviewer");
378
379        assert!(CeremonyEventReader::read(
380            AuditEventType::StepStarted,
381            EventSchemaVersion::V3,
382            raw,
383        )
384        .is_err());
385    }
386
387    #[test]
388    fn version_three_refuses_two_role_seal_markers() {
389        let mut raw = step_started_json();
390        raw["state_iteration"] = json!(1);
391        raw["role_from"] = json!("next_role");
392        raw["sealed_role"] = json!("writer");
393
394        assert!(CeremonyEventReader::read(
395            AuditEventType::StepStarted,
396            EventSchemaVersion::V3,
397            raw,
398        )
399        .is_err());
400    }
401
402    #[test]
403    fn version_four_reads_a_budgeted_ceremony_start() {
404        let raw = json!({
405            "type": "ceremony_instance_started",
406            "ceremony_id": "budgeted-reader",
407            "definition_name": "reader_fixture",
408            "definition_version": "1.0",
409            "initial_state": "OPEN",
410            "step_ids": ["work"],
411            "context": {},
412            "bound_definition": null,
413            "budget_account_id": "budgeted-reader",
414            "created_at": "2026-07-29T09:00:00Z"
415        });
416
417        let event = CeremonyEventReader::read(
418            AuditEventType::CeremonyInstanceStarted,
419            EventSchemaVersion::V4,
420            raw,
421        )
422        .unwrap();
423
424        assert_eq!(event.schema_version(), EventSchemaVersion::V4);
425    }
426
427    #[test]
428    fn version_six_reads_a_budgeted_step_claim() {
429        let mut raw = step_started_json();
430        raw["state_iteration"] = json!(1);
431        raw["state_visit"] = json!(1);
432        raw["budget_reservation_id"] = json!("reservation-reader-fixture");
433
434        let event =
435            CeremonyEventReader::read(AuditEventType::StepStarted, EventSchemaVersion::V6, raw)
436                .unwrap();
437
438        assert_eq!(event.schema_version(), EventSchemaVersion::V6);
439    }
440}