use std::collections::HashMap;
use std::fmt::Display;
use std::time::Duration;
use serde::de::DeserializeOwned;
use serde::{Deserialize, Serialize};
use taquba::QueueView;
use tracing::warn;
use crate::error::Result;
pub(crate) fn encode<T: Serialize>(record: &T) -> Vec<u8> {
rmp_serde::to_vec_named(record).expect("a durable record encodes")
}
pub(crate) fn decode<T: DeserializeOwned>(bytes: &[u8]) -> Result<T> {
Ok(rmp_serde::from_slice(bytes)?)
}
pub(crate) fn decode_or_absent<T: DeserializeOwned>(
bytes: &[u8],
record: &'static str,
id: &dyn Display,
) -> Option<T> {
match rmp_serde::from_slice(bytes) {
Ok(value) => Some(value),
Err(err) => {
warn!(record, %id, error = %err, "{record} failed to decode; treated as absent");
None
}
}
}
pub(crate) async fn kv_record<T: DeserializeOwned>(
view: &QueueView,
key: &[u8],
) -> Result<Option<T>> {
match view.kv_get(key).await? {
Some(bytes) => decode(&bytes).map(Some),
None => Ok(None),
}
}
use crate::effects::StagedEffects;
use crate::keys::RunId;
use crate::runner::{StepErrorKind, StepOutcome, Trigger};
use crate::terminal::{RunOutcome, TerminalStatus};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableRunRecord {
pub(crate) run_id: RunId,
pub(crate) submitted_at_ms: u64,
pub(crate) input_hash: [u8; 32],
pub(crate) cancel_requested: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableCurrentStep {
pub(crate) step_number: u32,
pub(crate) job_id: String,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
pub(crate) struct DurableDuration {
secs: u64,
nanos: u32,
}
impl From<Duration> for DurableDuration {
fn from(duration: Duration) -> Self {
Self {
secs: duration.as_secs(),
nanos: duration.subsec_nanos(),
}
}
}
impl From<DurableDuration> for Duration {
fn from(duration: DurableDuration) -> Self {
Duration::new(duration.secs, duration.nanos)
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) enum DurableTrigger {
Immediate,
After(DurableDuration),
OnSignal {
correlation_key: String,
timeout: DurableDuration,
},
}
impl From<&Trigger> for DurableTrigger {
fn from(trigger: &Trigger) -> Self {
match trigger {
Trigger::Immediate => Self::Immediate,
Trigger::After(delay) => Self::After((*delay).into()),
Trigger::OnSignal {
correlation_key,
timeout,
} => Self::OnSignal {
correlation_key: correlation_key.clone(),
timeout: (*timeout).into(),
},
}
}
}
impl From<DurableTrigger> for Trigger {
fn from(trigger: DurableTrigger) -> Self {
match trigger {
DurableTrigger::Immediate => Self::Immediate,
DurableTrigger::After(delay) => Self::After(delay.into()),
DurableTrigger::OnSignal {
correlation_key,
timeout,
} => Self::OnSignal {
correlation_key,
timeout: timeout.into(),
},
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) enum DurableStepOutcome {
Continue {
#[serde(with = "serde_bytes")]
payload: Vec<u8>,
when: DurableTrigger,
},
Succeed {
#[serde(with = "serde_bytes")]
result: Vec<u8>,
},
Fail {
reason: String,
},
Cancel {
reason: String,
},
}
impl From<&StepOutcome> for DurableStepOutcome {
fn from(outcome: &StepOutcome) -> Self {
match outcome {
StepOutcome::Continue { payload, when } => Self::Continue {
payload: payload.clone(),
when: when.into(),
},
StepOutcome::Succeed { result } => Self::Succeed {
result: result.clone(),
},
StepOutcome::Fail { reason } => Self::Fail {
reason: reason.clone(),
},
StepOutcome::Cancel { reason } => Self::Cancel {
reason: reason.clone(),
},
}
}
}
impl From<DurableStepOutcome> for StepOutcome {
fn from(outcome: DurableStepOutcome) -> Self {
match outcome {
DurableStepOutcome::Continue { payload, when } => Self::Continue {
payload,
when: when.into(),
},
DurableStepOutcome::Succeed { result } => Self::Succeed { result },
DurableStepOutcome::Fail { reason } => Self::Fail { reason },
DurableStepOutcome::Cancel { reason } => Self::Cancel { reason },
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableStepOutcomeRecord {
pub(crate) stored_at_ms: u64,
pub(crate) outcome: DurableStepOutcome,
pub(crate) effects: StagedEffects,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) enum DurableTerminalStatus {
Succeeded,
Failed,
Cancelled,
}
impl From<TerminalStatus> for DurableTerminalStatus {
fn from(status: TerminalStatus) -> Self {
match status {
TerminalStatus::Succeeded => Self::Succeeded,
TerminalStatus::Failed => Self::Failed,
TerminalStatus::Cancelled => Self::Cancelled,
}
}
}
impl From<DurableTerminalStatus> for TerminalStatus {
fn from(status: DurableTerminalStatus) -> Self {
match status {
DurableTerminalStatus::Succeeded => Self::Succeeded,
DurableTerminalStatus::Failed => Self::Failed,
DurableTerminalStatus::Cancelled => Self::Cancelled,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) struct DurableTermination {
pub(crate) status: DurableTerminalStatus,
pub(crate) error: Option<String>,
pub(crate) error_kind: Option<DurableErrorKind>,
pub(crate) final_step: u32,
pub(crate) terminated_at_ms: u64,
pub(crate) input_hash: [u8; 32],
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableMember {
pub(crate) run_id: RunId,
pub(crate) terminated: Option<DurableTermination>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableRunOutcome {
run_id: RunId,
status: DurableTerminalStatus,
#[serde(with = "serde_bytes")]
result: Option<Vec<u8>>,
error: Option<String>,
headers: HashMap<String, String>,
final_step: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub(crate) enum DurableErrorKind {
Transient,
Permanent,
}
impl From<StepErrorKind> for DurableErrorKind {
fn from(kind: StepErrorKind) -> Self {
match kind {
StepErrorKind::Transient => Self::Transient,
StepErrorKind::Permanent => Self::Permanent,
}
}
}
impl From<DurableErrorKind> for StepErrorKind {
fn from(kind: DurableErrorKind) -> Self {
match kind {
DurableErrorKind::Transient => Self::Transient,
DurableErrorKind::Permanent => Self::Permanent,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub(crate) struct DurableRunResult {
pub(crate) termination: DurableTermination,
pub(crate) outcome: DurableRunOutcome,
}
impl From<&RunOutcome> for DurableRunOutcome {
fn from(outcome: &RunOutcome) -> Self {
Self {
run_id: outcome.run_id.clone(),
status: outcome.status.into(),
result: outcome.result.clone(),
error: outcome.error.clone(),
headers: outcome.headers.clone(),
final_step: outcome.final_step,
}
}
}
impl From<DurableRunOutcome> for RunOutcome {
fn from(outcome: DurableRunOutcome) -> Self {
Self {
run_id: outcome.run_id,
status: outcome.status.into(),
result: outcome.result,
error: outcome.error,
headers: outcome.headers,
final_step: outcome.final_step,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_util::rid;
#[test]
fn payload_bytes_are_stored_as_binary_strings() {
let payload: Vec<u8> = (0..=255).collect();
let is_contiguous = |bytes: &[u8]| bytes.windows(payload.len()).any(|w| w == payload);
assert!(is_contiguous(&encode(&DurableStepOutcome::Continue {
payload: payload.clone(),
when: DurableTrigger::Immediate,
})));
assert!(is_contiguous(&encode(&DurableStepOutcome::Succeed {
result: payload.clone(),
})));
let outcome = DurableRunOutcome {
run_id: rid("run"),
status: DurableTerminalStatus::Succeeded,
result: Some(payload.clone()),
error: None,
headers: HashMap::new(),
final_step: 0,
};
let stored = encode(&outcome);
assert!(is_contiguous(&stored));
let decoded: DurableRunOutcome = decode(&stored).unwrap();
assert_eq!(decoded.result, Some(payload));
}
}