Skip to main content

sim_lib_operation_gate/
lifecycle_engine.rs

1//! Journal-backed operation execution and recovery coordinator.
2
3use std::sync::Arc;
4
5use sim_kernel::{Datum, Symbol};
6use sim_lib_journal::{Journal, JournalBackend, JournalEntry, JournalObject, Lease};
7
8use crate::{
9    OperationAttempt, OperationError, OperationGrant, OperationId, OperationIntent, ReplayPolicy,
10    lifecycle::{
11        FencedDispatch, LeaseWindow, LifecyclePerformer, LifecyclePerformerResponse,
12        LifecycleReceipt, OperationLease, OperationObservation, OperationOutcome, OperationStep,
13        PostconditionObserver, PostconditionRequest, PostconditionResponse,
14    },
15    lifecycle_project::project,
16    lifecycle_record::{IdentifiedOutcome, OperationLifecycleRecord},
17};
18
19/// Journal-backed M5 coordinator for fenced dispatch and independent reconciliation.
20pub struct OperationLifecycle<B: JournalBackend> {
21    journal: Journal<Arc<B>>,
22    lease: Option<Lease>,
23}
24
25impl<B: JournalBackend> OperationLifecycle<B> {
26    /// Creates a lifecycle over one backend value.
27    pub fn new(backend: B) -> Self {
28        Self::from_shared(Arc::new(backend))
29    }
30    /// Creates a lifecycle over a shared backend for crash/reopen recovery.
31    pub fn from_shared(backend: Arc<B>) -> Self {
32        Self {
33            journal: Journal::new(backend),
34            lease: None,
35        }
36    }
37    /// Reconstructs one complete lifecycle solely from the verified journal.
38    pub fn record(
39        &self,
40        operation: &OperationId,
41    ) -> Result<Option<OperationLifecycleRecord>, OperationError> {
42        Ok(project(self.journal.verified_snapshot()?)?.remove(operation))
43    }
44
45    /// Advances one operation, observing before any new or repeated performance.
46    pub fn run(
47        &mut self,
48        intent: &OperationIntent,
49        grant: &OperationGrant,
50        window: LeaseWindow,
51        performer: &mut dyn LifecyclePerformer,
52        observer: &mut dyn PostconditionObserver,
53    ) -> Result<OperationOutcome, OperationError> {
54        intent.verify()?;
55        grant.verify()?;
56        if grant.operation != intent.id {
57            return Err(OperationError::GrantMismatch);
58        }
59        let performer_identity = performer.identity();
60        let observer_identity = observer.identity();
61        if performer_identity == observer_identity {
62            return Err(OperationError::ObserverNotIndependent);
63        }
64        self.lease = Some(self.journal.acquire_lease()?);
65
66        let mut records = project(self.journal.verified_snapshot()?)?;
67        if let Some(existing) = records.get(intent.id()) {
68            if existing.intent != *intent {
69                return Err(OperationError::ContradictoryIntent);
70            }
71            if existing.grant != *grant {
72                return Err(OperationError::GrantMismatch);
73            }
74            if existing
75                .dispatches
76                .iter()
77                .any(|dispatch| dispatch.performer() != &performer_identity)
78            {
79                return Err(OperationError::PerformerMismatch);
80            }
81            if let Some(
82                outcome
83                @ (OperationOutcome::AlreadyTrue { .. } | OperationOutcome::Verified { .. }),
84            ) = existing.outcome()
85            {
86                return Ok(outcome.clone());
87            }
88            if existing
89                .leases
90                .last()
91                .is_some_and(|lease| window.acquired_at() < lease.acquired_at())
92            {
93                return Err(OperationError::InvalidLease);
94            }
95            if intent.replay_policy() == ReplayPolicy::ExactlyOnce
96                && matches!(existing.outcome(), Some(OperationOutcome::Diverged { .. }))
97            {
98                return Ok(existing.outcome().expect("matched outcome").clone());
99            }
100        } else {
101            self.append(
102                "lifecycle-intent-persisted",
103                vec![intent.canonical_datum(), grant.canonical_datum()],
104            )?;
105            records = project(self.journal.verified_snapshot()?)?;
106        }
107
108        let record = records
109            .get(intent.id())
110            .expect("intent was persisted")
111            .clone();
112        let observation =
113            self.observe(&record, &performer_identity, window.acquired_at(), observer)?;
114        match observation.response() {
115            PostconditionResponse::Satisfied { .. } => {
116                let outcome = if record.dispatches.is_empty() {
117                    OperationOutcome::AlreadyTrue {
118                        evidence: observation.evidence.clone(),
119                    }
120                } else {
121                    OperationOutcome::Verified {
122                        evidence: observation.evidence.clone(),
123                    }
124                };
125                return self.persist_outcome(intent.id(), outcome);
126            }
127            PostconditionResponse::Unavailable { .. } | PostconditionResponse::Disputed { .. } => {
128                return self.persist_outcome(
129                    intent.id(),
130                    OperationOutcome::Uncertain {
131                        last_durable_step: OperationStep::ObservationPersisted,
132                    },
133                );
134            }
135            PostconditionResponse::NotSatisfied { observed, .. }
136                if !record.dispatches.is_empty() =>
137            {
138                if intent.replay_policy() == ReplayPolicy::ExactlyOnce {
139                    return self.persist_outcome(
140                        intent.id(),
141                        OperationOutcome::Diverged {
142                            observed: observed.clone(),
143                            expected: intent.intended_result().clone(),
144                        },
145                    );
146                }
147                if record
148                    .leases
149                    .last()
150                    .is_some_and(|lease| window.acquired_at < lease.expires_at())
151                {
152                    return self.persist_outcome(
153                        intent.id(),
154                        OperationOutcome::Uncertain {
155                            last_durable_step: OperationStep::ObservationPersisted,
156                        },
157                    );
158                }
159            }
160            PostconditionResponse::NotSatisfied { .. } => {}
161        }
162
163        let writer_fence = self.lease.as_ref().expect("writer lease acquired").fence();
164        let operation_lease = OperationLease::new(
165            intent.id().clone(),
166            window.holder().clone(),
167            writer_fence,
168            window.acquired_at(),
169            window.expires_at(),
170        )?;
171        self.append(
172            "lifecycle-lease-acquired",
173            vec![operation_lease.canonical_datum()],
174        )?;
175        let ordinal =
176            u64::try_from(record.attempts.len()).map_err(|_| OperationError::SequenceExhausted)?;
177        let attempt = OperationAttempt::new(intent.id().clone(), ordinal)?;
178        let dispatch = FencedDispatch::new(
179            intent.id().clone(),
180            grant.id().clone(),
181            attempt.id().clone(),
182            operation_lease.id().clone(),
183            performer_identity.clone(),
184        )?;
185        self.append(
186            "lifecycle-dispatch-persisted",
187            vec![dispatch.canonical_datum(), attempt.canonical_datum()],
188        )?;
189
190        if let LifecyclePerformerResponse::Receipt(raw) = performer.perform(&dispatch) {
191            let receipt = LifecycleReceipt::new(dispatch.id().clone(), raw)?;
192            self.append(
193                "lifecycle-receipt-persisted",
194                vec![receipt.canonical_datum()],
195            )?;
196        }
197        let updated = self
198            .record(intent.id())?
199            .expect("lifecycle remains present");
200        let observation = self.observe(
201            &updated,
202            &performer_identity,
203            window.acquired_at(),
204            observer,
205        )?;
206        let outcome = match observation.response() {
207            PostconditionResponse::Satisfied { .. } => OperationOutcome::Verified {
208                evidence: observation.evidence.clone(),
209            },
210            PostconditionResponse::NotSatisfied { observed, .. } => OperationOutcome::Diverged {
211                observed: observed.clone(),
212                expected: intent.intended_result().clone(),
213            },
214            PostconditionResponse::Unavailable { .. } | PostconditionResponse::Disputed { .. } => {
215                OperationOutcome::Uncertain {
216                    last_durable_step: OperationStep::ObservationPersisted,
217                }
218            }
219        };
220        self.persist_outcome(intent.id(), outcome)
221    }
222
223    fn observe(
224        &mut self,
225        record: &OperationLifecycleRecord,
226        performer: &Datum,
227        observed_at: u64,
228        observer: &mut dyn PostconditionObserver,
229    ) -> Result<OperationObservation, OperationError> {
230        let request = PostconditionRequest {
231            operation: record.intent.id().clone(),
232            target: record.intent.target().clone(),
233            expected: record.intent.intended_result().clone(),
234            dispatch: record.dispatches.last().map(|value| value.id.clone()),
235            receipt: record
236                .dispatches
237                .last()
238                .and_then(|dispatch| {
239                    record
240                        .receipts
241                        .iter()
242                        .rev()
243                        .find(|receipt| receipt.dispatch == dispatch.id)
244                })
245                .map(|value| value.id.clone()),
246            last_durable_step: record.last_step(),
247            observed_at,
248        };
249        let observer_identity = observer.identity();
250        if &observer_identity == performer {
251            return Err(OperationError::ObserverNotIndependent);
252        }
253        let observation =
254            OperationObservation::new(&request, observer_identity, observer.observe(&request))?;
255        self.append(
256            "lifecycle-observation-persisted",
257            vec![observation.stored_datum()],
258        )?;
259        Ok(observation)
260    }
261
262    fn persist_outcome(
263        &mut self,
264        operation: &OperationId,
265        outcome: OperationOutcome,
266    ) -> Result<OperationOutcome, OperationError> {
267        let identified = IdentifiedOutcome::new(operation.clone(), outcome.clone())?;
268        self.append(
269            "lifecycle-outcome-persisted",
270            vec![identified.canonical_datum()],
271        )?;
272        Ok(outcome)
273    }
274
275    fn append(&self, kind: &'static str, datums: Vec<Datum>) -> Result<(), OperationError> {
276        let lease = self.lease.as_ref().ok_or(OperationError::NotResumed)?;
277        let objects = datums
278            .into_iter()
279            .map(JournalObject::from_datum)
280            .collect::<Result<Vec<_>, _>>()?;
281        let payloads = objects.iter().map(|object| object.id.clone()).collect();
282        let expected = self.journal.head()?;
283        let sequence = expected
284            .as_ref()
285            .map_or(Some(0), |head| head.sequence.checked_add(1))
286            .ok_or(OperationError::SequenceExhausted)?;
287        let entry = JournalEntry::new(
288            sequence,
289            expected.as_ref().map(|head| head.entry.clone()),
290            Symbol::qualified("operation", kind),
291            payloads,
292        );
293        self.journal
294            .publish(lease, expected.as_ref(), objects, vec![entry])?;
295        Ok(())
296    }
297}