use async_trait::async_trait;
use serde_json::json;
use crate::core::{Effect, EffectDescriptor, EffectError, Recovery, Sensitivity, Timestamp};
#[derive(Debug, Clone, Copy)]
pub struct Clock;
#[async_trait]
impl Effect for Clock {
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
type Output = Timestamp;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::nullary("clock.now")
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn max_sensitivity(&self) -> Sensitivity {
Sensitivity::Public
}
#[allow(clippy::disallowed_methods)]
async fn perform(&self) -> Result<Timestamp, EffectError> {
Ok(Timestamp::now_utc())
}
}
#[derive(Debug, Clone)]
pub struct Recorded {
pub name: String,
pub payload: serde_json::Value,
pub calls: std::sync::Arc<std::sync::atomic::AtomicUsize>,
pub mutates: bool,
pub max_sensitivity: Sensitivity,
}
impl Recorded {
pub fn new(name: impl Into<String>) -> Self {
Self {
name: name.into(),
payload: json!(null),
calls: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)),
mutates: true,
max_sensitivity: Sensitivity::Secret,
}
}
#[must_use]
pub fn counter(mut self, c: std::sync::Arc<std::sync::atomic::AtomicUsize>) -> Self {
self.calls = c;
self
}
#[must_use]
pub fn payload(mut self, v: serde_json::Value) -> Self {
self.payload = v;
self
}
#[must_use]
pub fn ceiling(mut self, s: Sensitivity) -> Self {
self.max_sensitivity = s;
self
}
#[must_use]
pub fn read_only(mut self) -> Self {
self.mutates = false;
self
}
}
#[async_trait]
impl Effect for Recorded {
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
type Output = serde_json::Value;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(format!("test.{}", self.name), self.payload.clone())
}
fn mutates(&self) -> bool {
self.mutates
}
fn max_sensitivity(&self) -> Sensitivity {
self.max_sensitivity
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<serde_json::Value, EffectError> {
let n = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
Ok(json!({ "call": n, "payload": self.payload }))
}
}
#[derive(Debug, Clone)]
pub struct ResolveDeadline {
pub(crate) calendar: std::sync::Arc<dyn crate::core::Calendar>,
pub(crate) name: String,
pub(crate) from: Timestamp,
pub(crate) spec: crate::core::DeadlineSpec,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct ResolvedDeadline {
#[serde(with = "time::serde::rfc3339")]
pub at: Timestamp,
pub calendar_digest: crate::core::Digest,
}
#[async_trait]
impl Effect for ResolveDeadline {
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
type Output = ResolvedDeadline;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"deadline.resolve",
json!({
"name": self.name,
"spec": self.spec,
"from": self.from.unix_timestamp(),
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<ResolvedDeadline, EffectError> {
let at = self
.calendar
.resolve(self.from, &self.spec)
.map_err(|e| EffectError::Rejected(e.to_string()))?;
let at = at
.replace_nanosecond(0)
.map_err(|e| EffectError::Rejected(format!("deadline instant out of range: {e}")))?;
Ok(ResolvedDeadline {
at,
calendar_digest: self.calendar.digest(),
})
}
}
#[derive(Debug, Clone)]
pub struct ReadCaseState {
pub(crate) cases: std::sync::Arc<dyn crate::case::CaseStore>,
pub(crate) case: crate::core::CaseId,
}
#[derive(Debug, Clone, PartialEq, serde::Serialize, serde::Deserialize)]
pub struct CaseSnapshot {
pub state: serde_json::Value,
pub version: crate::core::CaseVersion,
}
#[async_trait]
impl Effect for ReadCaseState {
type Output = CaseSnapshot;
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("case.read_state", json!({ "case": self.case.to_string() }))
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<CaseSnapshot, EffectError> {
let case = self
.cases
.case(self.case)
.await
.map_err(|e| EffectError::Other(e.to_string()))?
.ok_or_else(|| EffectError::Rejected(format!("no case {}", self.case)))?;
Ok(CaseSnapshot {
state: case.state,
version: case.version,
})
}
}
#[derive(Debug, Clone)]
pub struct WriteCaseState {
pub(crate) cases: std::sync::Arc<dyn crate::case::CaseStore>,
pub(crate) case: crate::core::CaseId,
pub(crate) expected: crate::core::CaseVersion,
pub(crate) state: serde_json::Value,
}
#[async_trait]
impl Effect for WriteCaseState {
type Output = crate::core::CaseVersion;
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"case.write_state",
json!({
"case": self.case.to_string(),
"expected": self.expected.0,
"state": self.state,
}),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Reconcile
}
async fn perform(&self) -> Result<crate::core::CaseVersion, EffectError> {
self.cases
.put_state(self.case, self.expected, self.state.clone())
.await
.map_err(|e| match e {
crate::core::StoreError::CaseConflict { .. } => {
EffectError::Rejected(e.to_string())
}
other => EffectError::Other(other.to_string()),
})
}
}