1use crate::error::DomainError;
5use crate::value_objects::{AuditEventType, EventSchemaVersion};
6
7use super::CeremonyEvent;
8
9#[derive(Debug)]
20pub struct CeremonyEventReader;
21
22impl CeremonyEventReader {
23 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}