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::StepStarted
40 | AuditEventType::StepCompleted
41 | AuditEventType::StepFailed
42 | AuditEventType::TransitionApplied
43 | AuditEventType::ContextWritten
44 | AuditEventType::StateIterationStarted
45 ),
46 EventSchemaVersion::V3 => matches!(
47 event_type,
48 AuditEventType::StepStarted
49 | AuditEventType::StepCompleted
50 | AuditEventType::StepFailed
51 | AuditEventType::TransitionApplied
52 ),
53 EventSchemaVersion::V4 => event_type == AuditEventType::StepStarted,
54 _ => false,
55 };
56 if !supported {
57 return Err(unreadable(
58 event_type,
59 version,
60 "no reader exists for this schema version",
61 ));
62 }
63 let event = serde_json::from_value::<CeremonyEvent>(raw).map_err(|_| {
64 unreadable(
65 event_type,
66 version,
67 "the payload does not deserialize as a ceremony event",
68 )
69 })?;
70 if event.event_type() != event_type {
71 return Err(unreadable(
72 event_type,
73 version,
74 "the payload's tag names a different event type",
75 ));
76 }
77 if let CeremonyEvent::StepStarted(started) = &event {
78 if started.role_from.is_some() && started.sealed_role.is_some() {
79 return Err(unreadable(
80 event_type,
81 version,
82 "a step start cannot carry both dynamic and static role seals",
83 ));
84 }
85 if started
86 .sealed_role
87 .as_ref()
88 .is_some_and(|sealed| sealed != &started.started_by)
89 {
90 return Err(unreadable(
91 event_type,
92 version,
93 "the sealed static role differs from started_by",
94 ));
95 }
96 }
97 let coordinates_valid = match &event {
98 CeremonyEvent::StepStarted(e) => e.state_visit.is_none() || e.state_iteration.is_some(),
99 CeremonyEvent::StepCompleted(e) => {
100 e.state_visit.is_none() || e.state_iteration.is_some()
101 }
102 CeremonyEvent::StepFailed(e) => e.state_visit.is_none() || e.state_iteration.is_some(),
103 CeremonyEvent::TransitionApplied(e) => match &e.destination {
104 Some(destination) => {
105 e.transition.has_explicit_state_visit()
106 && e.transition.has_explicit_state_iteration()
107 && e.transition.state_visit().next().ok() == Some(destination.state_visit)
108 && destination
109 .step_ids
110 .iter()
111 .collect::<std::collections::BTreeSet<_>>()
112 .len()
113 == destination.step_ids.len()
114 }
115 None => !e.transition.has_explicit_state_visit(),
116 },
117 _ => true,
118 };
119 if !coordinates_valid {
120 return Err(unreadable(
121 event_type,
122 version,
123 "inconsistent state visit coordinates or destination reset",
124 ));
125 }
126 if event.schema_version() != version {
127 return Err(unreadable(
128 event_type,
129 version,
130 "the payload shape does not match its schema version",
131 ));
132 }
133 Ok(event)
134 }
135}
136
137fn unreadable(
138 event_type: AuditEventType,
139 version: EventSchemaVersion,
140 reason: &'static str,
141) -> DomainError {
142 DomainError::UnreadableCeremonyEvent {
143 event_type: event_type.as_str(),
144 version: version.get(),
145 reason,
146 }
147}
148
149#[cfg(test)]
150mod tests {
151 use super::*;
152 use crate::entities::ceremony_events::CeremonyCompleted;
153 use crate::value_objects::StateId;
154 use serde_json::json;
155 use time::macros::datetime;
156
157 fn completed_json() -> serde_json::Value {
158 json!({
159 "type": "ceremony_completed",
160 "final_state": "DONE",
161 "completed_at": "2026-07-29T09:00:00Z",
162 })
163 }
164
165 fn step_started_json() -> serde_json::Value {
166 json!({
167 "type": "step_started",
168 "step_id": "draft",
169 "iteration": 1,
170 "attempt": 1,
171 "lease": {
172 "owner_id": "host-1",
173 "idempotency_key": "ceremony-1:draft:1",
174 "acquired_at": "2026-07-29T09:00:00Z",
175 "expires_at": "2026-07-29T09:01:00Z"
176 },
177 "started_by": "writer",
178 "started_at": "2026-07-29T09:00:00Z"
179 })
180 }
181
182 #[test]
183 fn version_one_reads_directly() {
184 let event = CeremonyEventReader::read(
185 AuditEventType::CeremonyCompleted,
186 EventSchemaVersion::V1,
187 completed_json(),
188 )
189 .unwrap();
190
191 assert_eq!(
192 event,
193 CeremonyEvent::CeremonyCompleted(CeremonyCompleted {
194 final_state: StateId::new("DONE").unwrap(),
195 completed_at: datetime!(2026-07-29 09:00:00 UTC),
196 })
197 );
198 }
199
200 #[test]
201 fn an_unknown_schema_version_is_refused_by_name() {
202 let error = CeremonyEventReader::read(
203 AuditEventType::CeremonyCompleted,
204 EventSchemaVersion::new(2).unwrap(),
205 completed_json(),
206 )
207 .unwrap_err();
208
209 assert!(matches!(
210 error,
211 DomainError::UnreadableCeremonyEvent {
212 event_type: "ceremony_completed",
213 version: 2,
214 ..
215 }
216 ));
217 assert!(error.to_string().contains("ceremony_completed"));
218 assert!(error.to_string().contains("version 2"));
219 }
220
221 #[test]
222 fn a_tag_that_disagrees_with_the_record_is_refused() {
223 let error = CeremonyEventReader::read(
224 AuditEventType::StepCompleted,
225 EventSchemaVersion::V1,
226 completed_json(),
227 )
228 .unwrap_err();
229
230 assert!(matches!(
231 error,
232 DomainError::UnreadableCeremonyEvent {
233 event_type: "step_completed",
234 version: 1,
235 ..
236 }
237 ));
238 }
239
240 #[test]
241 fn a_payload_that_is_not_an_event_is_refused() {
242 let error = CeremonyEventReader::read(
243 AuditEventType::CeremonyCompleted,
244 EventSchemaVersion::V1,
245 json!({ "type": "ceremony_completed" }),
246 )
247 .unwrap_err();
248
249 assert!(matches!(
250 error,
251 DomainError::UnreadableCeremonyEvent {
252 event_type: "ceremony_completed",
253 ..
254 }
255 ));
256 }
257
258 #[test]
259 fn version_one_refuses_a_version_two_coordinate() {
260 let mut raw = step_started_json();
261 raw["state_iteration"] = json!(1);
262
263 let error =
264 CeremonyEventReader::read(AuditEventType::StepStarted, EventSchemaVersion::V1, raw)
265 .unwrap_err();
266
267 assert!(matches!(
268 error,
269 DomainError::UnreadableCeremonyEvent {
270 event_type: "step_started",
271 version: 1,
272 ..
273 }
274 ));
275 }
276
277 #[test]
278 fn version_two_requires_an_explicit_coordinate() {
279 let error = CeremonyEventReader::read(
280 AuditEventType::StepStarted,
281 EventSchemaVersion::V2,
282 step_started_json(),
283 )
284 .unwrap_err();
285
286 assert!(matches!(
287 error,
288 DomainError::UnreadableCeremonyEvent {
289 event_type: "step_started",
290 version: 2,
291 ..
292 }
293 ));
294 }
295
296 #[test]
297 fn version_three_refuses_contradictory_static_role_seals() {
298 let mut raw = step_started_json();
299 raw["state_iteration"] = json!(1);
300 raw["sealed_role"] = json!("reviewer");
301
302 assert!(CeremonyEventReader::read(
303 AuditEventType::StepStarted,
304 EventSchemaVersion::V3,
305 raw,
306 )
307 .is_err());
308 }
309
310 #[test]
311 fn version_three_refuses_two_role_seal_markers() {
312 let mut raw = step_started_json();
313 raw["state_iteration"] = json!(1);
314 raw["role_from"] = json!("next_role");
315 raw["sealed_role"] = json!("writer");
316
317 assert!(CeremonyEventReader::read(
318 AuditEventType::StepStarted,
319 EventSchemaVersion::V3,
320 raw,
321 )
322 .is_err());
323 }
324}