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> {
210 let now = request.now;
211 match &request.lease {
212 Some(lease) => self.check_lease(lease, now)?,
213 None => {
214 if let Some(owner) = self.live_lease_owner(now) {
215 return Err(StoreError::LeaseHeld {
216 owner: owner.clone(),
217 expires_at: self.lease_expires_at.unwrap_or(now),
218 });
219 }
220 }
221 }
222 if request.expected_version != self.version {
223 return Err(StoreError::VersionConflict {
224 expected: request.expected_version,
225 actual: self.version,
226 });
227 }
228 let from = self.status;
229 let to = from.apply(request.transition)?;
230 if request.transition == Transition::FailedDefinitively
231 && self.may_have_applied
232 && self.kind != EffectKind::Read
233 {
234 return Err(InvalidTransition {
235 from,
236 event: request.transition,
237 }
238 .into());
239 }
240
241 match request.transition {
242 Transition::StartAttempt => {
243 self.attempt_count = self.attempt_count.saturating_add(1);
244 self.attempt_started_at = Some(now);
245 self.attempt_ended_at = None;
246 self.next_attempt_at = None;
247 self.output = None;
248 }
249 Transition::ScheduleRetry
250 | Transition::ResolvedRetry
251 | Transition::ScheduleCompensationRetry => {
252 self.next_attempt_at = Some(request.next_attempt_at.unwrap_or(now));
253 }
254 Transition::Approve => self.approved = true,
255 Transition::StartCompensation => {
256 self.compensation_attempts = 1;
257 self.next_attempt_at = None;
258 }
259 Transition::StartCompensationRetry => {
260 self.compensation_attempts = self.compensation_attempts.saturating_add(1);
261 self.next_attempt_at = None;
262 }
263 _ => {}
264 }
265 if from == EffectStatus::Executing {
266 self.attempt_ended_at = Some(now);
267 }
268 if to == EffectStatus::Committed {
269 self.committed_at = Some(now);
270 }
271 if to == EffectStatus::Unknown {
272 self.may_have_applied = true;
273 }
274 let shown_not_applied = matches!(
275 request.transition,
276 Transition::VerificationNotApplied | Transition::ResolvedNotApplied
277 ) || (from == EffectStatus::Verifying
278 && request.transition == Transition::ScheduleRetry);
279 if shown_not_applied {
280 self.may_have_applied = false;
281 }
282 if let Some(output) = request.output {
283 self.output = Some(output);
284 }
285 if let Some(error) = request.error {
286 self.last_error = Some(error);
287 }
288 self.status = to;
289 self.version += 1;
290 self.updated_at = now;
291
292 Ok(EffectEvent {
293 effect_id: self.id,
294 sequence: self.version,
295 transition: request.transition,
296 from,
297 to,
298 attempt: self.attempt_count,
299 actor: request.actor,
300 payload: request.payload,
301 at: now,
302 })
303 }
304
305 fn check_lease(&self, lease: &Lease, now: SystemTime) -> Result<(), StoreError> {
306 let valid = lease.effect_id == self.id
307 && self.lease_epoch == lease.epoch
308 && self.lease_owner.as_ref() == Some(&lease.owner)
309 && self
310 .lease_expires_at
311 .is_some_and(|expires_at| expires_at > now);
312 if valid {
313 Ok(())
314 } else {
315 Err(StoreError::LeaseLost)
316 }
317 }
318}
319
320fn expiry(now: SystemTime, ttl: Duration) -> SystemTime {
323 now.checked_add(ttl).unwrap_or(now)
324}
325
326#[cfg(test)]
327mod tests {
328 use super::*;
329 use crate::id::{EffectName, LogicalKey};
330
331 const TTL: Duration = Duration::from_secs(30);
332
333 fn t(secs: u64) -> SystemTime {
334 SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000 + secs)
335 }
336
337 fn record() -> EffectRecord {
338 let key = EffectKey::new(
339 EffectName::new("payment.charge").unwrap(),
340 LogicalKey::new("order_1").unwrap(),
341 );
342 EffectRecord::new(NewEffect::new(key, EffectKind::IrreversibleWrite, t(0)))
343 }
344
345 fn worker(name: &str) -> WorkerId {
346 WorkerId::new(name)
347 }
348
349 #[test]
350 fn failed_apply_leaves_the_record_unchanged() {
351 let mut rec = record();
352 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
353 let before = rec.clone();
354
355 let mut bad = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(1));
356 bad.output = Some(Value::from(1));
357 assert!(matches!(
358 rec.apply(bad),
359 Err(StoreError::InvalidTransition(_))
360 ));
361 assert_eq!(rec, before);
362 }
363
364 #[test]
365 fn stale_epoch_is_fenced_off_after_takeover() {
366 let mut rec = record();
367 let old = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
368 let new = rec.acquire_lease(&worker("b"), t(30), TTL).unwrap();
369 assert_eq!(new.epoch, old.epoch + 1);
370
371 let request = TransitionRequest::new(&rec, Some(&old), Transition::StartAttempt, t(31));
372 assert!(matches!(rec.apply(request), Err(StoreError::LeaseLost)));
373 assert!(
374 !rec.release_lease(&old),
375 "stale release must not clear b's lease"
376 );
377 assert_eq!(rec.live_lease_owner(t(31)), Some(&worker("b")));
378 }
379
380 #[test]
381 fn same_owner_cannot_double_acquire() {
382 let mut rec = record();
383 rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
384 assert!(matches!(
385 rec.acquire_lease(&worker("a"), t(1), TTL),
386 Err(StoreError::LeaseHeld { .. })
387 ));
388 }
389
390 #[test]
391 fn overflowing_ttl_fails_safe() {
392 let mut rec = record();
393 let lease = rec
394 .acquire_lease(&worker("a"), t(0), Duration::MAX)
395 .unwrap();
396 assert_eq!(lease.expires_at, t(0));
397 assert_eq!(rec.live_lease_owner(t(0)), None);
398 }
399
400 #[test]
401 fn a_failed_attempt_cannot_fail_an_effect_that_may_have_applied() {
402 let mut rec = record();
403 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
404 for transition in [
405 Transition::StartAttempt,
406 Transition::OutcomeUnknown,
407 Transition::ScheduleRetry,
408 Transition::StartAttempt,
409 ] {
410 rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(1)))
411 .unwrap();
412 }
413 assert!(
414 rec.may_have_applied,
415 "the unknown outcome is remembered across retries"
416 );
417 let before = rec.clone();
418 let fail = TransitionRequest::new(&rec, Some(&lease), Transition::FailedDefinitively, t(2));
419 assert!(matches!(
420 rec.apply(fail),
421 Err(StoreError::InvalidTransition(_))
422 ));
423 assert_eq!(rec, before);
424
425 for transition in [
427 Transition::StartVerification,
428 Transition::ScheduleRetry,
429 Transition::StartAttempt,
430 ] {
431 rec.apply(TransitionRequest::new(&rec, Some(&lease), transition, t(3)))
432 .unwrap();
433 }
434 assert!(!rec.may_have_applied);
435 rec.apply(TransitionRequest::new(
436 &rec,
437 Some(&lease),
438 Transition::FailedDefinitively,
439 t(4),
440 ))
441 .unwrap();
442 assert_eq!(rec.status, EffectStatus::Failed);
443 }
444
445 #[test]
446 fn bookkeeping_follows_transitions() {
447 let mut rec = record();
448 let lease = rec.acquire_lease(&worker("a"), t(0), TTL).unwrap();
449
450 let event = rec
451 .apply(TransitionRequest::new(
452 &rec,
453 Some(&lease),
454 Transition::StartAttempt,
455 t(1),
456 ))
457 .unwrap();
458 assert_eq!((event.sequence, event.attempt), (1, 1));
459 assert_eq!(rec.attempt_started_at, Some(t(1)));
460
461 let mut retry = TransitionRequest::new(&rec, Some(&lease), Transition::ScheduleRetry, t(2));
462 retry.next_attempt_at = Some(t(10));
463 retry.error = Some(ErrorRecord {
464 class: Some(crate::FailureClass::Transient),
465 message: "connection refused".into(),
466 });
467 rec.apply(retry).unwrap();
468 assert_eq!(rec.next_attempt_at, Some(t(10)));
469 assert_eq!(rec.attempt_ended_at, Some(t(2)));
470 assert!(rec.last_error.is_some());
471
472 rec.apply(TransitionRequest::new(
473 &rec,
474 Some(&lease),
475 Transition::StartAttempt,
476 t(10),
477 ))
478 .unwrap();
479 assert_eq!(rec.attempt_count, 2);
480 assert_eq!(rec.next_attempt_at, None);
481 assert_eq!(rec.attempt_ended_at, None);
482
483 let mut done = TransitionRequest::new(&rec, Some(&lease), Transition::Succeeded, t(11));
484 done.output = Some(serde_json::json!({ "payment": "pi_1" }));
485 let event = rec.apply(done).unwrap();
486 assert_eq!(event.to, EffectStatus::Committed);
487 assert_eq!(rec.committed_at, Some(t(11)));
488 assert_eq!(rec.version, 4);
489 }
490}