1use std::time::{Duration, SystemTime};
2
3use serde::{Deserialize, Serialize};
4use serde_json::Value;
5
6use crate::id::{EffectId, EffectKey, WorkerId};
7use crate::kind::EffectKind;
8use crate::state::{EffectStatus, InvalidTransition, Transition};
9use crate::{EffectEvent, ErrorRecord, Lease, NewEffect, StoreError, TransitionRequest};
10
11#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
17pub struct EffectRecord {
18 pub id: EffectId,
20 pub key: EffectKey,
22 pub kind: EffectKind,
24 pub status: EffectStatus,
26 pub input: Option<Value>,
28 pub input_fingerprint: Option<String>,
30 pub output: Option<Value>,
32 pub last_error: Option<ErrorRecord>,
34 pub created_by: Option<String>,
36 pub attempt_count: u32,
38 pub may_have_applied: bool,
43 pub compensation_attempts: u32,
45 pub approved: bool,
48 pub next_attempt_at: Option<SystemTime>,
50 pub attempt_started_at: Option<SystemTime>,
52 pub attempt_ended_at: Option<SystemTime>,
56 pub lease_owner: Option<WorkerId>,
58 pub lease_epoch: u64,
60 pub lease_expires_at: Option<SystemTime>,
62 pub version: u64,
64 pub created_at: SystemTime,
66 pub updated_at: SystemTime,
68 pub committed_at: Option<SystemTime>,
70}
71
72impl EffectRecord {
73 pub fn new(new: NewEffect) -> Self {
75 Self {
76 id: new.id,
77 key: new.key,
78 kind: new.kind,
79 status: EffectStatus::Pending,
80 input: new.input,
81 input_fingerprint: new.input_fingerprint,
82 output: None,
83 last_error: None,
84 created_by: new.created_by,
85 attempt_count: 0,
86 may_have_applied: false,
87 compensation_attempts: 0,
88 approved: false,
89 next_attempt_at: None,
90 attempt_started_at: None,
91 attempt_ended_at: None,
92 lease_owner: None,
93 lease_epoch: 0,
94 lease_expires_at: None,
95 version: 0,
96 created_at: new.now,
97 updated_at: new.now,
98 committed_at: None,
99 }
100 }
101
102 pub fn live_lease_owner(&self, now: SystemTime) -> Option<&WorkerId> {
104 match (&self.lease_owner, self.lease_expires_at) {
105 (Some(owner), Some(expires_at)) if expires_at > now => Some(owner),
106 _ => None,
107 }
108 }
109
110 pub fn acquire_lease(
117 &mut self,
118 owner: &WorkerId,
119 now: SystemTime,
120 ttl: Duration,
121 ) -> Result<Lease, StoreError> {
122 if let Some(holder) = self.live_lease_owner(now) {
123 return Err(StoreError::LeaseHeld {
124 owner: holder.clone(),
125 expires_at: self.lease_expires_at.unwrap_or(now),
126 });
127 }
128 let expires_at = expiry(now, ttl);
129 self.lease_epoch += 1;
130 self.lease_owner = Some(owner.clone());
131 self.lease_expires_at = Some(expires_at);
132 Ok(Lease {
133 effect_id: self.id,
134 owner: owner.clone(),
135 epoch: self.lease_epoch,
136 expires_at,
137 })
138 }
139
140 pub fn renew_lease(
148 &mut self,
149 lease: &Lease,
150 now: SystemTime,
151 ttl: Duration,
152 ) -> Result<Lease, StoreError> {
153 self.check_lease(lease, now)?;
154 let expires_at = expiry(now, ttl);
155 self.lease_expires_at = Some(expires_at);
156 Ok(Lease {
157 expires_at,
158 ..lease.clone()
159 })
160 }
161
162 pub fn release_lease(&mut self, lease: &Lease) -> bool {
165 let current = lease.effect_id == self.id
166 && self.lease_epoch == lease.epoch
167 && self.lease_owner.as_ref() == Some(&lease.owner);
168 if current {
169 self.lease_owner = None;
170 self.lease_expires_at = None;
171 }
172 current
173 }
174
175 pub fn apply(&mut self, request: TransitionRequest) -> Result<EffectEvent, StoreError> {
208 let now = request.now;
209 match &request.lease {
210 Some(lease) => self.check_lease(lease, now)?,
211 None => {
212 if let Some(owner) = self.live_lease_owner(now) {
213 return Err(StoreError::LeaseHeld {
214 owner: owner.clone(),
215 expires_at: self.lease_expires_at.unwrap_or(now),
216 });
217 }
218 }
219 }
220 if request.expected_version != self.version {
221 return Err(StoreError::VersionConflict {
222 expected: request.expected_version,
223 actual: self.version,
224 });
225 }
226 let from = self.status;
227 let to = from.apply(request.transition)?;
228 if request.transition == Transition::FailedDefinitively
229 && self.may_have_applied
230 && self.kind != EffectKind::Read
231 {
232 return Err(InvalidTransition {
233 from,
234 event: request.transition,
235 }
236 .into());
237 }
238
239 match request.transition {
240 Transition::StartAttempt => {
241 self.attempt_count = self.attempt_count.saturating_add(1);
242 self.attempt_started_at = Some(now);
243 self.attempt_ended_at = None;
244 self.next_attempt_at = None;
245 }
246 Transition::ScheduleRetry
247 | Transition::ResolvedRetry
248 | Transition::ScheduleCompensationRetry => {
249 self.next_attempt_at = Some(request.next_attempt_at.unwrap_or(now));
250 }
251 Transition::Approve => self.approved = true,
252 Transition::StartCompensation => {
253 self.compensation_attempts = 1;
254 self.next_attempt_at = None;
255 }
256 Transition::StartCompensationRetry => {
257 self.compensation_attempts = self.compensation_attempts.saturating_add(1);
258 self.next_attempt_at = None;
259 }
260 _ => {}
261 }
262 if from == EffectStatus::Executing {
263 self.attempt_ended_at = Some(now);
264 }
265 if to == EffectStatus::Committed {
266 self.committed_at = Some(now);
267 }
268 if to == EffectStatus::Unknown {
269 self.may_have_applied = true;
270 }
271 let shown_not_applied = matches!(
272 request.transition,
273 Transition::VerificationNotApplied | Transition::ResolvedNotApplied
274 ) || (from == EffectStatus::Verifying
275 && request.transition == Transition::ScheduleRetry);
276 if shown_not_applied {
277 self.may_have_applied = false;
278 }
279 if let Some(output) = request.output {
280 self.output = Some(output);
281 }
282 if let Some(error) = request.error {
283 self.last_error = Some(error);
284 }
285 self.status = to;
286 self.version += 1;
287 self.updated_at = now;
288
289 Ok(EffectEvent {
290 effect_id: self.id,
291 sequence: self.version,
292 transition: request.transition,
293 from,
294 to,
295 attempt: self.attempt_count,
296 actor: request.actor,
297 payload: request.payload,
298 at: now,
299 })
300 }
301
302 fn check_lease(&self, lease: &Lease, now: SystemTime) -> Result<(), StoreError> {
303 let valid = lease.effect_id == self.id
304 && self.lease_epoch == lease.epoch
305 && self.lease_owner.as_ref() == Some(&lease.owner)
306 && self
307 .lease_expires_at
308 .is_some_and(|expires_at| expires_at > now);
309 if valid {
310 Ok(())
311 } else {
312 Err(StoreError::LeaseLost)
313 }
314 }
315}
316
317fn expiry(now: SystemTime, ttl: Duration) -> SystemTime {
320 now.checked_add(ttl).unwrap_or(now)
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326 use crate::id::{EffectName, LogicalKey};
327
328 const TTL: Duration = Duration::from_secs(30);
329
330 fn t(secs: u64) -> SystemTime {
331 SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000 + secs)
332 }
333
334 fn record() -> EffectRecord {
335 let key = EffectKey::new(
336 EffectName::new("payment.charge").unwrap(),
337 LogicalKey::new("order_1").unwrap(),
338 );
339 EffectRecord::new(NewEffect::new(key, EffectKind::IrreversibleWrite, t(0)))
340 }
341
342 fn worker(name: &str) -> WorkerId {
343 WorkerId::new(name)
344 }
345
346 #[test]
347 fn failed_apply_leaves_the_record_unchanged() {
348 let mut rec = record();
349 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
350 let before = rec.clone();
351
352 let mut bad = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(1));
353 bad.output = Some(Value::from(1));
354 assert!(matches!(
355 rec.apply(bad),
356 Err(StoreError::InvalidTransition(_))
357 ));
358 assert_eq!(rec, before);
359 }
360
361 #[test]
362 fn stale_epoch_is_fenced_off_after_takeover() {
363 let mut rec = record();
364 let old = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
365 let new = rec.acquire_lease(&worker("b"), t(30), TTL).unwrap();
366 assert_eq!(new.epoch, old.epoch + 1);
367
368 let request = TransitionRequest::new(&rec, Some(&old), Transition::StartAttempt, t(31));
369 assert!(matches!(rec.apply(request), Err(StoreError::LeaseLost)));
370 assert!(
371 !rec.release_lease(&old),
372 "stale release must not clear b's lease"
373 );
374 assert_eq!(rec.live_lease_owner(t(31)), Some(&worker("b")));
375 }
376
377 #[test]
378 fn same_owner_cannot_double_acquire() {
379 let mut rec = record();
380 rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
381 assert!(matches!(
382 rec.acquire_lease(&worker("a"), t(1), TTL),
383 Err(StoreError::LeaseHeld { .. })
384 ));
385 }
386
387 #[test]
388 fn overflowing_ttl_fails_safe() {
389 let mut rec = record();
390 let lease = rec
391 .acquire_lease(&worker("a"), t(0), Duration::MAX)
392 .unwrap();
393 assert_eq!(lease.expires_at, t(0));
394 assert_eq!(rec.live_lease_owner(t(0)), None);
395 }
396
397 #[test]
398 fn a_failed_attempt_cannot_fail_an_effect_that_may_have_applied() {
399 let mut rec = record();
400 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
401 for transition in [
402 Transition::StartAttempt,
403 Transition::OutcomeUnknown,
404 Transition::ScheduleRetry,
405 Transition::StartAttempt,
406 ] {
407 rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(1)))
408 .unwrap();
409 }
410 assert!(
411 rec.may_have_applied,
412 "the unknown outcome is remembered across retries"
413 );
414 let before = rec.clone();
415 let fail = TransitionRequest::new(&rec, Some(&lease), Transition::FailedDefinitively, t(2));
416 assert!(matches!(
417 rec.apply(fail),
418 Err(StoreError::InvalidTransition(_))
419 ));
420 assert_eq!(rec, before);
421
422 for transition in [
424 Transition::StartVerification,
425 Transition::ScheduleRetry,
426 Transition::StartAttempt,
427 ] {
428 rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(3)))
429 .unwrap();
430 }
431 assert!(!rec.may_have_applied);
432 rec.apply(TransitionRequest::new(
433 &rec,
434 Some(&lease),
435 Transition::FailedDefinitively,
436 t(4),
437 ))
438 .unwrap();
439 assert_eq!(rec.status, EffectStatus::Failed);
440 }
441
442 #[test]
443 fn bookkeeping_follows_transitions() {
444 let mut rec = record();
445 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
446
447 let event = rec
448 .apply(TransitionRequest::new(
449 &rec,
450 Some(&lease),
451 Transition::StartAttempt,
452 t(1),
453 ))
454 .unwrap();
455 assert_eq!((event.sequence, event.attempt), (1, 1));
456 assert_eq!(rec.attempt_started_at, Some(t(1)));
457
458 let mut retry = TransitionRequest::new(&rec, Some(&lease), Transition::ScheduleRetry, t(2));
459 retry.next_attempt_at = Some(t(10));
460 retry.error = Some(ErrorRecord {
461 class: Some(crate::FailureClass::Transient),
462 message: "connection refused".into(),
463 });
464 rec.apply(retry).unwrap();
465 assert_eq!(rec.next_attempt_at, Some(t(10)));
466 assert_eq!(rec.attempt_ended_at, Some(t(2)));
467 assert!(rec.last_error.is_some());
468
469 rec.apply(TransitionRequest::new(
470 &rec,
471 Some(&lease),
472 Transition::StartAttempt,
473 t(10),
474 ))
475 .unwrap();
476 assert_eq!(rec.attempt_count, 2);
477 assert_eq!(rec.next_attempt_at, None);
478 assert_eq!(rec.attempt_ended_at, None);
479
480 let mut done = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(11));
481 done.output = Some(serde_json::json!({ "payment": "pi_1" }));
482 let event = rec.apply(done).unwrap();
483 assert_eq!(event.to, EffectStatus::Committed);
484 assert_eq!(rec.committed_at, Some(t(11)));
485 assert_eq!(rec.version, 4);
486 }
487}