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())
}
}
#[doc(hidden)]
#[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 sink_arguments(&self) -> Option<&serde_json::Value> {
Some(&self.payload)
}
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()),
})
}
}
#[derive(Debug, Clone)]
pub struct RecallMemory {
pub(crate) memories: std::sync::Arc<dyn crate::memory::MemoryStore>,
pub(crate) query: crate::memory::Recall,
}
#[derive(Debug, Clone)]
pub struct Embed {
pub(crate) embedder: std::sync::Arc<dyn crate::memory::Embedder>,
pub(crate) text: String,
pub(crate) arguments: serde_json::Value,
}
#[async_trait]
impl Effect for Embed {
type Output = Vec<f32>;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"memory.embed",
json!({
"revision": self.embedder.revision(),
"text": self.text,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn sink_arguments(&self) -> Option<&serde_json::Value> {
Some(&self.arguments)
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
self.embedder
.embed(&self.text)
.await
.map_err(|error| EffectError::Other(error.to_string()))
}
}
#[derive(Debug, Clone)]
pub struct SemanticRecall {
pub(crate) retriever: std::sync::Arc<dyn crate::memory::SemanticRetriever>,
pub(crate) query: crate::memory::SemanticQuery,
pub(crate) arguments: serde_json::Value,
}
#[async_trait]
impl Effect for SemanticRecall {
type Output = Vec<crate::memory::SemanticHit>;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"memory.semantic-recall",
json!({
"retriever": self.retriever.profile(),
"query": self.query,
}),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn max_sensitivity(&self) -> crate::core::Sensitivity {
self.query.max_sensitivity
}
fn sink_arguments(&self) -> Option<&serde_json::Value> {
Some(&self.arguments)
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
self.retriever
.search(&self.query)
.await
.map_err(|error| EffectError::Other(error.to_string()))
}
}
#[async_trait]
impl Effect for RecallMemory {
type Output = Vec<crate::memory::Selected>;
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"memory.recall",
serde_json::to_value(&self.query).unwrap_or(json!(null)),
)
}
fn mutates(&self) -> bool {
false
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
let found = self
.memories
.recall(&self.query)
.await
.map_err(|e| EffectError::Other(e.to_string()))?;
Ok(found
.iter()
.map(|item| crate::memory::Selected {
id: item.id.clone(),
version: item.version,
digest: item.selection_digest(),
})
.collect())
}
}
#[derive(Debug, Clone)]
pub struct TouchMemory {
pub(crate) memories: std::sync::Arc<dyn crate::memory::MemoryStore>,
pub(crate) ids: Vec<String>,
pub(crate) at: crate::core::Timestamp,
}
#[async_trait]
impl Effect for TouchMemory {
type Output = ();
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("memory.touch", json!({ "ids": self.ids, "at": self.at }))
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
self.memories
.touch(&self.ids, self.at)
.await
.map_err(|error| EffectError::Other(error.to_string()))
}
}
#[derive(Debug, Clone)]
pub struct RememberMemory {
pub(crate) memories: std::sync::Arc<dyn crate::memory::MemoryStore>,
pub(crate) item: crate::memory::MemoryItem,
}
#[async_trait]
impl Effect for RememberMemory {
type Output = u64;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"memory.remember",
json!({
"id": self.item.id,
"subject": self.item.subject,
"purpose": self.item.purpose,
"provenance": self.item.provenance,
"sensitivity": self.item.sensitivity,
"trust": self.item.trust,
"created_at": self.item.created_at,
"expires_at": self.item.expires_at,
"derived_from": self.item.derived_from,
"content": self.item.digest().to_hex(),
}),
)
}
async fn perform(&self) -> Result<u64, EffectError> {
self.memories
.remember(&self.item)
.await
.map_err(|e| EffectError::Other(e.to_string()))
}
}
#[derive(Debug, Clone)]
pub struct SweepExpiredMemory {
pub(crate) memories: std::sync::Arc<dyn crate::memory::MemoryStore>,
pub(crate) at: crate::core::Timestamp,
}
#[async_trait]
impl Effect for SweepExpiredMemory {
type Output = usize;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new("memory.sweep-expired", json!({ "at": self.at }))
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
self.memories
.sweep_expired(self.at)
.await
.map_err(|error| EffectError::Other(error.to_string()))
}
}
#[derive(Debug, Clone)]
pub struct SetCaseStatus {
pub(crate) cases: std::sync::Arc<dyn crate::case::CaseStore>,
pub(crate) case: crate::core::CaseId,
pub(crate) status: crate::core::CaseStatus,
}
#[async_trait]
impl Effect for SetCaseStatus {
type Output = ();
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"case.set_status",
json!({ "case": self.case.to_string(), "status": self.status }),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<(), EffectError> {
self.cases
.set_status(self.case, self.status)
.await
.map_err(|e| EffectError::Other(e.to_string()))
}
}
#[derive(Debug, Clone)]
pub struct TransitionDeadline {
pub(crate) cases: std::sync::Arc<dyn crate::case::CaseStore>,
pub(crate) case: crate::core::CaseId,
pub(crate) name: String,
pub(crate) to: crate::core::DeadlineState,
}
#[async_trait]
impl Effect for TransitionDeadline {
type Output = crate::core::DeadlineState;
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"case.transition_deadline",
json!({
"case": self.case.to_string(),
"name": self.name,
"to": self.to,
}),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<crate::core::DeadlineState, EffectError> {
let before = self
.cases
.deadlines(self.case)
.await
.map_err(|e| EffectError::Other(e.to_string()))?
.into_iter()
.find(|d| d.name == self.name)
.map_or(crate::core::DeadlineState::Pending, |d| d.state);
self.cases
.set_deadline_state(self.case, &self.name, self.to)
.await
.map_err(|e| EffectError::Other(e.to_string()))?;
Ok(before)
}
}
#[derive(Debug, Clone)]
pub struct DrawOnAuthority {
pub(crate) authorities: std::sync::Arc<dyn crate::authority::AuthorityStore>,
pub(crate) id: crate::authority::AuthorityId,
pub(crate) amount: crate::core::Spend,
pub(crate) at: Timestamp,
pub(crate) key: Option<crate::core::EffectKey>,
}
#[async_trait]
impl Effect for DrawOnAuthority {
type Output = crate::authority::Drawn;
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"authority.draw",
json!({
"authority": self.id,
"amount": self.amount,
"at": self.at,
}),
)
}
fn attach(&mut self, provenance: &crate::core::Provenance) {
self.key = Some(provenance.dispatch.unwrap_or(provenance.effect));
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
let key = self.key.ok_or_else(|| {
EffectError::Other(
"a draw reached the store without its dispatch key — deduplicating on \
anything else would let a retry spend the authority twice"
.to_owned(),
)
})?;
self.authorities
.draw(&self.id, key, self.amount, self.at)
.await
.map_err(|error| EffectError::Other(error.to_string()))
}
}