1use 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
19pub struct OperationLifecycle<B: JournalBackend> {
21 journal: Journal<Arc<B>>,
22 lease: Option<Lease>,
23}
24
25impl<B: JournalBackend> OperationLifecycle<B> {
26 pub fn new(backend: B) -> Self {
28 Self::from_shared(Arc::new(backend))
29 }
30 pub fn from_shared(backend: Arc<B>) -> Self {
32 Self {
33 journal: Journal::new(backend),
34 lease: None,
35 }
36 }
37 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 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}