pub mod failure;
pub mod id;
pub mod kind;
pub mod state;
#[cfg(feature = "testkit")]
pub mod testkit;
mod error;
mod record;
use std::future::Future;
use std::time::{Duration, SystemTime};
use serde::{Deserialize, Serialize};
use serde_json::Value;
pub use error::StoreError;
pub use failure::{Disposition, FailureClass};
pub use id::{
EffectId, EffectKey, EffectName, IdempotencyKey, IdentityError, LogicalKey, WorkerId,
};
pub use kind::EffectKind;
pub use record::EffectRecord;
pub use state::{EffectStatus, InvalidTransition, Transition};
pub trait EffectStore: Send + Sync + 'static {
fn insert_or_get(
&self,
new: NewEffect,
) -> impl Future<Output = Result<InsertOutcome, StoreError>> + Send;
fn get(
&self,
id: EffectId,
) -> impl Future<Output = Result<Option<EffectRecord>, StoreError>> + Send;
fn get_by_key(
&self,
key: &EffectKey,
) -> impl Future<Output = Result<Option<EffectRecord>, StoreError>> + Send;
fn acquire_lease(
&self,
id: EffectId,
owner: &WorkerId,
now: SystemTime,
ttl: Duration,
) -> impl Future<Output = Result<Lease, StoreError>> + Send;
fn renew_lease(
&self,
lease: &Lease,
now: SystemTime,
ttl: Duration,
) -> impl Future<Output = Result<Lease, StoreError>> + Send;
fn release_lease(&self, lease: &Lease) -> impl Future<Output = Result<(), StoreError>> + Send;
fn transition(
&self,
request: TransitionRequest,
) -> impl Future<Output = Result<EffectRecord, StoreError>> + Send;
fn list(
&self,
query: ListQuery,
) -> impl Future<Output = Result<Vec<EffectRecord>, StoreError>> + Send;
fn events(
&self,
id: EffectId,
) -> impl Future<Output = Result<Vec<EffectEvent>, StoreError>> + Send;
fn prune(&self, query: PruneQuery) -> impl Future<Output = Result<u64, StoreError>> + Send;
}
#[derive(Clone, Debug, PartialEq)]
pub struct NewEffect {
pub id: EffectId,
pub key: EffectKey,
pub kind: EffectKind,
pub input: Option<Value>,
pub input_fingerprint: Option<String>,
pub created_by: Option<String>,
pub now: SystemTime,
}
impl NewEffect {
pub fn new(key: EffectKey, kind: EffectKind, now: SystemTime) -> Self {
Self {
id: EffectId::new(),
key,
kind,
input: None,
input_fingerprint: None,
created_by: None,
now,
}
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct InsertOutcome {
pub record: EffectRecord,
pub inserted: bool,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct Lease {
pub effect_id: EffectId,
pub owner: WorkerId,
pub epoch: u64,
pub expires_at: SystemTime,
}
#[derive(Clone, Debug, PartialEq)]
pub struct TransitionRequest {
pub id: EffectId,
pub expected_version: u64,
pub lease: Option<Lease>,
pub transition: Transition,
pub now: SystemTime,
pub output: Option<Value>,
pub error: Option<ErrorRecord>,
pub next_attempt_at: Option<SystemTime>,
pub actor: Option<String>,
pub payload: Option<Value>,
}
impl TransitionRequest {
pub fn new(
record: &EffectRecord,
lease: Option<&Lease>,
transition: Transition,
now: SystemTime,
) -> Self {
Self {
id: record.id,
expected_version: record.version,
lease: lease.cloned(),
transition,
now,
output: None,
error: None,
next_attempt_at: None,
actor: None,
payload: None,
}
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct ErrorRecord {
pub class: Option<FailureClass>,
pub message: String,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct EffectEvent {
pub effect_id: EffectId,
pub sequence: u64,
pub transition: Transition,
pub from: EffectStatus,
pub to: EffectStatus,
pub attempt: u32,
pub actor: Option<String>,
pub payload: Option<Value>,
pub at: SystemTime,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PruneQuery {
pub status: EffectStatus,
pub older_than: Duration,
pub now: SystemTime,
pub limit: usize,
}
impl PruneQuery {
pub const DEFAULT_LIMIT: usize = 500;
pub fn new(status: EffectStatus, older_than: Duration, now: SystemTime) -> Self {
Self {
status,
older_than,
now,
limit: Self::DEFAULT_LIMIT,
}
}
#[must_use]
pub fn limit(mut self, limit: usize) -> Self {
self.limit = limit;
self
}
pub fn cutoff(&self) -> Option<SystemTime> {
self.now.checked_sub(self.older_than)
}
pub fn matches(&self, record: &EffectRecord) -> bool {
self.status.is_settled()
&& record.status == self.status
&& self
.cutoff()
.is_some_and(|cutoff| record.updated_at <= cutoff)
&& record.live_lease_owner(self.now).is_none()
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ListQuery {
pub statuses: Vec<EffectStatus>,
pub lease_expired_at: Option<SystemTime>,
pub after: Option<EffectId>,
pub limit: usize,
}
impl ListQuery {
pub const DEFAULT_LIMIT: usize = 100;
pub fn statuses(statuses: impl IntoIterator<Item = EffectStatus>) -> Self {
Self {
statuses: statuses.into_iter().collect(),
lease_expired_at: None,
after: None,
limit: Self::DEFAULT_LIMIT,
}
}
pub fn expired_leases(now: SystemTime) -> Self {
Self {
lease_expired_at: Some(now),
..Self::statuses([EffectStatus::Executing, EffectStatus::Verifying])
}
}
#[must_use]
pub fn after(mut self, id: EffectId) -> Self {
self.after = Some(id);
self
}
#[must_use]
pub fn limit(mut self, limit: usize) -> Self {
self.limit = limit;
self
}
pub fn matches(&self, record: &EffectRecord) -> bool {
(self.statuses.is_empty() || self.statuses.contains(&record.status))
&& self
.lease_expired_at
.is_none_or(|now| record.live_lease_owner(now).is_none())
&& self.after.is_none_or(|after| record.id > after)
}
}