use std::{fmt, str::FromStr, time::Duration};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value as JsonValue};
use uuid::Uuid;
pub type Json = JsonValue;
pub type JsonObject = Map<String, Json>;
#[derive(
Debug, Clone, Copy, Serialize, Deserialize, sqlx::Type, PartialEq, Eq, Hash, PartialOrd, Ord,
)]
#[serde(transparent)]
#[sqlx(transparent)]
pub struct TaskId(Uuid);
impl TaskId {
pub const fn from_uuid(value: Uuid) -> Self {
Self(value)
}
pub const fn into_uuid(self) -> Uuid {
self.0
}
}
impl From<Uuid> for TaskId {
fn from(value: Uuid) -> Self {
Self::from_uuid(value)
}
}
impl From<TaskId> for Uuid {
fn from(value: TaskId) -> Self {
value.into_uuid()
}
}
impl fmt::Display for TaskId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl FromStr for TaskId {
type Err = uuid::Error;
fn from_str(value: &str) -> Result<Self, Self::Err> {
value.parse().map(Self::from_uuid)
}
}
#[derive(
Debug, Clone, Copy, Serialize, Deserialize, sqlx::Type, PartialEq, Eq, Hash, PartialOrd, Ord,
)]
#[serde(transparent)]
#[sqlx(transparent)]
pub struct RunId(Uuid);
impl RunId {
pub const fn from_uuid(value: Uuid) -> Self {
Self(value)
}
pub const fn into_uuid(self) -> Uuid {
self.0
}
}
impl From<Uuid> for RunId {
fn from(value: Uuid) -> Self {
Self::from_uuid(value)
}
}
impl From<RunId> for Uuid {
fn from(value: RunId) -> Self {
value.into_uuid()
}
}
impl fmt::Display for RunId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
self.0.fmt(f)
}
}
impl FromStr for RunId {
type Err = uuid::Error;
fn from_str(value: &str) -> Result<Self, Self::Err> {
value.parse().map(Self::from_uuid)
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum RetryStrategy {
Fixed {
delay: Duration,
},
Exponential {
initial_delay: Duration,
factor: f64,
max_delay: Option<Duration>,
},
None,
}
impl RetryStrategy {
pub const fn fixed(delay: Duration) -> Self {
Self::Fixed { delay }
}
pub const fn exponential(
initial_delay: Duration,
factor: f64,
max_delay: Option<Duration>,
) -> Self {
Self::Exponential { initial_delay, factor, max_delay }
}
pub const fn none() -> Self {
Self::None
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct CancellationPolicy {
pub(crate) max_duration: Option<Duration>,
pub(crate) max_delay: Option<Duration>,
}
impl CancellationPolicy {
pub const fn new() -> Self {
Self { max_duration: None, max_delay: None }
}
#[must_use]
pub const fn max_duration(mut self, duration: Duration) -> Self {
self.max_duration = Some(duration);
self
}
#[must_use]
pub const fn max_delay(mut self, delay: Duration) -> Self {
self.max_delay = Some(delay);
self
}
}
#[derive(Debug, Clone, Default)]
pub(crate) struct SpawnConfig {
pub max_attempts: Option<u32>,
pub retry_strategy: Option<RetryStrategy>,
pub headers: Option<JsonObject>,
pub cancellation: Option<CancellationPolicy>,
pub idempotency_key: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct SpawnResult {
pub task_id: TaskId,
pub created: bool,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct QueuePolicyOptions {
pub(crate) cleanup_ttl: Option<Duration>,
pub(crate) cleanup_limit: Option<u32>,
}
impl QueuePolicyOptions {
pub const fn new() -> Self {
Self { cleanup_ttl: None, cleanup_limit: None }
}
#[must_use]
pub const fn cleanup_ttl(mut self, cleanup_ttl: Duration) -> Self {
self.cleanup_ttl = Some(cleanup_ttl);
self
}
#[must_use]
pub const fn cleanup_limit(mut self, cleanup_limit: u32) -> Self {
self.cleanup_limit = Some(cleanup_limit);
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueuePolicy {
pub queue_name: String,
pub cleanup_ttl: Duration,
pub cleanup_limit: u32,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueueCleanup {
pub queue_name: String,
pub tasks_deleted: u32,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum TaskState {
Pending,
Running,
Sleeping,
Completed,
Failed,
Cancelled,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum TaskResultSnapshot {
Pending,
Running,
Sleeping,
Completed {
result: Json,
},
Failed {
failure: Json,
},
Cancelled,
}
impl TaskResultSnapshot {
pub(crate) const fn is_terminal(&self) -> bool {
matches!(self, Self::Completed { .. } | Self::Failed { .. } | Self::Cancelled)
}
}
#[derive(Debug, Clone)]
pub(crate) struct ClaimedTask {
pub(crate) run_id: RunId,
pub(crate) task_id: TaskId,
pub(crate) task_name: String,
pub(crate) attempt: u32,
pub(crate) params: Json,
pub(crate) headers: Option<JsonObject>,
}