use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, Default)]
pub enum QueueOrdering {
#[default]
Fifo,
Fastest,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct QueueConfig {
pub ordering: QueueOrdering,
}
impl Default for QueueConfig {
fn default() -> Self {
Self {
ordering: QueueOrdering::Fifo,
}
}
}
impl QueueConfig {
pub fn new(ordering: QueueOrdering) -> Self {
Self { ordering }
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Queue {
pub name: String,
pub config: QueueConfig,
}
impl Queue {
pub fn new(name: impl Into<String>, config: QueueConfig) -> Self {
Self {
name: name.into(),
config,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct Worker {
pub id: String,
}
impl Worker {
pub fn new(id: impl Into<String>) -> Self {
Self { id: id.into() }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum OrderingCapability {
None,
FifoLeasing,
FifoAndFastestLeasing,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum BackpressureCapability {
BacklogOnly,
ExecutionRateLimit,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum DurabilityCapability {
ProcessLocal,
Persistent,
ConfigurationDependent,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum TransactionalEnqueueCapability {
BackendOperationOnly,
SameDatabase,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum NotificationCapability {
ProcessLocalHint,
BestEffortHint,
BestEffortHintWithPolling,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum RetentionCapability {
ProcessLifetime,
ExplicitPruning,
BackendConfigured,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum ConsumerGroupCapability {
OffsetsOnly,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct BackendSemanticCapabilities {
pub durability: DurabilityCapability,
pub transactional_enqueue_scope: TransactionalEnqueueCapability,
pub notification_delivery: NotificationCapability,
pub job_retention: RetentionCapability,
pub stream_retention: RetentionCapability,
pub consumer_group_coordination: ConsumerGroupCapability,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct BackendCapabilities {
pub transactional_enqueue: bool,
pub durable_jobs: bool,
pub notifications: bool,
pub streams: bool,
pub consumer_groups: bool,
pub distributed_workers: bool,
pub ordering: OrderingCapability,
pub backpressure: BackpressureCapability,
}
impl BackendCapabilities {
pub const fn memory() -> Self {
Self {
transactional_enqueue: false,
durable_jobs: false,
notifications: true,
streams: true,
consumer_groups: true,
distributed_workers: false,
ordering: OrderingCapability::FifoAndFastestLeasing,
backpressure: BackpressureCapability::BacklogOnly,
}
}
pub const fn sqlite() -> Self {
Self {
transactional_enqueue: true,
durable_jobs: true,
notifications: true,
streams: true,
consumer_groups: true,
distributed_workers: false,
ordering: OrderingCapability::FifoAndFastestLeasing,
backpressure: BackpressureCapability::BacklogOnly,
}
}
pub const fn postgres() -> Self {
Self {
transactional_enqueue: true,
durable_jobs: true,
notifications: true,
streams: true,
consumer_groups: true,
distributed_workers: true,
ordering: OrderingCapability::FifoAndFastestLeasing,
backpressure: BackpressureCapability::ExecutionRateLimit,
}
}
pub const fn redis() -> Self {
Self {
transactional_enqueue: false,
durable_jobs: true,
notifications: true,
streams: true,
consumer_groups: true,
distributed_workers: true,
ordering: OrderingCapability::FifoLeasing,
backpressure: BackpressureCapability::BacklogOnly,
}
}
pub fn supports_portable_job_api(&self) -> bool {
self.durable_jobs || !self.distributed_workers
}
pub const fn semantics(&self) -> Option<BackendSemanticCapabilities> {
match (
self.transactional_enqueue,
self.durable_jobs,
self.notifications,
self.streams,
self.consumer_groups,
self.distributed_workers,
self.ordering,
self.backpressure,
) {
(
false,
false,
true,
true,
true,
false,
OrderingCapability::FifoAndFastestLeasing,
BackpressureCapability::BacklogOnly,
) => Some(BackendSemanticCapabilities::memory()),
(
true,
true,
true,
true,
true,
false,
OrderingCapability::FifoAndFastestLeasing,
BackpressureCapability::BacklogOnly,
) => Some(BackendSemanticCapabilities::sqlite()),
(
true,
true,
true,
true,
true,
true,
OrderingCapability::FifoAndFastestLeasing,
BackpressureCapability::ExecutionRateLimit,
) => Some(BackendSemanticCapabilities::postgres()),
(
false,
true,
true,
true,
true,
true,
OrderingCapability::FifoLeasing,
BackpressureCapability::BacklogOnly,
) => Some(BackendSemanticCapabilities::redis()),
_ => None,
}
}
}
impl BackendSemanticCapabilities {
pub const fn memory() -> Self {
Self {
durability: DurabilityCapability::ProcessLocal,
transactional_enqueue_scope: TransactionalEnqueueCapability::BackendOperationOnly,
notification_delivery: NotificationCapability::ProcessLocalHint,
job_retention: RetentionCapability::ProcessLifetime,
stream_retention: RetentionCapability::ProcessLifetime,
consumer_group_coordination: ConsumerGroupCapability::OffsetsOnly,
}
}
pub const fn sqlite() -> Self {
Self {
durability: DurabilityCapability::Persistent,
transactional_enqueue_scope: TransactionalEnqueueCapability::SameDatabase,
notification_delivery: NotificationCapability::BestEffortHintWithPolling,
job_retention: RetentionCapability::ExplicitPruning,
stream_retention: RetentionCapability::ExplicitPruning,
consumer_group_coordination: ConsumerGroupCapability::OffsetsOnly,
}
}
pub const fn postgres() -> Self {
Self {
durability: DurabilityCapability::Persistent,
transactional_enqueue_scope: TransactionalEnqueueCapability::SameDatabase,
notification_delivery: NotificationCapability::BestEffortHint,
job_retention: RetentionCapability::ExplicitPruning,
stream_retention: RetentionCapability::ExplicitPruning,
consumer_group_coordination: ConsumerGroupCapability::OffsetsOnly,
}
}
pub const fn redis() -> Self {
Self {
durability: DurabilityCapability::ConfigurationDependent,
transactional_enqueue_scope: TransactionalEnqueueCapability::BackendOperationOnly,
notification_delivery: NotificationCapability::BestEffortHintWithPolling,
job_retention: RetentionCapability::BackendConfigured,
stream_retention: RetentionCapability::BackendConfigured,
consumer_group_coordination: ConsumerGroupCapability::OffsetsOnly,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "sqlx", derive(sqlx::FromRow))]
pub struct JobListItem {
pub id: Uuid,
pub idempotency_key: Option<String>,
pub queue: String,
pub job_type: String,
pub status: String,
pub run_at: DateTime<Utc>,
#[serde(default)]
pub deadline_at: Option<DateTime<Utc>>,
#[serde(default)]
pub timeout_seconds: Option<i64>,
#[serde(default)]
pub recurring_interval_seconds: Option<i64>,
pub priority: i32,
pub max_attempts: i32,
pub last_error_code: Option<String>,
pub last_error_message: Option<String>,
pub dlq_reason_code: Option<String>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[cfg_attr(feature = "sqlx", derive(sqlx::FromRow))]
pub struct Job {
pub dataset_id: String,
pub replay_of_job_id: Option<Uuid>,
pub idempotency_key: Option<String>,
pub id: Uuid,
pub queue: String,
pub job_type: String,
#[cfg_attr(feature = "sqlx", sqlx(rename = "payload_json"))]
pub payload: Value,
pub run_at: DateTime<Utc>,
#[serde(default)]
pub deadline_at: Option<DateTime<Utc>>,
#[serde(default)]
pub timeout_seconds: Option<i64>,
#[serde(default)]
pub recurring_interval_seconds: Option<i64>,
pub status: String,
pub priority: i32,
pub max_attempts: i32,
pub locked_at: Option<DateTime<Utc>>,
pub locked_by: Option<String>,
pub lock_expires_at: Option<DateTime<Utc>>,
pub dlq_reason_code: Option<String>,
pub dlq_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
pub updated_at: DateTime<Utc>,
}
impl Job {
pub fn new(job_type: impl Into<String>, payload: Value) -> Self {
let now = Utc::now();
Self {
dataset_id: "default".to_string(),
replay_of_job_id: None,
idempotency_key: None,
id: Uuid::new_v4(),
queue: "default".to_string(),
job_type: job_type.into(),
payload,
run_at: now,
deadline_at: None,
timeout_seconds: None,
recurring_interval_seconds: None,
status: JobStatus::Queued.as_str().to_string(),
priority: 0,
max_attempts: 25,
locked_at: None,
locked_by: None,
lock_expires_at: None,
dlq_reason_code: None,
dlq_at: None,
created_at: now,
updated_at: now,
}
}
pub fn queue(mut self, queue: impl Into<String>) -> Self {
self.queue = queue.into();
self
}
pub fn priority(mut self, priority: i32) -> Self {
self.priority = priority;
self
}
pub fn max_attempts(mut self, max_attempts: i32) -> Self {
self.max_attempts = max_attempts;
self
}
pub fn idempotency_key(mut self, idempotency_key: impl Into<String>) -> Self {
self.idempotency_key = Some(idempotency_key.into());
self
}
pub fn run_at(mut self, run_at: DateTime<Utc>) -> Self {
self.run_at = run_at;
self
}
pub fn deadline_at(mut self, deadline_at: DateTime<Utc>) -> Self {
self.deadline_at = Some(deadline_at);
self
}
pub fn timeout_seconds(mut self, timeout_seconds: i64) -> Self {
self.timeout_seconds = Some(timeout_seconds.max(0));
self
}
pub fn recurring_interval_seconds(mut self, interval_seconds: i64) -> Self {
self.recurring_interval_seconds = Some(interval_seconds.max(1));
self
}
pub fn payload_json(&self) -> &Value {
&self.payload
}
pub fn lifecycle_state_at(
&self,
now: DateTime<Utc>,
failed_attempts: usize,
) -> Result<JobLifecycleState, crate::error::Error> {
JobLifecycleState::from_persisted(
JobStatus::parse(&self.status)?,
self.run_at,
now,
failed_attempts,
)
}
pub fn payload_typed<T: serde::de::DeserializeOwned>(&self) -> Result<T, crate::error::Error> {
serde_json::from_value(self.payload.clone())
.map_err(crate::error::Error::PayloadDeserialization)
}
}
#[async_trait::async_trait]
pub trait JobProcessor: Send + Sync {
async fn process(&self, job: Job) -> anyhow::Result<()>;
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct NewJob {
pub queue: String,
pub job_type: String,
pub payload_json: Value,
pub idempotency_key: Option<String>,
pub run_at: DateTime<Utc>,
#[serde(default)]
pub deadline_at: Option<DateTime<Utc>>,
#[serde(default)]
pub timeout_seconds: Option<i64>,
#[serde(default)]
pub recurring_interval_seconds: Option<i64>,
pub priority: i32,
pub max_attempts: i32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobExecution {
pub job_id: Uuid,
pub attempt_id: Uuid,
pub attempt_no: i32,
pub worker_id: String,
pub lease_expires_at: DateTime<Utc>,
pub started_at: DateTime<Utc>,
}
impl From<Job> for NewJob {
fn from(job: Job) -> Self {
NewJob {
queue: job.queue,
job_type: job.job_type,
payload_json: job.payload,
idempotency_key: job.idempotency_key,
run_at: job.run_at,
deadline_at: job.deadline_at,
timeout_seconds: job.timeout_seconds,
recurring_interval_seconds: job.recurring_interval_seconds,
priority: job.priority,
max_attempts: job.max_attempts,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum JobStatus {
Queued,
Running,
Completed,
Succeeded,
Failed,
Dlq,
Cancelled,
Canceled,
}
impl JobStatus {
pub fn as_str(&self) -> &'static str {
match self {
JobStatus::Queued => "queued",
JobStatus::Running => "running",
JobStatus::Completed | JobStatus::Succeeded => "succeeded",
JobStatus::Failed => "failed",
JobStatus::Dlq => "dlq",
JobStatus::Cancelled | JobStatus::Canceled => "canceled",
}
}
pub fn parse(status: &str) -> Result<Self, crate::error::Error> {
match status {
"queued" => Ok(JobStatus::Queued),
"running" => Ok(JobStatus::Running),
"succeeded" | "completed" => Ok(JobStatus::Completed),
"failed" => Ok(JobStatus::Failed),
"dlq" => Ok(JobStatus::Dlq),
"canceled" | "cancelled" => Ok(JobStatus::Cancelled),
other => Err(crate::error::Error::InvalidState(format!(
"unknown job status '{other}'"
))),
}
}
pub fn is_terminal(&self) -> bool {
matches!(
self,
JobStatus::Completed
| JobStatus::Succeeded
| JobStatus::Dlq
| JobStatus::Cancelled
| JobStatus::Canceled
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub enum JobLifecycleState {
Scheduled,
Queued,
Running,
Completed,
RetryWait,
Cancelled,
Dlq,
}
impl JobLifecycleState {
pub fn as_str(&self) -> &'static str {
match self {
JobLifecycleState::Scheduled => "scheduled",
JobLifecycleState::Queued => "queued",
JobLifecycleState::Running => "running",
JobLifecycleState::Completed => "completed",
JobLifecycleState::RetryWait => "retry_wait",
JobLifecycleState::Cancelled => "cancelled",
JobLifecycleState::Dlq => "dlq",
}
}
pub fn is_terminal(&self) -> bool {
matches!(
self,
JobLifecycleState::Completed | JobLifecycleState::Cancelled | JobLifecycleState::Dlq
)
}
pub fn legal_successors(&self) -> &'static [JobLifecycleState] {
use JobLifecycleState::*;
match self {
Scheduled => &[Queued],
Queued => &[Running],
Running => &[Completed, RetryWait, Cancelled, Dlq],
RetryWait => &[Queued],
Completed | Cancelled | Dlq => &[],
}
}
pub fn can_transition_to(&self, next: JobLifecycleState) -> bool {
self.legal_successors().contains(&next)
}
pub fn ensure_transition_to(&self, next: JobLifecycleState) -> Result<(), crate::error::Error> {
if self.can_transition_to(next) {
Ok(())
} else {
Err(crate::error::Error::InvalidState(format!(
"illegal job state transition: {} -> {}",
self.as_str(),
next.as_str()
)))
}
}
pub fn from_persisted(
status: JobStatus,
run_at: DateTime<Utc>,
now: DateTime<Utc>,
failed_attempts: usize,
) -> Result<Self, crate::error::Error> {
match status {
JobStatus::Queued if run_at > now && failed_attempts > 0 => {
Ok(JobLifecycleState::RetryWait)
}
JobStatus::Queued if run_at > now => Ok(JobLifecycleState::Scheduled),
JobStatus::Queued => Ok(JobLifecycleState::Queued),
JobStatus::Running => Ok(JobLifecycleState::Running),
JobStatus::Completed | JobStatus::Succeeded => Ok(JobLifecycleState::Completed),
JobStatus::Dlq => Ok(JobLifecycleState::Dlq),
JobStatus::Cancelled | JobStatus::Canceled => Ok(JobLifecycleState::Cancelled),
JobStatus::Failed => Err(crate::error::Error::InvalidState(
"job status 'failed' is legacy; failures belong to JobAttempt".to_string(),
)),
}
}
}
pub type JobHandler = std::sync::Arc<
dyn Fn(Job) -> std::pin::Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send>>
+ Send
+ Sync,
>;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "sqlx", derive(sqlx::FromRow))]
pub struct Event {
pub sequence_no: i64,
pub stream_name: String,
pub event_type: String,
pub payload_json: serde_json::Value,
pub created_at: DateTime<Utc>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct NewEvent {
pub event_type: String,
pub payload_json: serde_json::Value,
}
impl NewEvent {
pub fn new(event_type: impl Into<String>, payload_json: serde_json::Value) -> Self {
Self {
event_type: event_type.into(),
payload_json,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "sqlx", derive(sqlx::FromRow))]
pub struct ConsumerGroupStatus {
pub consumer_group: String,
pub stream_name: String,
pub last_acked_seq: i64,
pub updated_at: DateTime<Utc>,
}