Skip to main content

made_core/entities/
ceremony_intervention.rs

1//! Dynamic intervention owned by a running ceremony instance.
2
3use 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    // Everything below is additive and absent by default, so an item
33    // sealed before interventions could be routed re-serializes with
34    // exactly the bytes it was sealed with.
35    #[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    /// The same item, said to be for one of the four reasons.
89    #[must_use]
90    pub fn with_intent(mut self, intent: CeremonyInterventionIntent) -> Self {
91        self.intent = Some(intent);
92        self
93    }
94
95    /// The same item, offered to its host on the given terms.
96    #[must_use]
97    pub fn with_delivery(mut self, delivery: InterventionDeliveryPolicy) -> Self {
98        self.delivery = Some(delivery);
99        self
100    }
101
102    /// The same item, asked by somebody who holds no seat at the table.
103    #[must_use]
104    pub fn with_supervisor(mut self, supervisor: SupervisorPrincipal) -> Self {
105        self.supervisor = Some(supervisor);
106        self
107    }
108
109    /// Record that a named agent said it saw this item.
110    ///
111    /// Appending is the fold's job; the rule about what agrees with
112    /// what lives on the acknowledgement, and repeating an identical
113    /// one changes nothing rather than growing the list.
114    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    /// What this item's host said about one offer of it, if anything.
132    #[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    /// Whether a recipient has already said it saw this item.
143    #[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    /// Append an answer exactly as it was sealed.
160    ///
161    /// The door a replay comes through. A fold that rebuilt the answer
162    /// from its parts would silently drop whatever the value object
163    /// grew since — which agent gave it, which offer it closes — and
164    /// the replayed session would disagree with its own journal about
165    /// who answered.
166    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}