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::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}