#![forbid(unsafe_code)]
#![warn(missing_docs)]
use serde::{Deserialize, Serialize};
use uuid::Uuid;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct Pact {
pub id: Uuid,
pub docket: String,
pub kind: String,
pub clause: Vec<u8>,
}
impl Pact {
#[must_use]
pub fn new(id: Uuid, docket: String, kind: String, clause: Vec<u8>) -> Self {
Self {
id,
docket,
kind,
clause,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct Retainer(Uuid);
impl Retainer {
#[must_use]
pub fn new(id: Uuid) -> Self {
Self(id)
}
#[must_use]
pub fn id(&self) -> Uuid {
self.0
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub struct Timestamp(u64);
impl Timestamp {
#[must_use]
pub fn from_millis(millis: u64) -> Self {
Self(millis)
}
#[must_use]
pub fn as_millis(self) -> u64 {
self.0
}
#[must_use]
pub fn plus_millis(self, millis: u64) -> Self {
Self(self.0.saturating_add(millis))
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[non_exhaustive]
pub struct Claim {
pub pact: Pact,
pub retainer: Retainer,
pub lease_expiry: Timestamp,
}
impl Claim {
#[must_use]
pub fn new(pact: Pact, retainer: Retainer, lease_expiry: Timestamp) -> Self {
Self {
pact,
retainer,
lease_expiry,
}
}
}
const _: fn() = || {
fn assert_key<T: Eq + std::hash::Hash>() {}
assert_key::<Retainer>();
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Outcome {
Fulfilled,
Breached,
}
pub type Settlement = Outcome;
pub mod lifecycle {
use crate::{Retainer, Timestamp};
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum State {
Available,
Held {
retainer: Retainer,
expiry: Timestamp,
},
Deferred {
reclaimable_at: Timestamp,
},
Settled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct NotCurrentHolder;
impl std::fmt::Display for NotCurrentHolder {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "retainer is not the current holder of this pact")
}
}
impl std::error::Error for NotCurrentHolder {}
#[must_use]
pub fn lease_expiry(now: Timestamp, lease_millis: u64) -> Timestamp {
now.plus_millis(lease_millis)
}
#[must_use]
pub fn is_claimable(state: &State, now: Timestamp) -> bool {
match state {
State::Available => true,
State::Held { expiry, .. } => *expiry < now,
State::Deferred { reclaimable_at } => *reclaimable_at <= now,
State::Settled => false,
}
}
#[must_use]
pub fn on_claim(retainer: &Retainer, now: Timestamp, lease_millis: u64) -> State {
State::Held {
retainer: retainer.clone(),
expiry: lease_expiry(now, lease_millis),
}
}
pub fn on_heartbeat(
state: &State,
retainer: &Retainer,
now: Timestamp,
lease_millis: u64,
) -> Result<State, NotCurrentHolder> {
match state {
State::Held {
retainer: held,
expiry,
} if held == retainer && *expiry >= now => Ok(State::Held {
retainer: retainer.clone(),
expiry: lease_expiry(now, lease_millis),
}),
_ => Err(NotCurrentHolder),
}
}
pub fn on_settle(state: &State, retainer: &Retainer) -> Result<State, NotCurrentHolder> {
if is_current_holder(state, retainer) {
Ok(State::Settled)
} else {
Err(NotCurrentHolder)
}
}
pub fn on_release(
state: &State,
retainer: &Retainer,
reclaimable_at: Timestamp,
) -> Result<State, NotCurrentHolder> {
if is_current_holder(state, retainer) {
Ok(State::Deferred { reclaimable_at })
} else {
Err(NotCurrentHolder)
}
}
fn is_current_holder(state: &State, retainer: &Retainer) -> bool {
matches!(state, State::Held { retainer: held, .. } if held == retainer)
}
#[cfg(test)]
mod tests {
use super::*;
use uuid::Uuid;
fn retainer() -> Retainer {
Retainer::new(Uuid::new_v4())
}
#[test]
fn eligibility_covers_each_state() {
let now = Timestamp::from_millis(100);
assert!(is_claimable(&State::Available, now));
assert!(!is_claimable(
&State::Held {
retainer: retainer(),
expiry: Timestamp::from_millis(101)
},
now
));
assert!(is_claimable(
&State::Held {
retainer: retainer(),
expiry: Timestamp::from_millis(99)
},
now
));
assert!(!is_claimable(
&State::Deferred {
reclaimable_at: Timestamp::from_millis(101)
},
now
));
assert!(is_claimable(
&State::Deferred {
reclaimable_at: Timestamp::from_millis(100)
},
now
));
assert!(!is_claimable(&State::Settled, now));
}
#[test]
fn transitions_require_the_current_holder() {
let holder = retainer();
let held = State::Held {
retainer: holder.clone(),
expiry: Timestamp::from_millis(200),
};
let stranger = retainer();
assert_eq!(on_settle(&held, &stranger), Err(NotCurrentHolder));
assert_eq!(
on_release(&held, &stranger, Timestamp::from_millis(0)),
Err(NotCurrentHolder)
);
assert_eq!(on_settle(&State::Settled, &holder), Err(NotCurrentHolder));
assert_eq!(on_settle(&held, &holder), Ok(State::Settled));
assert_eq!(
on_release(&held, &holder, Timestamp::from_millis(500)),
Ok(State::Deferred {
reclaimable_at: Timestamp::from_millis(500)
})
);
}
#[test]
fn heartbeat_refreshes_but_does_not_revive_a_lapsed_lease() {
let holder = retainer();
let held = State::Held {
retainer: holder.clone(),
expiry: Timestamp::from_millis(200),
};
assert_eq!(
on_heartbeat(&held, &holder, Timestamp::from_millis(150), 100),
Ok(State::Held {
retainer: holder.clone(),
expiry: Timestamp::from_millis(250)
})
);
assert_eq!(
on_heartbeat(&held, &holder, Timestamp::from_millis(201), 100),
Err(NotCurrentHolder)
);
}
}
}
pub type Transition<'a> = dyn Fn(&lifecycle::State) -> Result<lifecycle::State, lifecycle::NotCurrentHolder>
+ Send
+ Sync
+ 'a;
#[cfg(feature = "async")]
mod async_registry;
#[cfg(feature = "async")]
pub use async_registry::{AsyncRegistry, apply_via_cas};
pub trait Registry: Send + Sync {
type Error: std::error::Error;
fn claim(&self, dockets: &[&str], now: Timestamp) -> Result<Option<Claim>, Self::Error>;
fn lease_millis(&self) -> u64;
fn apply(&self, retainer: &Retainer, transition: &Transition<'_>) -> Result<(), Self::Error>;
fn heartbeat(&self, retainer: &Retainer, now: Timestamp) -> Result<(), Self::Error> {
let lease = self.lease_millis();
self.apply(retainer, &|state| {
lifecycle::on_heartbeat(state, retainer, now, lease)
})
}
fn fulfill(&self, retainer: &Retainer) -> Result<(), Self::Error> {
self.apply(retainer, &|state| lifecycle::on_settle(state, retainer))
}
fn breach(&self, retainer: &Retainer) -> Result<(), Self::Error> {
self.apply(retainer, &|state| lifecycle::on_settle(state, retainer))
}
fn release(&self, retainer: &Retainer, reclaimable_at: Timestamp) -> Result<(), Self::Error> {
self.apply(retainer, &|state| {
lifecycle::on_release(state, retainer, reclaimable_at)
})
}
}
pub mod kernel {
use crate::{Claim, Outcome, Pact, Retainer};
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum Directive {
Claim,
Execute(Pact),
Settle(Retainer, Outcome),
Idle,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum Notice {
Claimed(Option<Claim>),
Executed(Outcome),
ExecutionFailed,
Settled,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum StepResult {
Idle,
Settled(Outcome),
Unsettled,
}
#[derive(Debug)]
enum Phase {
Claiming,
Executing {
pact: Pact,
retainer: Retainer,
},
Settling {
retainer: Retainer,
outcome: Outcome,
},
DoneIdle,
DoneSettled(Outcome),
DoneUnsettled,
}
#[derive(Debug)]
pub struct Kernel {
phase: Phase,
}
impl Kernel {
#[must_use]
pub fn new() -> Self {
Self {
phase: Phase::Claiming,
}
}
#[must_use]
pub fn poll(&self) -> Directive {
match &self.phase {
Phase::Claiming => Directive::Claim,
Phase::Executing { pact, .. } => Directive::Execute(pact.clone()),
Phase::Settling { retainer, outcome } => {
Directive::Settle(retainer.clone(), *outcome)
}
Phase::DoneIdle | Phase::DoneSettled(_) | Phase::DoneUnsettled => Directive::Idle,
}
}
pub fn on_event(&mut self, notice: Notice) {
let phase = std::mem::replace(&mut self.phase, Phase::DoneIdle);
self.phase = match (phase, notice) {
(Phase::Claiming, Notice::Claimed(Some(claim))) => Phase::Executing {
pact: claim.pact,
retainer: claim.retainer,
},
(Phase::Claiming, Notice::Claimed(None)) => Phase::DoneIdle,
(Phase::Executing { retainer, .. }, Notice::Executed(outcome)) => {
Phase::Settling { retainer, outcome }
}
(Phase::Executing { .. }, Notice::ExecutionFailed) => Phase::DoneUnsettled,
(Phase::Settling { outcome, .. }, Notice::Settled) => Phase::DoneSettled(outcome),
(other, _) => other,
};
}
#[must_use]
pub fn result(&self) -> Option<StepResult> {
match &self.phase {
Phase::DoneIdle => Some(StepResult::Idle),
Phase::DoneSettled(outcome) => Some(StepResult::Settled(*outcome)),
Phase::DoneUnsettled => Some(StepResult::Unsettled),
_ => None,
}
}
}
impl Default for Kernel {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Timestamp;
use uuid::Uuid;
fn claim() -> Claim {
Claim::new(
Pact::new(
Uuid::new_v4(),
"default".to_string(),
"example".to_string(),
Vec::new(),
),
Retainer::new(Uuid::new_v4()),
Timestamp::from_millis(0),
)
}
fn drive(
kernel: &mut Kernel,
execution: Result<Outcome, ()>,
claimable: bool,
) -> StepResult {
loop {
if let Some(result) = kernel.result() {
return result;
}
match kernel.poll() {
Directive::Claim => {
let notice = if claimable {
Notice::Claimed(Some(claim()))
} else {
Notice::Claimed(None)
};
kernel.on_event(notice);
}
Directive::Execute(_) => kernel.on_event(match execution {
Ok(outcome) => Notice::Executed(outcome),
Err(()) => Notice::ExecutionFailed,
}),
Directive::Settle(_, _) => kernel.on_event(Notice::Settled),
Directive::Idle => return StepResult::Idle,
}
}
}
#[test]
fn fulfilled_execution_settles_fulfilled() {
let mut kernel = Kernel::new();
assert_eq!(
drive(&mut kernel, Ok(Outcome::Fulfilled), true),
StepResult::Settled(Outcome::Fulfilled)
);
}
#[test]
fn breached_execution_settles_breached() {
let mut kernel = Kernel::new();
assert_eq!(
drive(&mut kernel, Ok(Outcome::Breached), true),
StepResult::Settled(Outcome::Breached)
);
}
#[test]
fn infrastructure_error_is_unsettled() {
let mut kernel = Kernel::new();
assert_eq!(drive(&mut kernel, Err(()), true), StepResult::Unsettled);
}
#[test]
fn empty_claim_is_idle() {
let mut kernel = Kernel::new();
assert_eq!(
drive(&mut kernel, Ok(Outcome::Fulfilled), false),
StepResult::Idle
);
}
}
}