use std::future::Future;
use bytes::Bytes;
use crate::cipher::EntryKindTag;
use crate::config::RetentionPolicy;
use crate::effect::EffectClass;
use crate::error::DurableError;
use crate::ids::{
ExecutionId, ExecutionKind, IdempotencyKey, JournalSeq, PromiseId, StepId, TimerId,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionStatus {
Running,
Completed,
Failed,
Aborted,
Canceled,
}
impl ExecutionStatus {
#[must_use]
pub fn as_str(self) -> &'static str {
match self {
Self::Running => "running",
Self::Completed => "completed",
Self::Failed => "failed",
Self::Aborted => "aborted",
Self::Canceled => "canceled",
}
}
#[must_use]
pub fn from_tag(tag: &str) -> Option<Self> {
match tag {
"running" => Some(Self::Running),
"completed" => Some(Self::Completed),
"failed" => Some(Self::Failed),
"aborted" => Some(Self::Aborted),
"canceled" => Some(Self::Canceled),
_ => None,
}
}
#[must_use]
pub fn is_running(self) -> bool {
matches!(self, Self::Running)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum EntryKind {
StepResult {
idempotency_key: IdempotencyKey,
payload: Bytes,
effect: EffectClass,
payload_version: u8,
},
EffectIntent {
idempotency_key: IdempotencyKey,
effect: EffectClass,
hmac: Option<[u8; 32]>,
},
PromiseCreated {
promise_id: PromiseId,
resolver_token_hash: [u8; 32],
hmac: Option<[u8; 32]>,
},
PromiseResolved {
promise_id: PromiseId,
payload: Bytes,
},
TimerArmed {
timer_id: TimerId,
due_at_ms: i64,
hmac: Option<[u8; 32]>,
},
TimerFired {
timer_id: TimerId,
},
Checkpoint {
up_to_step: u32,
snapshot: Bytes,
},
}
impl EntryKind {
#[must_use]
pub fn tag_enum(&self) -> EntryKindTag {
match self {
Self::StepResult { .. } => EntryKindTag::StepResult,
Self::EffectIntent { .. } => EntryKindTag::EffectIntent,
Self::PromiseCreated { .. } => EntryKindTag::PromiseCreated,
Self::PromiseResolved { .. } => EntryKindTag::PromiseResolved,
Self::TimerArmed { .. } => EntryKindTag::TimerArmed,
Self::TimerFired { .. } => EntryKindTag::TimerFired,
Self::Checkpoint { .. } => EntryKindTag::Checkpoint,
}
}
#[must_use]
pub fn tag(&self) -> &'static str {
self.tag_enum().as_str()
}
#[must_use]
pub fn idempotency_key(&self) -> Option<IdempotencyKey> {
match self {
Self::StepResult {
idempotency_key, ..
}
| Self::EffectIntent {
idempotency_key, ..
} => Some(*idempotency_key),
Self::PromiseCreated { .. }
| Self::PromiseResolved { .. }
| Self::TimerArmed { .. }
| Self::TimerFired { .. }
| Self::Checkpoint { .. } => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct JournalEntry {
pub seq: Option<JournalSeq>,
pub execution_id: ExecutionId,
pub kind: ExecutionKind,
pub step_id: StepId,
pub entry: EntryKind,
pub created_at_ms: i64,
}
pub trait Journal: Send + Sync {
fn append(
&self,
entry: JournalEntry,
) -> impl Future<Output = Result<JournalSeq, DurableError>> + Send;
fn read_execution(
&self,
id: ExecutionId,
) -> impl Future<Output = Result<Vec<JournalEntry>, DurableError>> + Send;
fn read_execution_range(
&self,
id: ExecutionId,
from_step_id: u32,
limit: usize,
) -> impl Future<Output = Result<Vec<JournalEntry>, DurableError>> + Send;
fn finalize(
&self,
id: ExecutionId,
status: ExecutionStatus,
) -> impl Future<Output = Result<(), DurableError>> + Send;
fn prune(
&self,
policy: &RetentionPolicy,
) -> impl Future<Output = Result<u64, DurableError>> + Send;
fn sweep_orphans(
&self,
policy: &RetentionPolicy,
) -> impl Future<Output = Result<u64, DurableError>> + Send;
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ids::ExecutionId;
fn sample_entry(entry: EntryKind) -> JournalEntry {
JournalEntry {
seq: None,
execution_id: ExecutionId::new(),
kind: ExecutionKind::AgentTurn,
step_id: StepId::new(0),
entry,
created_at_ms: 0,
}
}
#[test]
fn entry_kind_match_is_exhaustive() {
let key = IdempotencyKey::derive(ExecutionId::new(), StepId::new(0), b"op");
for entry in [
EntryKind::StepResult {
idempotency_key: key,
payload: Bytes::from_static(b"x"),
effect: EffectClass::Idempotent,
payload_version: 1,
},
EntryKind::EffectIntent {
idempotency_key: key,
effect: EffectClass::ExactlyOnceGuarded,
hmac: None,
},
EntryKind::PromiseCreated {
promise_id: PromiseId::new(),
resolver_token_hash: [0u8; 32],
hmac: Some([1u8; 32]),
},
EntryKind::PromiseResolved {
promise_id: PromiseId::new(),
payload: Bytes::new(),
},
EntryKind::TimerArmed {
timer_id: TimerId::new(),
due_at_ms: 100,
hmac: None,
},
EntryKind::TimerFired {
timer_id: TimerId::new(),
},
EntryKind::Checkpoint {
up_to_step: 3,
snapshot: Bytes::new(),
},
] {
let tag = match &entry {
EntryKind::StepResult { .. } => "step_result",
EntryKind::EffectIntent { .. } => "effect_intent",
EntryKind::PromiseCreated { .. } => "promise_created",
EntryKind::PromiseResolved { .. } => "promise_resolved",
EntryKind::TimerArmed { .. } => "timer_armed",
EntryKind::TimerFired { .. } => "timer_fired",
EntryKind::Checkpoint { .. } => "checkpoint",
};
assert_eq!(tag, entry.tag());
}
}
#[test]
fn journal_entry_is_clonable_and_comparable() {
let entry = sample_entry(EntryKind::TimerFired {
timer_id: TimerId::new(),
});
assert_eq!(entry, entry.clone());
}
#[test]
fn execution_status_round_trips_through_str() {
for status in [
ExecutionStatus::Running,
ExecutionStatus::Completed,
ExecutionStatus::Failed,
ExecutionStatus::Aborted,
ExecutionStatus::Canceled,
] {
assert!(!status.as_str().is_empty());
assert_eq!(ExecutionStatus::from_tag(status.as_str()), Some(status));
}
assert!(ExecutionStatus::Running.is_running());
assert!(!ExecutionStatus::Aborted.is_running());
assert!(!ExecutionStatus::Canceled.is_running());
}
}