1use serde::{Deserialize, Serialize};
4use time::OffsetDateTime;
5
6use crate::entities::CeremonyEvidencePack;
7use crate::error::DomainError;
8use crate::value_objects::{
9 CeremonyInterventionContent, CeremonyInterventionId, CeremonyInterventionIntent,
10 CeremonyInterventionKind, CeremonyInterventionProvenance, CeremonyInterventionResponse,
11 CeremonyInterventionStatus, CeremonyInterventionTarget, DeliveryRecipient,
12 InterventionDeliveryAck, InterventionDeliveryPolicy, RoleId, SupervisorPrincipal,
13};
14
15#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
16pub struct CeremonyIntervention {
17 id: CeremonyInterventionId,
18 kind: CeremonyInterventionKind,
19 requested_by: RoleId,
20 target: CeremonyInterventionTarget,
21 request: CeremonyInterventionContent,
22 #[serde(default)]
23 provenance: Option<CeremonyInterventionProvenance>,
24 responses: Vec<CeremonyInterventionResponse>,
25 status: CeremonyInterventionStatus,
26 #[serde(with = "time::serde::rfc3339")]
27 created_at: OffsetDateTime,
28 #[serde(with = "time::serde::rfc3339")]
29 updated_at: OffsetDateTime,
30 #[serde(with = "time::serde::rfc3339::option")]
31 closed_at: Option<OffsetDateTime>,
32 #[serde(default, skip_serializing_if = "Option::is_none")]
36 intent: Option<CeremonyInterventionIntent>,
37 #[serde(default, skip_serializing_if = "Option::is_none")]
38 delivery: Option<InterventionDeliveryPolicy>,
39 #[serde(default, skip_serializing_if = "Option::is_none")]
40 supervisor: Option<SupervisorPrincipal>,
41 #[serde(default, skip_serializing_if = "Vec::is_empty")]
42 deliveries: Vec<InterventionDeliveryAck>,
43}
44
45impl CeremonyIntervention {
46 #[must_use]
47 pub fn open(
48 id: CeremonyInterventionId,
49 kind: CeremonyInterventionKind,
50 requested_by: RoleId,
51 target: CeremonyInterventionTarget,
52 request: CeremonyInterventionContent,
53 now: OffsetDateTime,
54 ) -> Self {
55 Self::open_with_provenance(id, kind, requested_by, target, request, None, now)
56 }
57
58 #[allow(clippy::too_many_arguments)]
59 #[must_use]
60 pub fn open_with_provenance(
61 id: CeremonyInterventionId,
62 kind: CeremonyInterventionKind,
63 requested_by: RoleId,
64 target: CeremonyInterventionTarget,
65 request: CeremonyInterventionContent,
66 provenance: Option<CeremonyInterventionProvenance>,
67 now: OffsetDateTime,
68 ) -> Self {
69 Self {
70 id,
71 kind,
72 requested_by,
73 target,
74 request,
75 provenance,
76 responses: Vec::new(),
77 status: CeremonyInterventionStatus::Open,
78 created_at: now,
79 updated_at: now,
80 closed_at: None,
81 intent: None,
82 delivery: None,
83 supervisor: None,
84 deliveries: Vec::new(),
85 }
86 }
87
88 #[must_use]
90 pub fn with_intent(mut self, intent: CeremonyInterventionIntent) -> Self {
91 self.intent = Some(intent);
92 self
93 }
94
95 #[must_use]
97 pub fn with_delivery(mut self, delivery: InterventionDeliveryPolicy) -> Self {
98 self.delivery = Some(delivery);
99 self
100 }
101
102 #[must_use]
104 pub fn with_supervisor(mut self, supervisor: SupervisorPrincipal) -> Self {
105 self.supervisor = Some(supervisor);
106 self
107 }
108
109 pub fn acknowledge_delivery(
115 &mut self,
116 ack: InterventionDeliveryAck,
117 ) -> Result<(), DomainError> {
118 if let Some(existing) = self.delivery_ack(ack.delivery_id()) {
119 if existing.agrees_with(&ack) {
120 return Ok(());
121 }
122 return Err(DomainError::Conflict {
123 what: "ceremony_intervention.delivery_acknowledgement",
124 });
125 }
126 self.updated_at = ack.acknowledged_at();
127 self.deliveries.push(ack);
128 Ok(())
129 }
130
131 #[must_use]
133 pub fn delivery_ack(
134 &self,
135 delivery_id: &crate::value_objects::HostDeliveryId,
136 ) -> Option<&InterventionDeliveryAck> {
137 self.deliveries
138 .iter()
139 .find(|ack| ack.delivery_id() == delivery_id)
140 }
141
142 #[must_use]
144 pub fn was_acknowledged_by(&self, recipient: &DeliveryRecipient) -> bool {
145 self.deliveries
146 .iter()
147 .any(|ack| ack.recipient() == recipient)
148 }
149
150 pub fn respond(
151 &mut self,
152 role_id: RoleId,
153 content: CeremonyInterventionContent,
154 now: OffsetDateTime,
155 ) -> Result<(), DomainError> {
156 self.accept_response(CeremonyInterventionResponse::new(role_id, content, now))
157 }
158
159 pub fn accept_response(
167 &mut self,
168 response: CeremonyInterventionResponse,
169 ) -> Result<(), DomainError> {
170 self.ensure_can_respond(response.role_id())?;
171 self.updated_at = response.responded_at();
172 self.responses.push(response);
173 Ok(())
174 }
175
176 pub fn respond_with_evidence(
177 &mut self,
178 role_id: RoleId,
179 evidence_pack: CeremonyEvidencePack,
180 now: OffsetDateTime,
181 ) -> Result<(), DomainError> {
182 self.accept_response(CeremonyInterventionResponse::from_evidence(
183 role_id,
184 evidence_pack,
185 now,
186 )?)
187 }
188
189 pub(crate) fn ensure_can_respond(&self, role_id: &RoleId) -> Result<(), DomainError> {
190 if !self.status.is_open() {
191 return Err(DomainError::InvariantViolated {
192 reason: "closed ceremony interventions cannot receive responses",
193 });
194 }
195 if !self.target.accepts(role_id) {
196 return Err(DomainError::InvariantViolated {
197 reason: "ceremony intervention does not target responding role",
198 });
199 }
200 if self
201 .responses
202 .iter()
203 .any(|response| response.role_id() == role_id)
204 {
205 return Err(DomainError::AlreadyExists {
206 what: "ceremony_intervention.response_role",
207 });
208 }
209 Ok(())
210 }
211
212 pub fn close(&mut self, role_id: &RoleId, now: OffsetDateTime) -> Result<(), DomainError> {
213 if !self.status.is_open() {
214 return Err(DomainError::InvariantViolated {
215 reason: "ceremony intervention is already closed",
216 });
217 }
218 if role_id != &self.requested_by {
219 return Err(DomainError::InvariantViolated {
220 reason: "only the requesting role can close a ceremony intervention",
221 });
222 }
223 self.status = CeremonyInterventionStatus::Closed;
224 self.updated_at = now;
225 self.closed_at = Some(now);
226 Ok(())
227 }
228
229 #[must_use]
230 pub fn id(&self) -> &CeremonyInterventionId {
231 &self.id
232 }
233
234 #[must_use]
235 pub const fn kind(&self) -> CeremonyInterventionKind {
236 self.kind
237 }
238
239 #[must_use]
240 pub fn requested_by(&self) -> &RoleId {
241 &self.requested_by
242 }
243
244 #[must_use]
245 pub fn target(&self) -> &CeremonyInterventionTarget {
246 &self.target
247 }
248
249 #[must_use]
250 pub fn request(&self) -> &CeremonyInterventionContent {
251 &self.request
252 }
253
254 #[must_use]
255 pub fn provenance(&self) -> Option<&CeremonyInterventionProvenance> {
256 self.provenance.as_ref()
257 }
258
259 #[must_use]
260 pub fn responses(&self) -> &[CeremonyInterventionResponse] {
261 &self.responses
262 }
263
264 #[must_use]
265 pub const fn status(&self) -> CeremonyInterventionStatus {
266 self.status
267 }
268
269 #[must_use]
270 pub fn created_at(&self) -> OffsetDateTime {
271 self.created_at
272 }
273
274 #[must_use]
275 pub fn updated_at(&self) -> OffsetDateTime {
276 self.updated_at
277 }
278
279 #[must_use]
280 pub fn closed_at(&self) -> Option<OffsetDateTime> {
281 self.closed_at
282 }
283
284 #[must_use]
285 pub const fn intent(&self) -> Option<CeremonyInterventionIntent> {
286 self.intent
287 }
288
289 #[must_use]
290 pub const fn delivery(&self) -> Option<&InterventionDeliveryPolicy> {
291 self.delivery.as_ref()
292 }
293
294 #[must_use]
295 pub const fn supervisor(&self) -> Option<&SupervisorPrincipal> {
296 self.supervisor.as_ref()
297 }
298
299 #[must_use]
300 pub fn deliveries(&self) -> &[InterventionDeliveryAck] {
301 &self.deliveries
302 }
303}
304
305#[cfg(test)]
306mod tests {
307 use time::macros::datetime;
308
309 use crate::value_objects::Attributes;
310
311 use super::*;
312
313 fn content(message: &str) -> CeremonyInterventionContent {
314 CeremonyInterventionContent::new(message, Attributes::empty()).unwrap()
315 }
316
317 #[test]
318 fn accepts_one_response_per_target_role_and_requester_controls_close() {
319 let now = datetime!(2026-07-20 12:00:00 UTC);
320 let engineer = RoleId::new("ENGINEER").unwrap();
321 let observer = RoleId::new("OBSERVER").unwrap();
322 let mut intervention = CeremonyIntervention::open(
323 CeremonyInterventionId::new("intervention-1").unwrap(),
324 CeremonyInterventionKind::Investigation,
325 engineer.clone(),
326 CeremonyInterventionTarget::roles([observer.clone()]).unwrap(),
327 content("Inspect the queue without consuming messages."),
328 now,
329 );
330
331 intervention
332 .respond(observer.clone(), content("Depth is stable."), now)
333 .unwrap();
334
335 assert!(intervention
336 .respond(observer.clone(), content("Duplicate."), now)
337 .is_err());
338 assert!(intervention.close(&observer, now).is_err());
339 intervention.close(&engineer, now).unwrap();
340 assert_eq!(intervention.status(), CeremonyInterventionStatus::Closed);
341 }
342}