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,
pub(crate) max_sensitivity: Sensitivity,
}
#[async_trait]
impl Effect for Embed {
type Output = crate::memory::Embedding;
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 max_sensitivity(&self) -> Sensitivity {
self.max_sensitivity
}
fn sink_arguments(&self) -> Option<&serde_json::Value> {
Some(&self.arguments)
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
let vector = self
.embedder
.embed(&self.text)
.await
.map_err(|error| EffectError::Other(error.to_string()))?;
Ok(crate::memory::Embedding {
vector,
revision: self.embedder.revision(),
})
}
}
#[derive(Debug, Clone)]
pub struct SemanticRecall {
pub(crate) retriever: std::sync::Arc<dyn crate::memory::SemanticRetriever>,
pub(crate) memories: std::sync::Arc<dyn crate::memory::MemoryStore>,
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> {
let hits = self
.retriever
.search(&self.query)
.await
.map_err(|error| EffectError::Other(error.to_string()))?;
if hits.len() > self.query.limit {
return Err(EffectError::Other(format!(
"semantic retriever returned {} hits for a limit of {}",
hits.len(),
self.query.limit
)));
}
let mut survivors = Vec::with_capacity(hits.len());
for hit in hits {
if !hit.score.is_finite() {
return Err(EffectError::Other(
"semantic retriever returned a non-finite score".to_owned(),
));
}
let current = self
.memories
.current(&hit.selected.id, Some(self.query.as_of))
.await
.map_err(|error| EffectError::Other(error.to_string()))?;
if current.is_some_and(|item| item.version == hit.selected.version) {
survivors.push(hit);
}
}
Ok(survivors)
}
}
#[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": crate::core::format_timestamp(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": crate::core::format_timestamp(self.item.created_at),
"expires_at": self.item.expires_at.map(crate::core::format_timestamp),
"access_retention_seconds": self.item.access_retention_seconds,
"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": crate::core::format_timestamp(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 OpenTask {
pub(crate) tasks: std::sync::Arc<dyn crate::case::TaskStore>,
pub(crate) run: crate::core::RunId,
pub(crate) case: crate::core::CaseId,
pub(crate) spec: crate::core::TaskSpec,
pub(crate) at: Timestamp,
pub(crate) due_at: Timestamp,
pub(crate) key: Option<crate::core::EffectKey>,
}
#[async_trait]
impl Effect for OpenTask {
type Output = crate::core::TaskId;
fn trust(&self) -> crate::core::Trust {
crate::core::Trust::Trusted
}
fn descriptor(&self) -> EffectDescriptor {
EffectDescriptor::new(
"task.open",
json!({
"case": self.case.to_string(),
"kind": self.spec.kind,
"deadline": self.spec.deadline,
"roles": self.spec.candidate_roles,
"escalate_to": self.spec.escalate_to,
"priority": self.spec.priority.as_str(),
"justification": self.spec.justification,
}),
)
}
fn mutates(&self) -> bool {
true
}
fn recovery(&self) -> Recovery {
Recovery::Retry
}
fn max_sensitivity(&self) -> Sensitivity {
Sensitivity::Secret
}
fn attach(&mut self, provenance: &crate::core::Provenance) {
self.key = Some(provenance.dispatch.unwrap_or(provenance.effect));
}
async fn perform(&self) -> Result<Self::Output, EffectError> {
let key = self.key.ok_or_else(|| {
EffectError::Other(
"a task was opened without its effect key — an id that is not derived \
from the run would open a second row on every resume"
.to_owned(),
)
})?;
let id = crate::core::TaskId::derive(self.run, key);
self.tasks
.open(&crate::core::Task {
id,
run: self.run,
case: Some(self.case),
kind: self.spec.kind.clone(),
justification: self.spec.justification.clone(),
candidate_roles: self.spec.candidate_roles.clone(),
escalate_to: self.spec.escalate_to.clone(),
excluded_actors: self.spec.excluded_actors.clone(),
assignee: None,
priority: self.spec.priority,
state: crate::core::TaskState::Open,
on_expiry: self.spec.on_expiry,
created_at: self.at,
due_at: Some(self.due_at),
})
.await
.map_err(|e| EffectError::Other(e.to_string()))?;
Ok(id)
}
}
#[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": crate::core::format_timestamp(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| match &error {
crate::authority::AuthorityError::Unavailable(_) => EffectError::Unavailable {
driver: "authority-store".to_owned(),
detail: error.to_string(),
},
_ => EffectError::Refused(error.to_string()),
})
}
}