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