use std::collections::HashMap;
use std::future::Future;
use std::ops::{Deref, DerefMut};
use std::sync::Arc;
use std::time::Duration;
use taquba::object_store::memory::InMemory;
use taquba::{LeaseHandle, PermanentFailure, WorkerError};
use tokio_util::sync::CancellationToken;
use crate::effects::EffectsHandle;
use crate::keys::RunId;
use crate::kv::KvReadHandle;
use crate::memo::{Memo, MemoStore};
#[derive(Debug, Clone)]
pub struct Delivery {
pub run_id: RunId,
pub headers: HashMap<String, String>,
pub job_id: String,
pub attempts: u32,
pub max_attempts: u32,
pub cancel_token: CancellationToken,
pub lease: LeaseHandle,
pub memo: Memo,
pub run_memo: Memo,
pub effects: EffectsHandle,
pub kv: KvReadHandle,
}
impl Delivery {
pub fn is_last_attempt(&self) -> bool {
self.attempts >= self.max_attempts
}
pub fn detached() -> Self {
let run_id = RunId::new("detached").expect("a literal run id");
let memo_store = MemoStore::new(Arc::new(InMemory::new()), "memo");
Self {
run_id: run_id.clone(),
headers: HashMap::new(),
job_id: "detached".to_string(),
attempts: 1,
max_attempts: 3,
cancel_token: CancellationToken::new(),
lease: LeaseHandle::detached(),
memo: memo_store.new_memo(&run_id, 0),
run_memo: memo_store.new_run_memo(&run_id),
effects: EffectsHandle::detached(),
kv: KvReadHandle::detached(),
}
}
}
#[derive(Debug, Clone)]
pub struct Step {
pub delivery: Delivery,
pub step_number: u32,
pub payload: Vec<u8>,
pub signal: Option<Vec<u8>>,
}
impl Deref for Step {
type Target = Delivery;
fn deref(&self) -> &Delivery {
&self.delivery
}
}
impl DerefMut for Step {
fn deref_mut(&mut self) -> &mut Delivery {
&mut self.delivery
}
}
impl Step {
pub fn detached(payload: impl Into<Vec<u8>>) -> Self {
Self {
delivery: Delivery::detached(),
step_number: 0,
payload: payload.into(),
signal: None,
}
}
}
#[derive(Debug, Clone)]
pub enum Trigger {
Immediate,
After(Duration),
OnSignal {
correlation_key: String,
timeout: Duration,
},
}
#[derive(Debug, Clone)]
pub enum StepOutcome {
Continue {
payload: Vec<u8>,
when: Trigger,
},
Succeed {
result: Vec<u8>,
},
Fail {
reason: String,
},
Cancel {
reason: String,
},
}
impl StepOutcome {
pub fn continue_now(payload: Vec<u8>) -> Self {
Self::Continue {
payload,
when: Trigger::Immediate,
}
}
pub fn continue_after(payload: Vec<u8>, delay: Duration) -> Self {
Self::Continue {
payload,
when: Trigger::After(delay),
}
}
pub fn continue_on_signal(
payload: Vec<u8>,
correlation_key: impl Into<String>,
timeout: Duration,
) -> Self {
Self::Continue {
payload,
when: Trigger::OnSignal {
correlation_key: correlation_key.into(),
timeout,
},
}
}
}
#[derive(Debug, Clone)]
pub struct StepError {
pub message: String,
pub kind: StepErrorKind,
}
impl StepError {
pub fn transient(message: impl Into<String>) -> Self {
Self {
message: message.into(),
kind: StepErrorKind::Transient,
}
}
pub fn permanent(message: impl Into<String>) -> Self {
Self {
message: message.into(),
kind: StepErrorKind::Permanent,
}
}
pub(crate) fn into_worker_error(self) -> WorkerError {
match self.kind {
StepErrorKind::Permanent => PermanentFailure::new(self.message).into(),
StepErrorKind::Transient => self.message.into(),
}
}
}
impl std::fmt::Display for StepError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl std::error::Error for StepError {}
impl From<crate::Error> for StepError {
fn from(err: crate::Error) -> Self {
let permanent = err.is_permanent();
let message = err.to_string();
if permanent {
Self::permanent(message)
} else {
Self::transient(message)
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum StepErrorKind {
Transient,
Permanent,
}
pub trait StepRunner: Send + Sync {
fn run_step(
&self,
step: &Step,
) -> impl Future<Output = std::result::Result<StepOutcome, StepError>> + Send;
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_util::rid;
#[test]
fn from_workflow_error_maps_via_is_permanent() {
let permanent: StepError = crate::Error::InputMismatch(rid("run-1")).into();
assert_eq!(permanent.kind, StepErrorKind::Permanent);
let store_err = taquba::object_store::Error::NotFound {
path: "x".into(),
source: "missing".into(),
};
let transient: StepError = crate::Error::Store(store_err).into();
assert_eq!(transient.kind, StepErrorKind::Transient);
}
}