From f04e06e14110e83e46059a0f13e06a319eae7b1e Mon Sep 17 00:00:00 2001 From: RoyLin Date: Sun, 9 Aug 2026 08:55:51 +0800 Subject: [PATCH 1/2] feat(queue): add ORM-backed PostgreSQL backend --- Cargo.toml | 3 + README.md | 60 +- ROADMAP.md | 14 +- src/lib.rs | 9 +- src/queue.rs | 216 +++---- src/queue/module.rs | 114 ++++ src/queue/postgres.rs | 452 ++++++++++++++ src/queue/postgres/deduplication.rs | 157 +++++ src/queue/postgres/lifecycle.rs | 306 +++++++++ src/queue/postgres/migrations.rs | 73 +++ src/queue/postgres/store.rs | 922 ++++++++++++++++++++++++++++ tests/postgres_queue.rs | 454 ++++++++++++++ tests/postgres_queue_contract.rs | 166 +++++ tests/postgres_queue_shared.rs | 102 +++ tests/queue.rs | 58 +- 15 files changed, 2970 insertions(+), 136 deletions(-) create mode 100644 src/queue/module.rs create mode 100644 src/queue/postgres.rs create mode 100644 src/queue/postgres/deduplication.rs create mode 100644 src/queue/postgres/lifecycle.rs create mode 100644 src/queue/postgres/migrations.rs create mode 100644 src/queue/postgres/store.rs create mode 100644 tests/postgres_queue.rs create mode 100644 tests/postgres_queue_contract.rs create mode 100644 tests/postgres_queue_shared.rs diff --git a/Cargo.toml b/Cargo.toml index d030ce5..9824f32 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -40,6 +40,7 @@ logging = [] macros = ["dep:a3s-boot-macros"] openapi-schemas = ["dep:schemars"] queue = ["dep:a3s-lane", "dep:chrono", "dep:tokio"] +queue-postgres = ["queue", "dep:a3s-orm", "dep:uuid"] request-context = ["dep:tokio"] grpc-transport = ["dep:prost", "dep:tokio", "dep:tonic", "dep:tonic-prost"] kafka-transport = ["dep:chrono", "dep:rskafka", "dep:tokio"] @@ -59,6 +60,7 @@ a3s-acl = { version = "0.2.1", optional = true } a3s-boot-macros = { version = "0.1.2", path = "macros", optional = true } a3s-event = { version = "0.3.0", default-features = false, optional = true } a3s-lane = { version = "0.5.1", default-features = false, optional = true } +a3s-orm = { version = "0.2.0", default-features = false, features = ["postgres"], optional = true } async-trait = { version = "0.1", optional = true } async-nats = { version = "0.49.1", default-features = false, optional = true } axum = { version = "0.8", features = ["ws"], optional = true } @@ -90,6 +92,7 @@ tokio = { version = "1", features = ["fs", "io-util", "net", "rt", "sync", "time tonic = { version = "0.14.6", default-features = false, features = ["codegen", "transport"], optional = true } tonic-prost = { version = "0.14.6", optional = true } url = { version = "2", optional = true } +uuid = { version = "1", features = ["v4"], optional = true } zeroize = { version = "1", features = ["derive"], optional = true } [dev-dependencies] diff --git a/README.md b/README.md index 9c7e7f0..d094bac 100644 --- a/README.md +++ b/README.md @@ -119,7 +119,8 @@ are opt-in. | Database | `database` | Replaceable database facade and in-memory backend | | Events | `events` | A3S Event-backed emitter and listeners | | CQRS | `cqrs` | Command, query, and event buses | -| Queue | `queue` | A3S Lane-backed jobs, retries, priorities, and processors | +| Queue | `queue` | A3S Lane-backed in-process jobs, retries, priorities, and processors | +| Queue persistence | `queue-postgres` | A3S ORM-backed shared PostgreSQL leasing, recovery, fencing, and retention | | Scheduling | `schedule` | Cron, interval, and timeout jobs | | Observability | `logging`, `health` | Structured logging and health indicators | | HTTP utilities | `http-client`, `compression` | Outbound HTTP and gzip responses | @@ -133,8 +134,9 @@ are opt-in. A feature exposes the corresponding framework integration; external transports still require their broker or service to be available. Database, cache, session, -queue, and scheduler APIs are backend abstractions, and the bundled implementation -is not a claim of support for every production backend. +queue, and scheduler APIs are backend abstractions. The `queue-postgres` feature +is the durable shared queue implementation; other bundled implementations are +not a claim of support for every production backend. ## Quick Start @@ -249,6 +251,58 @@ Modules and providers can observe initialization, bootstrap, destruction, and application shutdown. Shutdown hooks can listen for SIGINT and SIGTERM when the `shutdown-hooks` feature is enabled. +### Durable PostgreSQL queues + +Enable `queue-postgres` when workers in multiple processes must share durable +work. `PostgresQueueBackend` stores jobs through A3S ORM, leases ready jobs with +PostgreSQL `SKIP LOCKED`, renews live leases, fences stale workers, and recovers +expired leases after worker or process death. + +```toml +[dependencies] +a3s-boot = { version = "0.1.3", features = ["queue-postgres"] } +serde_json = "1" +tokio = { version = "1", features = ["macros", "rt-multi-thread"] } +``` + +```rust,no_run +use std::time::Duration; + +use a3s_boot::{ + ModuleRef, PostgresQueueBackend, Queue, QueueContext, QueueJob, QueueOptions, + Result, +}; +use serde_json::json; + +async fn run_worker(database_url: &str) -> Result<()> { + let options = QueueOptions::new() + .with_worker_count(4) + .with_lease_duration(Duration::from_secs(30)); + let backend = PostgresQueueBackend::connect(database_url, "workflow", options).await?; + let queue = Queue::new("workflow", backend); + queue.process("resume", |job: QueueJob, _context: QueueContext| async move { + println!("resuming {}", job.data["runId"]); + Ok(()) + })?; + queue.start(ModuleRef::new()).await?; + queue.enqueue("resume", &json!({"runId": "run-42"})).await?; + queue.shutdown().await +} +``` + +Use a database URL whose search path selects a schema dedicated to Boot. The +host application creates that schema; Boot owns the queue tables and its A3S ORM +migration ledger inside it. Sharing the Flow or application schema can make one +component accept another component's migration history. + +The backend supports caller-assigned idempotency keys, priority and FIFO/LIFO +ordering, delay, retry, processor timeout, terminal retention, deduplication, +active keep-latest successors, and graceful lease release. Delivery is +at-least-once, so processors must make business effects idempotent. Repeat jobs +and Lane parent/child flow options are rejected explicitly. Async services can +use `jobs_async`, `failures_async`, `stats_async`, and `clear_async` on a retained +backend handle for non-blocking diagnostics. + ### Request pipeline HTTP handlers run through deterministic middleware and pipeline stages: diff --git a/ROADMAP.md b/ROADMAP.md index 3af0409..eb06671 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -249,9 +249,11 @@ Implemented today: in-process timeout/interval/cron jobs, named/global provider exports, Nest-style `#[schedule]` / `#[cron]` / `#[interval]` / `#[timeout]` macros, and lifecycle-managed shutdown. -- Provider-backed queues with `QueueModule`, `Queue`, `a3s-lane` backed job - storage and workers, typed serde JSON payloads, named/global provider - exports, and lifecycle-managed processors. +- Provider-backed queues with `QueueModule`, `Queue`, `a3s-lane` backed + in-process workers, and an optional A3S ORM-backed PostgreSQL backend with + shared leasing, fencing, process-death recovery, retry, timeout, retention, + idempotency, and deduplication. Queues retain typed serde JSON payloads, + named/global provider exports, and lifecycle-managed processors. - Provider-backed application events with `EventModule`, Nest-style `EventEmitter`, injectable `a3s-event` `EventBus`, listener macros, and pluggable providers. @@ -1008,7 +1010,11 @@ Acceptance: participate in module imports/exports. (Covered) - Queue can register typed providers, enqueue serde JSON jobs through `a3s-lane`, run named processors through lifecycle-managed workers, and - participate in module imports/exports. (Covered) + participate in module imports/exports. Its optional PostgreSQL backend uses + A3S ORM migrations and parameterized SQL, shares work across independent + workers, fences stale leases, recovers process death, and preserves typed + retry, timeout, retention, job-id idempotency, and deduplication semantics. + (Covered against PostgreSQL 17) - Application events can register an `a3s-event` backed `EventEmitter` provider, dispatch typed JSON payloads to exact or wildcard listeners, expose Nest-style listener macros, retain events through the underlying `EventBus`, diff --git a/src/lib.rs b/src/lib.rs index c7b4ff5..1823282 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -183,11 +183,14 @@ pub use provider::{ ProviderOnApplicationShutdown, ProviderOnModuleDestroy, ProviderOnModuleInit, ProviderRef, ProviderScope, ProviderToken, }; +#[cfg(feature = "queue-postgres")] +pub use queue::PostgresQueueBackend; #[cfg(feature = "queue")] pub use queue::{ - InProcessQueueBackend, Queue, QueueBackend, QueueContext, QueueJob, QueueJobFailure, - QueueJobInfo, QueueJobOptions, QueueJobPriority, QueueJobReceipt, QueueJobState, QueueModule, - QueueOptions, QueueProcessor, QueueRetryPolicy, QueueStats, + InProcessQueueBackend, Queue, QueueBackend, QueueContext, QueueDeduplicationOptions, QueueJob, + QueueJobFailure, QueueJobInfo, QueueJobOptions, QueueJobPriority, QueueJobReceipt, + QueueJobRetention, QueueJobState, QueueModule, QueueOptions, QueueProcessor, QueueRetryPolicy, + QueueStats, }; #[cfg(feature = "request-context")] pub use request_context::RequestContext; diff --git a/src/queue.rs b/src/queue.rs index 546dc44..70d1391 100644 --- a/src/queue.rs +++ b/src/queue.rs @@ -1,24 +1,35 @@ -use crate::{BootError, BoxFuture, Module, ModuleRef, ProviderDefinition, ProviderToken, Result}; +use crate::{BootError, BoxFuture, ModuleRef, Result}; +pub use a3s_lane::{ + DeduplicationOptions as QueueDeduplicationOptions, JobOptions as QueueJobOptions, + JobPriority as QueueJobPriority, JobRetention as QueueJobRetention, + RetryPolicy as QueueRetryPolicy, +}; use a3s_lane::{ InMemoryJobQueue, Job, JobListOptions, JobQueueBackend as LaneJobQueueBackend, JobQueueStats as LaneJobQueueStats, JobState as LaneJobState, LaneError, }; -pub use a3s_lane::{ - JobOptions as QueueJobOptions, JobPriority as QueueJobPriority, RetryPolicy as QueueRetryPolicy, -}; use chrono::Utc; use serde::de::DeserializeOwned; use serde::Serialize; use serde_json::Value; use std::collections::BTreeMap; use std::fmt; -use std::future::Future; +use std::future::{pending, Future}; use std::sync::{Arc, Mutex, RwLock}; use std::time::Duration; use tokio::runtime::{Builder as TokioRuntimeBuilder, Handle}; use tokio::sync::Notify; use tokio::task::JoinHandle; +#[cfg(feature = "queue-postgres")] +mod postgres; + +mod module; + +pub use module::QueueModule; +#[cfg(feature = "queue-postgres")] +pub use postgres::PostgresQueueBackend; + /// Queue runtime options shared by queue modules and Lane-backed queue backends. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct QueueOptions { @@ -705,6 +716,11 @@ async fn run_queue_worker_once( options: QueueOptions, worker_id: &str, ) -> Result<()> { + state + .backend + .recover_stalled_jobs(Utc::now()) + .await + .map_err(lane_error)?; state .backend .promote_due_jobs(Utc::now()) @@ -752,23 +768,79 @@ async fn process_claimed_job( queue_name: queue_name.to_string(), module_ref: module_ref.clone(), }; - match processor.process(queue_job, context).await { - Ok(()) => { - state - .backend - .complete_job(&job.id, &lock_token, Value::Null, Utc::now()) - .await - .map_err(lane_error)?; + let processing = processor.process(queue_job, context); + tokio::pin!(processing); + let timeout = job.options.timeout; + let timeout_future = async move { + match timeout { + Some(timeout) => tokio::time::sleep(timeout).await, + None => pending::<()>().await, } - Err(error) => { - state - .backend - .fail_job_discarding_retry(&job.id, &lock_token, error.to_string(), Utc::now()) - .await - .map_err(lane_error)?; + }; + tokio::pin!(timeout_future); + let heartbeat_every = queue_heartbeat_interval(options.lease_duration); + let mut heartbeat = tokio::time::interval_at( + tokio::time::Instant::now() + heartbeat_every, + heartbeat_every, + ); + loop { + tokio::select! { + result = &mut processing => { + match result { + Ok(()) => state + .backend + .complete_job(&job.id, &lock_token, Value::Null, Utc::now()) + .await + .map_err(lane_error)?, + Err(error) => state + .backend + .fail_job(&job.id, &lock_token, error.to_string(), Utc::now()) + .await + .map_err(lane_error)?, + }; + return Ok(()); + } + () = &mut timeout_future => { + state + .backend + .fail_job( + &job.id, + &lock_token, + format!("queue processor timed out after {timeout:?}"), + Utc::now(), + ) + .await + .map_err(lane_error)?; + return Ok(()); + } + _ = heartbeat.tick() => { + state + .backend + .renew_lease(&job.id, &lock_token, options.lease_duration, Utc::now()) + .await + .map_err(lane_error)?; + } + _ = state.notify.notified() => { + if !state.is_running() { + state + .backend + .release_active_job(&job.id, &lock_token, Utc::now()) + .await + .map_err(lane_error)?; + return Ok(()); + } + } } } - Ok(()) +} + +fn queue_heartbeat_interval(lease_duration: Duration) -> Duration { + let interval = lease_duration / 3; + if interval.is_zero() { + Duration::from_millis(1) + } else { + interval + } } async fn wait_for_processor( @@ -832,109 +904,3 @@ where .map_err(|error| BootError::Internal(format!("failed to create queue runtime: {error}")))?; runtime.block_on(future) } - -/// Module that registers and exports a [`Queue`] provider. -#[derive(Clone)] -pub struct QueueModule { - name: &'static str, - token: ProviderToken, - queue: Arc, - processors: Vec<(String, Arc)>, - global: bool, -} - -impl fmt::Debug for QueueModule { - fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { - f.debug_struct("QueueModule") - .field("name", &self.name) - .field("token", &self.token) - .field("queue", &self.queue) - .field("processors", &self.processors.len()) - .field("global", &self.global) - .finish_non_exhaustive() - } -} - -impl QueueModule { - pub fn in_process(name: &'static str) -> Self { - Self::from_queue(name, Queue::in_process(name)) - } - - pub fn in_process_with_options(name: &'static str, options: QueueOptions) -> Self { - Self::from_queue(name, Queue::in_process_with_options(name, options)) - } - - pub fn from_lane_backend_arc( - name: &'static str, - backend: Arc, - ) -> Self { - Self::from_queue(name, Queue::from_lane_backend_arc(name, backend)) - } - - pub fn from_queue(name: &'static str, queue: Queue) -> Self { - Self { - name, - token: ProviderToken::of::(), - queue: Arc::new(queue), - processors: Vec::new(), - global: false, - } - } - - pub fn processor

(mut self, name: impl Into, processor: P) -> Self - where - P: QueueProcessor, - { - self.processors.push((name.into(), Arc::new(processor))); - self - } - - pub fn named(mut self, token: impl Into) -> Self { - self.token = ProviderToken::named(token); - self - } - - pub fn global(mut self) -> Self { - self.global = true; - self - } -} - -impl Module for QueueModule { - fn name(&self) -> &'static str { - self.name - } - - fn providers(&self) -> Result> { - Ok(vec![ProviderDefinition::named_from_arc( - self.token.as_str(), - Arc::clone(&self.queue), - )]) - } - - fn exports(&self) -> Result> { - Ok(vec![self.token.clone()]) - } - - fn is_global(&self) -> bool { - self.global - } - - fn on_module_init(&self, _module_ref: &ModuleRef) -> Result<()> { - for (name, processor) in &self.processors { - self.queue - .process_arc(name.clone(), Arc::clone(processor))?; - } - Ok(()) - } - - fn on_application_bootstrap(&self, module_ref: ModuleRef) -> BoxFuture<'static, Result<()>> { - let queue = Arc::clone(&self.queue); - Box::pin(async move { queue.start(module_ref).await }) - } - - fn on_application_shutdown(&self, _module_ref: ModuleRef) -> BoxFuture<'static, Result<()>> { - let queue = Arc::clone(&self.queue); - Box::pin(async move { queue.shutdown().await }) - } -} diff --git a/src/queue/module.rs b/src/queue/module.rs new file mode 100644 index 0000000..cdb15cc --- /dev/null +++ b/src/queue/module.rs @@ -0,0 +1,114 @@ +use std::fmt; +use std::sync::Arc; + +use a3s_lane::JobQueueBackend as LaneJobQueueBackend; + +use crate::{BoxFuture, Module, ModuleRef, ProviderDefinition, ProviderToken, Result}; + +use super::{Queue, QueueOptions, QueueProcessor}; + +/// Module that registers and exports a [`Queue`] provider. +#[derive(Clone)] +pub struct QueueModule { + name: &'static str, + token: ProviderToken, + queue: Arc, + processors: Vec<(String, Arc)>, + global: bool, +} + +impl fmt::Debug for QueueModule { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_struct("QueueModule") + .field("name", &self.name) + .field("token", &self.token) + .field("queue", &self.queue) + .field("processors", &self.processors.len()) + .field("global", &self.global) + .finish_non_exhaustive() + } +} + +impl QueueModule { + pub fn in_process(name: &'static str) -> Self { + Self::from_queue(name, Queue::in_process(name)) + } + + pub fn in_process_with_options(name: &'static str, options: QueueOptions) -> Self { + Self::from_queue(name, Queue::in_process_with_options(name, options)) + } + + pub fn from_lane_backend_arc( + name: &'static str, + backend: Arc, + ) -> Self { + Self::from_queue(name, Queue::from_lane_backend_arc(name, backend)) + } + + pub fn from_queue(name: &'static str, queue: Queue) -> Self { + Self { + name, + token: ProviderToken::of::(), + queue: Arc::new(queue), + processors: Vec::new(), + global: false, + } + } + + pub fn processor

(mut self, name: impl Into, processor: P) -> Self + where + P: QueueProcessor, + { + self.processors.push((name.into(), Arc::new(processor))); + self + } + + pub fn named(mut self, token: impl Into) -> Self { + self.token = ProviderToken::named(token); + self + } + + pub fn global(mut self) -> Self { + self.global = true; + self + } +} + +impl Module for QueueModule { + fn name(&self) -> &'static str { + self.name + } + + fn providers(&self) -> Result> { + Ok(vec![ProviderDefinition::named_from_arc( + self.token.as_str(), + Arc::clone(&self.queue), + )]) + } + + fn exports(&self) -> Result> { + Ok(vec![self.token.clone()]) + } + + fn is_global(&self) -> bool { + self.global + } + + fn on_module_init(&self, _module_ref: &ModuleRef) -> Result<()> { + for (name, processor) in &self.processors { + self.queue + .process_arc(name.clone(), Arc::clone(processor))?; + } + Ok(()) + } + + fn on_application_bootstrap(&self, module_ref: ModuleRef) -> BoxFuture<'static, Result<()>> { + let queue = Arc::clone(&self.queue); + Box::pin(async move { queue.start(module_ref).await }) + } + + fn on_application_shutdown(&self, _module_ref: ModuleRef) -> BoxFuture<'static, Result<()>> { + let queue = Arc::clone(&self.queue); + Box::pin(async move { queue.shutdown().await }) + } +} diff --git a/src/queue/postgres.rs b/src/queue/postgres.rs new file mode 100644 index 0000000..65e3f80 --- /dev/null +++ b/src/queue/postgres.rs @@ -0,0 +1,452 @@ +use std::collections::BTreeMap; +use std::fmt; +use std::future::{pending, Future}; +use std::sync::{Arc, Mutex, RwLock}; +use std::time::Duration; + +use a3s_orm::PostgresExecutor; +use serde_json::Value; +use tokio::runtime::{Builder as TokioRuntimeBuilder, Handle}; +use tokio::sync::{watch, Notify}; +use tokio::task::JoinHandle; + +use crate::{BootError, BoxFuture, ModuleRef, Result}; + +use super::{ + QueueBackend, QueueContext, QueueJob, QueueJobFailure, QueueJobInfo, QueueJobOptions, + QueueJobReceipt, QueueOptions, QueueProcessor, QueueStats, +}; + +mod deduplication; +mod lifecycle; +mod migrations; +mod store; + +use store::{ClaimedJob, PostgresQueueStore}; + +const RECOVERY_BATCH_SIZE: usize = 100; + +/// A3S ORM-backed shared PostgreSQL queue backend. +/// +/// The database URL should select a schema dedicated to Boot migrations. Boot, +/// Flow, and host applications keep separate ORM migration ledgers so one +/// component cannot silently accept another component's migration history. +#[derive(Clone)] +pub struct PostgresQueueBackend { + state: Arc, + options: QueueOptions, +} + +impl fmt::Debug for PostgresQueueBackend { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + let processors = self + .state + .processors + .read() + .map(|processors| processors.len()) + .unwrap_or_default(); + let workers = self + .state + .lifecycle + .lock() + .map(|lifecycle| lifecycle.handles.len()) + .unwrap_or_default(); + formatter + .debug_struct("PostgresQueueBackend") + .field("queue_name", &self.state.store.queue_name()) + .field("options", &self.options) + .field("processors", &processors) + .field("workers", &workers) + .finish_non_exhaustive() + } +} + +impl PostgresQueueBackend { + /// Connect with an A3S ORM PostgreSQL executor and apply Boot migrations. + pub async fn connect( + database_url: impl AsRef, + queue_name: impl AsRef, + options: QueueOptions, + ) -> Result { + let store = PostgresQueueStore::connect(database_url.as_ref(), queue_name.as_ref()).await?; + Ok(Self::from_store(store, options)) + } + + /// Use a configured A3S ORM executor and apply Boot migrations. + pub async fn from_executor( + executor: PostgresExecutor, + queue_name: impl AsRef, + options: QueueOptions, + ) -> Result { + let store = PostgresQueueStore::from_executor(executor, queue_name.as_ref()).await?; + Ok(Self::from_store(store, options)) + } + + fn from_store(store: PostgresQueueStore, options: QueueOptions) -> Self { + Self { + state: Arc::new(PostgresQueueState { + store, + processors: RwLock::new(BTreeMap::new()), + lifecycle: Mutex::new(PostgresQueueLifecycle::default()), + last_worker_error: Mutex::new(None), + notify: Notify::new(), + }), + options, + } + } + + pub fn queue_name(&self) -> &str { + self.state.store.queue_name() + } + + /// Return the most recent background worker failure, if one was observed. + pub fn last_worker_error(&self) -> Result> { + self.state + .last_worker_error + .lock() + .map(|error| error.clone()) + .map_err(|_| { + BootError::Internal("PostgreSQL queue worker error lock is poisoned".to_string()) + }) + } + + /// Return the current jobs without blocking the calling Tokio runtime. + pub async fn jobs_async(&self) -> Result> { + self.state.store.jobs().await + } + + /// Return retained job failures without blocking the calling Tokio runtime. + pub async fn failures_async(&self) -> Result> { + self.state.store.failures().await + } + + /// Return current queue counters without blocking the calling Tokio runtime. + pub async fn stats_async(&self) -> Result { + self.state.store.stats().await + } + + /// Remove every job owned by this queue without blocking the calling Tokio runtime. + pub async fn clear_async(&self) -> Result<()> { + let result = self.state.store.clear().await; + self.state.notify.notify_waiters(); + result + } +} + +impl QueueBackend for PostgresQueueBackend { + fn enqueue(&self, name: String, data: Value) -> BoxFuture<'static, Result> { + self.enqueue_with_options(name, data, QueueJobOptions::new()) + } + + fn enqueue_with_options( + &self, + name: String, + data: Value, + options: QueueJobOptions, + ) -> BoxFuture<'static, Result> { + let backend = self.clone(); + Box::pin(async move { + let receipt = backend.state.store.enqueue(name, data, options).await?; + backend.state.notify.notify_waiters(); + Ok(receipt) + }) + } + + fn process(&self, name: String, processor: Arc) -> Result<()> { + let name = name.trim().to_string(); + if name.is_empty() { + return Err(BootError::BadRequest( + "PostgreSQL queue processor name cannot be empty".to_string(), + )); + } + let mut processors = self.state.write_processors()?; + if processors.contains_key(&name) { + return Err(BootError::Conflict(format!( + "PostgreSQL queue processor is already registered: {name}" + ))); + } + processors.insert(name, processor); + self.state.notify.notify_waiters(); + Ok(()) + } + + fn jobs(&self) -> Result> { + let store = self.state.store.clone(); + postgres_sync_wait(async move { store.jobs().await }) + } + + fn failures(&self) -> Result> { + let store = self.state.store.clone(); + postgres_sync_wait(async move { store.failures().await }) + } + + fn stats(&self) -> Result { + let store = self.state.store.clone(); + postgres_sync_wait(async move { store.stats().await }) + } + + fn clear(&self) -> Result<()> { + let store = self.state.store.clone(); + let result = postgres_sync_wait(async move { store.clear().await }); + self.state.notify.notify_waiters(); + result + } + + fn start(&self, queue_name: String, module_ref: ModuleRef) -> BoxFuture<'static, Result<()>> { + let backend = self.clone(); + Box::pin(async move { + backend.options.validate()?; + if queue_name != backend.queue_name() { + return Err(BootError::BadRequest(format!( + "Boot queue name {queue_name} does not match PostgreSQL queue {}", + backend.queue_name() + ))); + } + let runtime = Handle::try_current().map_err(|error| { + BootError::Internal(format!( + "PostgreSQL queue requires a running Tokio runtime: {error}" + )) + })?; + let mut lifecycle = backend.state.lock_lifecycle()?; + if lifecycle.shutdown.is_some() { + return Ok(()); + } + let (shutdown, _) = watch::channel(false); + for index in 0..backend.options.worker_count { + let state = Arc::clone(&backend.state); + let module_ref = module_ref.clone(); + let receiver = shutdown.subscribe(); + let options = backend.options; + let worker_id = format!("{queue_name}-worker-{}", index + 1); + lifecycle.handles.push(runtime.spawn(async move { + run_postgres_worker(state, module_ref, options, worker_id, receiver).await; + })); + } + lifecycle.shutdown = Some(shutdown); + drop(lifecycle); + backend.state.notify.notify_waiters(); + Ok(()) + }) + } + + fn shutdown(&self) -> BoxFuture<'static, Result<()>> { + let backend = self.clone(); + Box::pin(async move { + let (shutdown, handles) = { + let mut lifecycle = backend.state.lock_lifecycle()?; + ( + lifecycle.shutdown.take(), + std::mem::take(&mut lifecycle.handles), + ) + }; + if let Some(shutdown) = shutdown { + let _ = shutdown.send(true); + } + backend.state.notify.notify_waiters(); + for handle in handles { + handle.await.map_err(|error| { + BootError::Internal(format!("PostgreSQL queue worker failed: {error}")) + })?; + } + Ok(()) + }) + } +} + +struct PostgresQueueState { + store: PostgresQueueStore, + processors: RwLock>>, + lifecycle: Mutex, + last_worker_error: Mutex>, + notify: Notify, +} + +#[derive(Default)] +struct PostgresQueueLifecycle { + shutdown: Option>, + handles: Vec>, +} + +impl PostgresQueueState { + fn write_processors( + &self, + ) -> Result>>> { + self.processors.write().map_err(|_| { + BootError::Internal("PostgreSQL queue processor registry lock is poisoned".to_string()) + }) + } + + fn processor_for(&self, name: &str) -> Result>> { + Ok(self + .processors + .read() + .map_err(|_| { + BootError::Internal( + "PostgreSQL queue processor registry lock is poisoned".to_string(), + ) + })? + .get(name) + .map(Arc::clone)) + } + + fn lock_lifecycle(&self) -> Result> { + self.lifecycle.lock().map_err(|_| { + BootError::Internal("PostgreSQL queue lifecycle lock is poisoned".to_string()) + }) + } + + fn record_worker_error(&self, error: &BootError) { + if let Ok(mut last_error) = self.last_worker_error.lock() { + *last_error = Some(error.to_string()); + } + } +} + +async fn run_postgres_worker( + state: Arc, + module_ref: ModuleRef, + options: QueueOptions, + worker_id: String, + mut shutdown: watch::Receiver, +) { + loop { + if *shutdown.borrow() { + return; + } + let result = + run_postgres_worker_once(&state, &module_ref, options, &worker_id, &mut shutdown).await; + if let Err(error) = result { + state.record_worker_error(&error); + wait_for_work(&state, options.poll_interval, &mut shutdown).await; + } + } +} + +async fn run_postgres_worker_once( + state: &Arc, + module_ref: &ModuleRef, + options: QueueOptions, + worker_id: &str, + shutdown: &mut watch::Receiver, +) -> Result<()> { + state.store.recover_expired(RECOVERY_BATCH_SIZE).await?; + let Some(job) = state.store.claim(worker_id, options.lease_duration).await? else { + wait_for_work(state, options.poll_interval, shutdown).await; + return Ok(()); + }; + let Some(processor) = state.processor_for(&job.name)? else { + state.store.release(&job.id, &job.lock_token).await?; + wait_for_work(state, options.poll_interval, shutdown).await; + return Ok(()); + }; + process_claimed_job(state, module_ref, options, job, processor, shutdown).await +} + +async fn process_claimed_job( + state: &Arc, + module_ref: &ModuleRef, + options: QueueOptions, + job: ClaimedJob, + processor: Arc, + shutdown: &mut watch::Receiver, +) -> Result<()> { + let queue_job = QueueJob { + id: job.id.clone(), + name: job.name, + data: job.payload, + }; + let context = QueueContext { + queue_name: state.store.queue_name().to_string(), + module_ref: module_ref.clone(), + }; + let processing = processor.process(queue_job, context); + tokio::pin!(processing); + let timeout = job.options.timeout; + let timeout_future = async move { + match timeout { + Some(timeout) => tokio::time::sleep(timeout).await, + None => pending::<()>().await, + } + }; + tokio::pin!(timeout_future); + let heartbeat_every = heartbeat_interval(options.lease_duration); + let mut heartbeat = tokio::time::interval_at( + tokio::time::Instant::now() + heartbeat_every, + heartbeat_every, + ); + loop { + tokio::select! { + changed = shutdown.changed() => { + let _ = changed; + state.store.release(&job.id, &job.lock_token).await?; + return Ok(()); + } + result = &mut processing => { + return match result { + Ok(()) => state.store.complete(&job.id, &job.lock_token).await, + Err(error) => state.store.fail(&job.id, &job.lock_token, error.to_string()).await, + }; + } + () = &mut timeout_future => { + return state.store.fail( + &job.id, + &job.lock_token, + format!("queue processor timed out after {timeout:?}"), + ).await; + } + _ = heartbeat.tick() => { + state.store.heartbeat(&job.id, &job.lock_token, options.lease_duration).await?; + } + } + } +} + +async fn wait_for_work( + state: &PostgresQueueState, + poll_interval: Duration, + shutdown: &mut watch::Receiver, +) { + tokio::select! { + _ = state.notify.notified() => {} + _ = tokio::time::sleep(poll_interval) => {} + _ = shutdown.changed() => {} + } +} + +fn heartbeat_interval(lease_duration: Duration) -> Duration { + let interval = lease_duration / 3; + if interval.is_zero() { + Duration::from_millis(1) + } else { + interval + } +} + +fn postgres_sync_wait(future: Fut) -> Result +where + T: Send + 'static, + Fut: Future> + Send + 'static, +{ + if Handle::try_current().is_ok() { + std::thread::spawn(move || run_on_postgres_runtime(future)) + .join() + .map_err(|_| BootError::Internal("PostgreSQL queue runtime thread panicked".into()))? + } else { + run_on_postgres_runtime(future) + } +} + +fn run_on_postgres_runtime(future: Fut) -> Result +where + Fut: Future>, +{ + TokioRuntimeBuilder::new_current_thread() + .enable_all() + .build() + .map_err(|error| { + BootError::Internal(format!( + "could not create PostgreSQL queue runtime: {error}" + )) + })? + .block_on(future) +} diff --git a/src/queue/postgres/deduplication.rs b/src/queue/postgres/deduplication.rs new file mode 100644 index 0000000..7f26e89 --- /dev/null +++ b/src/queue/postgres/deduplication.rs @@ -0,0 +1,157 @@ +use a3s_orm::{sql_query, PostgresTransaction}; + +use crate::Result; + +use super::super::QueueJobOptions; +use super::store::{add_duration, execute_query, fetch_one_query, fetch_optional_query}; + +pub(super) async fn lock_deduplication( + transaction: &PostgresTransaction, + queue_name: &str, + deduplication_id: &str, +) -> Result<()> { + let key = format!("a3s-boot:{queue_name}:{deduplication_id}"); + fetch_one_query( + transaction, + sql_query::("SELECT 1 FROM pg_advisory_xact_lock(hashtextextended(") + .bind(key) + .append(", 0))"), + ) + .await?; + Ok(()) +} + +pub(super) async fn release_expired_deduplication( + transaction: &PostgresTransaction, + queue_name: &str, + deduplication_id: &str, + now: i64, +) -> Result<()> { + execute_query( + transaction, + sql_query::<()>( + "UPDATE boot_queue_jobs SET deduplication_id = NULL, \ + deduplication_expires_at_nanos = NULL, updated_at_nanos = ", + ) + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND deduplication_id = ") + .bind(deduplication_id) + .append(" AND state IN ('pending', 'active') AND deduplication_expires_at_nanos <= ") + .bind(now), + ) + .await?; + Ok(()) +} + +pub(super) async fn find_deduplication_owner( + transaction: &PostgresTransaction, + queue_name: &str, + deduplication_id: &str, +) -> Result> { + fetch_optional_query( + transaction, + sql_query::<(String, String, String, i64)>( + "SELECT job_id, job_name, state, available_at_nanos FROM boot_queue_jobs \ + WHERE queue_name = ", + ) + .bind(queue_name) + .append(" AND deduplication_id = ") + .bind(deduplication_id) + .append(" AND state IN ('pending', 'active') FOR UPDATE"), + ) + .await +} + +pub(super) async fn update_deduplication_expiry( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + expires_at: Option, + now: i64, +) -> Result<()> { + execute_query( + transaction, + sql_query::<()>("UPDATE boot_queue_jobs SET deduplication_expires_at_nanos = ") + .bind(expires_at) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(job_id), + ) + .await?; + Ok(()) +} + +#[allow(clippy::too_many_arguments)] +pub(super) async fn store_successor( + transaction: &PostgresTransaction, + queue_name: &str, + owner_id: &str, + job_id: &str, + name: &str, + payload_json: &str, + options_json: &str, + now: i64, +) -> Result<()> { + execute_query( + transaction, + sql_query::<()>("UPDATE boot_queue_jobs SET successor_job_id = ") + .bind(job_id) + .append(", successor_job_name = ") + .bind(name) + .append(", successor_payload_json = ") + .bind(payload_json) + .append(", successor_options_json = ") + .bind(options_json) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(owner_id) + .append(" AND state = 'active'"), + ) + .await?; + Ok(()) +} + +#[allow(clippy::too_many_arguments)] +pub(super) async fn replace_delayed_owner( + transaction: &PostgresTransaction, + queue_name: &str, + owner_id: &str, + name: &str, + payload_json: &str, + options_json: &str, + options: &QueueJobOptions, + now: i64, +) -> Result<()> { + execute_query( + transaction, + sql_query::<()>("UPDATE boot_queue_jobs SET job_name = ") + .bind(name) + .append(", payload_json = ") + .bind(payload_json) + .append(", options_json = ") + .bind(options_json) + .append(", priority = ") + .bind(i64::from(options.priority)) + .append(", lifo = ") + .bind(options.lifo) + .append(", available_at_nanos = ") + .bind(options.delay.map_or(now, |delay| add_duration(now, delay))) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(owner_id) + .append(" AND state = 'pending'"), + ) + .await?; + Ok(()) +} diff --git a/src/queue/postgres/lifecycle.rs b/src/queue/postgres/lifecycle.rs new file mode 100644 index 0000000..ea6692c --- /dev/null +++ b/src/queue/postgres/lifecycle.rs @@ -0,0 +1,306 @@ +use a3s_orm::{sql_query, DecodeError, FromRow, PostgresTransaction, Row}; + +use crate::Result; + +use super::super::QueueJobOptions; +use super::store::{ + decode, ensure_idempotent_job, execute_query, insert_job, require_fenced_row, + stored_json_error, subtract_duration, +}; + +pub(super) struct ActiveJobRow { + pub options_json: String, + pub attempts_made: u32, + pub successor_job_id: Option, + pub successor_job_name: Option, + pub successor_payload_json: Option, + pub successor_options_json: Option, +} + +impl FromRow for ActiveJobRow { + fn from_row(row: &impl Row) -> std::result::Result { + Ok(Self { + options_json: decode(row, 0)?, + attempts_made: decode(row, 1)?, + successor_job_id: decode(row, 2)?, + successor_job_name: decode(row, 3)?, + successor_payload_json: decode(row, 4)?, + successor_options_json: decode(row, 5)?, + }) + } +} + +pub(super) struct ExpiredJobRow { + pub id: String, + pub lock_token: String, + pub options_json: String, + pub stalled_count: u32, + pub successor_job_id: Option, + pub successor_job_name: Option, + pub successor_payload_json: Option, + pub successor_options_json: Option, +} + +impl FromRow for ExpiredJobRow { + fn from_row(row: &impl Row) -> std::result::Result { + Ok(Self { + id: decode(row, 0)?, + lock_token: decode(row, 1)?, + options_json: decode(row, 2)?, + stalled_count: decode(row, 3)?, + successor_job_id: decode(row, 4)?, + successor_job_name: decode(row, 5)?, + successor_payload_json: decode(row, 6)?, + successor_options_json: decode(row, 7)?, + }) + } +} + +#[derive(Clone)] +pub(super) struct Successor { + id: String, + name: String, + payload_json: String, + options_json: String, +} + +pub(super) fn successor_from_active(row: &ActiveJobRow) -> Option { + successor( + row.successor_job_id.clone(), + row.successor_job_name.clone(), + row.successor_payload_json.clone(), + row.successor_options_json.clone(), + ) +} + +pub(super) fn successor_from_expired(row: &ExpiredJobRow) -> Option { + successor( + row.successor_job_id.clone(), + row.successor_job_name.clone(), + row.successor_payload_json.clone(), + row.successor_options_json.clone(), + ) +} + +fn successor( + id: Option, + name: Option, + payload_json: Option, + options_json: Option, +) -> Option { + match (id, name, payload_json, options_json) { + (Some(id), Some(name), Some(payload_json), Some(options_json)) => Some(Successor { + id, + name, + payload_json, + options_json, + }), + _ => None, + } +} + +pub(super) async fn terminal_completion( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + lock_token: &str, + options: &QueueJobOptions, + successor: Option, + now: i64, +) -> Result<()> { + let remove = options.remove_on_complete + || options + .completion_retention + .as_ref() + .is_some_and(|retention| retention.count == Some(0)); + set_terminal( + transaction, + queue_name, + job_id, + lock_token, + "completed", + None, + remove, + successor, + now, + ) + .await +} + +#[allow(clippy::too_many_arguments)] +pub(super) async fn terminal_failure( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + lock_token: &str, + options: &QueueJobOptions, + message: &str, + successor: Option, + now: i64, +) -> Result<()> { + let remove = options.remove_on_fail + || options + .failure_retention + .as_ref() + .is_some_and(|retention| retention.count == Some(0)); + set_terminal( + transaction, + queue_name, + job_id, + lock_token, + "failed", + Some(message), + remove, + successor, + now, + ) + .await +} + +#[allow(clippy::too_many_arguments)] +async fn set_terminal( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + lock_token: &str, + state: &str, + failure: Option<&str>, + remove: bool, + successor: Option, + now: i64, +) -> Result<()> { + let rows = if remove { + execute_query( + transaction, + sql_query::<()>("DELETE FROM boot_queue_jobs WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(job_id) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token), + ) + .await? + } else { + execute_query( + transaction, + sql_query::<()>("UPDATE boot_queue_jobs SET state = ") + .bind(state) + .append( + ", worker_id = NULL, lock_token = NULL, lease_expires_at_nanos = NULL, \ + failed_reason = ", + ) + .bind(failure) + .append( + ", deduplication_id = NULL, deduplication_expires_at_nanos = NULL, \ + successor_job_id = NULL, successor_job_name = NULL, \ + successor_payload_json = NULL, successor_options_json = NULL, \ + finished_at_nanos = ", + ) + .bind(now) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(job_id) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token), + ) + .await? + }; + require_fenced_row(rows, job_id, state)?; + if let Some(successor) = successor { + materialize_successor(transaction, queue_name, successor, now).await?; + } + Ok(()) +} + +async fn materialize_successor( + transaction: &PostgresTransaction, + queue_name: &str, + successor: Successor, + now: i64, +) -> Result<()> { + let options: QueueJobOptions = + serde_json::from_str(&successor.options_json).map_err(stored_json_error)?; + let inserted = insert_job( + transaction, + queue_name, + &successor.id, + &successor.name, + &successor.payload_json, + &successor.options_json, + &options, + now, + ) + .await?; + if !inserted { + ensure_idempotent_job( + transaction, + queue_name, + &successor.id, + &successor.name, + &successor.payload_json, + &successor.options_json, + ) + .await?; + } + Ok(()) +} + +pub(super) async fn apply_retention( + transaction: &PostgresTransaction, + queue_name: &str, + state: &str, + options: &QueueJobOptions, + now: i64, +) -> Result<()> { + let retention = match state { + "completed" if !options.remove_on_complete => options.completion_retention.as_ref(), + "failed" if !options.remove_on_fail => options.failure_retention.as_ref(), + _ => None, + }; + let Some(retention) = retention else { + return Ok(()); + }; + if let Some(age) = retention.age { + let limit = i64::try_from(retention.limit.unwrap_or(usize::MAX)).unwrap_or(i64::MAX); + execute_query( + transaction, + sql_query::<()>( + "DELETE FROM boot_queue_jobs WHERE (queue_name, job_id) IN (SELECT queue_name, \ + job_id FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(queue_name) + .append(" AND state = ") + .bind(state) + .append(" AND finished_at_nanos < ") + .bind(subtract_duration(now, age)) + .append(" ORDER BY finished_at_nanos ASC, job_id ASC LIMIT ") + .bind(limit) + .append(")"), + ) + .await?; + } + if let Some(count) = retention.count { + let offset = i64::try_from(count).unwrap_or(i64::MAX); + let limit = i64::try_from(retention.limit.unwrap_or(usize::MAX)).unwrap_or(i64::MAX); + execute_query( + transaction, + sql_query::<()>( + "DELETE FROM boot_queue_jobs WHERE (queue_name, job_id) IN (SELECT queue_name, \ + job_id FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(queue_name) + .append(" AND state = ") + .bind(state) + .append(" ORDER BY finished_at_nanos DESC, job_id DESC OFFSET ") + .bind(offset) + .append(" LIMIT ") + .bind(limit) + .append(")"), + ) + .await?; + } + Ok(()) +} diff --git a/src/queue/postgres/migrations.rs b/src/queue/postgres/migrations.rs new file mode 100644 index 0000000..f307e2a --- /dev/null +++ b/src/queue/postgres/migrations.rs @@ -0,0 +1,73 @@ +use a3s_orm::Migration; + +const POSTGRES_QUEUE_SQL: &str = r#" +CREATE TABLE IF NOT EXISTS boot_queue_jobs ( + queue_name TEXT NOT NULL, + job_id TEXT NOT NULL, + job_name TEXT NOT NULL, + payload_json TEXT NOT NULL, + options_json TEXT NOT NULL, + state TEXT NOT NULL CHECK (state IN ('pending', 'active', 'completed', 'failed')), + priority BIGINT NOT NULL CHECK (priority BETWEEN 0 AND 4294967295), + lifo BOOLEAN NOT NULL, + available_at_nanos BIGINT NOT NULL, + attempts_made BIGINT NOT NULL DEFAULT 0 CHECK (attempts_made >= 0), + stalled_count BIGINT NOT NULL DEFAULT 0 CHECK (stalled_count >= 0), + worker_id TEXT, + lock_token TEXT, + lease_expires_at_nanos BIGINT, + failed_reason TEXT, + deduplication_id TEXT, + deduplication_expires_at_nanos BIGINT, + successor_job_id TEXT, + successor_job_name TEXT, + successor_payload_json TEXT, + successor_options_json TEXT, + created_at_nanos BIGINT NOT NULL, + updated_at_nanos BIGINT NOT NULL, + finished_at_nanos BIGINT, + PRIMARY KEY (queue_name, job_id), + CHECK ( + (state = 'active' AND worker_id IS NOT NULL AND lock_token IS NOT NULL + AND lease_expires_at_nanos IS NOT NULL) + OR + (state <> 'active' AND worker_id IS NULL AND lock_token IS NULL + AND lease_expires_at_nanos IS NULL) + ), + CHECK ( + (successor_job_id IS NULL AND successor_job_name IS NULL + AND successor_payload_json IS NULL AND successor_options_json IS NULL) + OR + (successor_job_id IS NOT NULL AND successor_job_name IS NOT NULL + AND successor_payload_json IS NOT NULL AND successor_options_json IS NOT NULL) + ) +); + +CREATE INDEX IF NOT EXISTS boot_queue_jobs_claim_idx + ON boot_queue_jobs ( + queue_name, + state, + available_at_nanos, + priority, + created_at_nanos, + job_id + ); + +CREATE INDEX IF NOT EXISTS boot_queue_jobs_lease_idx + ON boot_queue_jobs (queue_name, state, lease_expires_at_nanos, job_id); + +CREATE INDEX IF NOT EXISTS boot_queue_jobs_terminal_idx + ON boot_queue_jobs (queue_name, state, finished_at_nanos DESC, job_id DESC); + +CREATE UNIQUE INDEX IF NOT EXISTS boot_queue_jobs_active_dedup_idx + ON boot_queue_jobs (queue_name, deduplication_id) + WHERE deduplication_id IS NOT NULL AND state IN ('pending', 'active'); +"#; + +pub(super) fn postgres_queue_migrations() -> Vec { + vec![Migration::new( + "a3s-boot-0001-queue", + "create durable Boot queue jobs", + POSTGRES_QUEUE_SQL, + )] +} diff --git a/src/queue/postgres/store.rs b/src/queue/postgres/store.rs new file mode 100644 index 0000000..17c5aa5 --- /dev/null +++ b/src/queue/postgres/store.rs @@ -0,0 +1,922 @@ +use std::time::Duration; + +use a3s_orm::{ + sql_query, DecodeError, Executor, FromRow, FromValue, Migrator, PostgresDialect, PostgresError, + PostgresExecutor, PostgresRow, PostgresTransaction, PostgresTransactionError, Query, Row, + SqlQuery, +}; +use chrono::Utc; +use serde_json::Value; +use uuid::Uuid; + +use crate::{BootError, Result}; + +use super::super::{ + QueueJobFailure, QueueJobInfo, QueueJobOptions, QueueJobReceipt, QueueJobRetention, + QueueJobState, QueueStats, +}; +use super::deduplication::{ + find_deduplication_owner, lock_deduplication, release_expired_deduplication, + replace_delayed_owner, store_successor, update_deduplication_expiry, +}; +use super::lifecycle::{ + apply_retention, successor_from_active, successor_from_expired, terminal_completion, + terminal_failure, ActiveJobRow, ExpiredJobRow, +}; +use super::migrations::postgres_queue_migrations; + +const STALLED_FAILURE: &str = "queue job exceeded its stalled lease recovery limit"; + +#[derive(Clone)] +pub(super) struct PostgresQueueStore { + executor: PostgresExecutor, + queue_name: String, +} + +#[derive(Debug)] +pub(super) struct ClaimedJob { + pub id: String, + pub name: String, + pub payload: Value, + pub options: QueueJobOptions, + pub lock_token: String, +} + +struct ClaimedJobRow { + id: String, + name: String, + payload_json: String, + options_json: String, + lock_token: String, +} + +impl FromRow for ClaimedJobRow { + fn from_row(row: &impl Row) -> std::result::Result { + Ok(Self { + id: decode(row, 0)?, + name: decode(row, 1)?, + payload_json: decode(row, 2)?, + options_json: decode(row, 3)?, + lock_token: decode(row, 4)?, + }) + } +} + +struct JobInfoRow { + id: String, + name: String, + state: String, + payload_json: String, +} + +impl FromRow for JobInfoRow { + fn from_row(row: &impl Row) -> std::result::Result { + Ok(Self { + id: decode(row, 0)?, + name: decode(row, 1)?, + state: decode(row, 2)?, + payload_json: decode(row, 3)?, + }) + } +} + +struct FailureRow { + id: String, + name: String, + message: Option, +} + +impl FromRow for FailureRow { + fn from_row(row: &impl Row) -> std::result::Result { + Ok(Self { + id: decode(row, 0)?, + name: decode(row, 1)?, + message: decode(row, 2)?, + }) + } +} + +impl PostgresQueueStore { + pub(super) async fn connect(database_url: &str, queue_name: &str) -> Result { + let executor = PostgresExecutor::connect_no_tls(database_url, 5).map_err(|error| { + BootError::Internal(format!( + "could not configure PostgreSQL Boot queue: {error}" + )) + })?; + Self::from_executor(executor, queue_name).await + } + + pub(super) async fn from_executor( + executor: PostgresExecutor, + queue_name: &str, + ) -> Result { + let queue_name = queue_name.trim(); + if queue_name.is_empty() { + return Err(BootError::BadRequest( + "PostgreSQL queue name cannot be empty".to_string(), + )); + } + Migrator::new(executor.clone()) + .run(postgres_queue_migrations()) + .await + .map_err(|error| { + BootError::Internal(format!("PostgreSQL Boot queue migration failed: {error}")) + })?; + Ok(Self { + executor, + queue_name: queue_name.to_string(), + }) + } + + pub(super) fn queue_name(&self) -> &str { + &self.queue_name + } + + pub(super) async fn enqueue( + &self, + name: String, + payload: Value, + options: QueueJobOptions, + ) -> Result { + let name = validate_job(name)?; + validate_options(&options)?; + let payload_json = serde_json::to_string(&payload).map_err(json_error)?; + let options_json = serde_json::to_string(&options).map_err(json_error)?; + let requested_job_id = options + .job_id + .clone() + .unwrap_or_else(|| Uuid::new_v4().to_string()); + let queue_name = self.queue_name.clone(); + let result = self + .executor + .transaction(|transaction| { + let name = name.clone(); + let payload_json = payload_json.clone(); + let options_json = options_json.clone(); + let requested_job_id = requested_job_id.clone(); + let options = options.clone(); + let queue_name = queue_name.clone(); + Box::pin(async move { + let now = now_nanos(); + if let Some(deduplication) = options.deduplication.as_ref() { + lock_deduplication(transaction, &queue_name, &deduplication.id).await?; + release_expired_deduplication( + transaction, + &queue_name, + &deduplication.id, + now, + ) + .await?; + if let Some((owner_id, owner_name, owner_state, available_at)) = + find_deduplication_owner(transaction, &queue_name, &deduplication.id) + .await? + { + if deduplication.extend { + update_deduplication_expiry( + transaction, + &queue_name, + &owner_id, + deduplication.ttl.map(|ttl| add_duration(now, ttl)), + now, + ) + .await?; + } + if owner_state == "active" && deduplication.keep_last_if_active { + store_successor( + transaction, + &queue_name, + &owner_id, + &requested_job_id, + &name, + &payload_json, + &options_json, + now, + ) + .await?; + } else if owner_state == "pending" + && available_at > now + && deduplication.replace + { + replace_delayed_owner( + transaction, + &queue_name, + &owner_id, + &name, + &payload_json, + &options_json, + &options, + now, + ) + .await?; + } + return Ok(QueueJobReceipt { + id: owner_id, + name: owner_name, + }); + } + } + + let inserted = insert_job( + transaction, + &queue_name, + &requested_job_id, + &name, + &payload_json, + &options_json, + &options, + now, + ) + .await?; + if inserted { + return Ok(QueueJobReceipt { + id: requested_job_id, + name, + }); + } + ensure_idempotent_job( + transaction, + &queue_name, + &requested_job_id, + &name, + &payload_json, + &options_json, + ) + .await?; + Ok(QueueJobReceipt { + id: requested_job_id, + name, + }) + }) + }) + .await; + map_transaction(result) + } + + pub(super) async fn recover_expired(&self, limit: usize) -> Result { + let queue_name = self.queue_name.clone(); + let limit = i64::try_from(limit).map_err(|error| { + BootError::Internal(format!( + "PostgreSQL queue recovery limit is invalid: {error}" + )) + })?; + let result = self + .executor + .transaction(|transaction| { + let queue_name = queue_name.clone(); + Box::pin(async move { + let now = now_nanos(); + let expired = fetch_all_query( + transaction, + sql_query::( + "SELECT job_id, lock_token, options_json, stalled_count, \ + successor_job_id, successor_job_name, successor_payload_json, \ + successor_options_json FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(queue_name.clone()) + .append(" AND state = 'active' AND lease_expires_at_nanos <= ") + .bind(now) + .append(" ORDER BY lease_expires_at_nanos ASC, job_id ASC FOR UPDATE SKIP LOCKED LIMIT ") + .bind(limit), + ) + .await?; + for job in &expired { + let options: QueueJobOptions = + serde_json::from_str(&job.options_json).map_err(stored_json_error)?; + if job.stalled_count >= options.max_stalled_count { + terminal_failure( + transaction, + &queue_name, + &job.id, + &job.lock_token, + &options, + STALLED_FAILURE, + successor_from_expired(job), + now, + ) + .await?; + apply_retention( + transaction, + &queue_name, + "failed", + &options, + now, + ) + .await?; + } else { + execute_query( + transaction, + sql_query::<()>( + "UPDATE boot_queue_jobs SET state = 'pending', worker_id = NULL, \ + lock_token = NULL, lease_expires_at_nanos = NULL, \ + attempts_made = GREATEST(attempts_made - 1, 0), \ + stalled_count = stalled_count + 1, updated_at_nanos = ", + ) + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name.clone()) + .append(" AND job_id = ") + .bind(job.id.clone()) + .append(" AND state = 'active' AND lock_token = ") + .bind(job.lock_token.clone()), + ) + .await?; + } + } + Ok(expired.len()) + }) + }) + .await; + map_transaction(result) + } + + pub(super) async fn claim( + &self, + worker_id: &str, + lease_duration: Duration, + ) -> Result> { + let now = now_nanos(); + let lock_token = Uuid::new_v4().to_string(); + let lease_expires_at = add_duration(now, lease_duration); + let row = fetch_optional_query( + &self.executor, + sql_query::( + "WITH next_job AS (SELECT job_id FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(self.queue_name.clone()) + .append(" AND state = 'pending' AND available_at_nanos <= ") + .bind(now) + .append( + " ORDER BY priority ASC, \ + CASE WHEN lifo THEN created_at_nanos END DESC NULLS LAST, \ + CASE WHEN NOT lifo THEN created_at_nanos END ASC NULLS LAST, job_id ASC \ + FOR UPDATE SKIP LOCKED LIMIT 1) UPDATE boot_queue_jobs SET state = 'active', \ + attempts_made = attempts_made + 1, worker_id = ", + ) + .bind(worker_id) + .append(", lock_token = ") + .bind(lock_token) + .append(", lease_expires_at_nanos = ") + .bind(lease_expires_at) + .append(", updated_at_nanos = ") + .bind(now) + .append(" FROM next_job WHERE boot_queue_jobs.queue_name = ") + .bind(self.queue_name.clone()) + .append( + " AND boot_queue_jobs.job_id = next_job.job_id RETURNING \ + boot_queue_jobs.job_id, boot_queue_jobs.job_name, \ + boot_queue_jobs.payload_json, boot_queue_jobs.options_json, \ + boot_queue_jobs.lock_token", + ), + ) + .await?; + row.map(claimed_job_from_row).transpose() + } + + pub(super) async fn heartbeat( + &self, + job_id: &str, + lock_token: &str, + lease_duration: Duration, + ) -> Result<()> { + let now = now_nanos(); + let rows = execute_query( + &self.executor, + sql_query::<()>("UPDATE boot_queue_jobs SET lease_expires_at_nanos = ") + .bind(add_duration(now, lease_duration)) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(self.queue_name.clone()) + .append(" AND job_id = ") + .bind(job_id) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token), + ) + .await?; + require_fenced_row(rows, job_id, "heartbeat") + } + + pub(super) async fn release(&self, job_id: &str, lock_token: &str) -> Result<()> { + let rows = execute_query( + &self.executor, + sql_query::<()>( + "UPDATE boot_queue_jobs SET state = 'pending', worker_id = NULL, \ + lock_token = NULL, lease_expires_at_nanos = NULL, \ + attempts_made = GREATEST(attempts_made - 1, 0), updated_at_nanos = ", + ) + .bind(now_nanos()) + .append(" WHERE queue_name = ") + .bind(self.queue_name.clone()) + .append(" AND job_id = ") + .bind(job_id) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token), + ) + .await?; + require_fenced_row(rows, job_id, "release") + } + + pub(super) async fn complete(&self, job_id: &str, lock_token: &str) -> Result<()> { + self.finish(job_id, lock_token, None).await + } + + pub(super) async fn fail(&self, job_id: &str, lock_token: &str, message: String) -> Result<()> { + self.finish(job_id, lock_token, Some(message)).await + } + + async fn finish(&self, job_id: &str, lock_token: &str, failure: Option) -> Result<()> { + let queue_name = self.queue_name.clone(); + let job_id = job_id.to_string(); + let lock_token = lock_token.to_string(); + let result = self + .executor + .transaction(|transaction| { + let queue_name = queue_name.clone(); + let job_id = job_id.clone(); + let lock_token = lock_token.clone(); + let failure = failure.clone(); + Box::pin(async move { + let row = fetch_optional_query( + transaction, + sql_query::( + "SELECT options_json, attempts_made, successor_job_id, \ + successor_job_name, successor_payload_json, successor_options_json \ + FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(queue_name.clone()) + .append(" AND job_id = ") + .bind(job_id.clone()) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token.clone()) + .append(" FOR UPDATE"), + ) + .await? + .ok_or_else(|| lease_conflict(&job_id, "finish"))?; + let options: QueueJobOptions = + serde_json::from_str(&row.options_json).map_err(stored_json_error)?; + let now = now_nanos(); + if let Some(message) = failure { + if options.retry_policy.max_retries > 0 + && row.attempts_made <= options.retry_policy.max_retries + { + let delay = options.retry_policy.delay_for_attempt(row.attempts_made); + let rows = execute_query( + transaction, + sql_query::<()>( + "UPDATE boot_queue_jobs SET state = 'pending', \ + available_at_nanos = ", + ) + .bind(add_duration(now, delay)) + .append( + ", worker_id = NULL, lock_token = NULL, \ + lease_expires_at_nanos = NULL, failed_reason = ", + ) + .bind(message) + .append(", updated_at_nanos = ") + .bind(now) + .append(" WHERE queue_name = ") + .bind(queue_name) + .append(" AND job_id = ") + .bind(job_id.clone()) + .append(" AND state = 'active' AND lock_token = ") + .bind(lock_token), + ) + .await?; + return require_fenced_row(rows, &job_id, "retry"); + } + terminal_failure( + transaction, + &queue_name, + &job_id, + &lock_token, + &options, + &message, + successor_from_active(&row), + now, + ) + .await?; + apply_retention(transaction, &queue_name, "failed", &options, now).await?; + } else { + terminal_completion( + transaction, + &queue_name, + &job_id, + &lock_token, + &options, + successor_from_active(&row), + now, + ) + .await?; + apply_retention(transaction, &queue_name, "completed", &options, now) + .await?; + } + Ok(()) + }) + }) + .await; + map_transaction(result) + } + + pub(super) async fn jobs(&self) -> Result> { + fetch_all_query( + &self.executor, + sql_query::( + "SELECT job_id, job_name, state, payload_json FROM boot_queue_jobs \ + WHERE queue_name = ", + ) + .bind(self.queue_name.clone()) + .append(" ORDER BY created_at_nanos ASC, job_id ASC"), + ) + .await? + .into_iter() + .map(job_info_from_row) + .collect() + } + + pub(super) async fn failures(&self) -> Result> { + Ok(fetch_all_query( + &self.executor, + sql_query::( + "SELECT job_id, job_name, failed_reason FROM boot_queue_jobs \ + WHERE queue_name = ", + ) + .bind(self.queue_name.clone()) + .append(" AND state = 'failed' ORDER BY finished_at_nanos ASC, job_id ASC"), + ) + .await? + .into_iter() + .map(|row| QueueJobFailure { + id: row.id, + name: row.name, + message: row + .message + .unwrap_or_else(|| "job failed without a retained reason".to_string()), + }) + .collect()) + } + + pub(super) async fn stats(&self) -> Result { + let rows = fetch_all_query( + &self.executor, + sql_query::<(String, i64)>( + "SELECT state, COUNT(*)::BIGINT FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(self.queue_name.clone()) + .append(" GROUP BY state"), + ) + .await?; + let mut stats = QueueStats::default(); + for (state, count) in rows { + let count = usize::try_from(count).map_err(|error| { + BootError::Internal(format!("stored PostgreSQL queue count is invalid: {error}")) + })?; + match state.as_str() { + "pending" => stats.pending = count, + "active" => stats.active = count, + "completed" => stats.completed = count, + "failed" => stats.failed = count, + other => { + return Err(BootError::Internal(format!( + "stored PostgreSQL queue state is invalid: {other}" + ))) + } + } + } + Ok(stats) + } + + pub(super) async fn clear(&self) -> Result<()> { + execute_query( + &self.executor, + sql_query::<()>("DELETE FROM boot_queue_jobs WHERE queue_name = ") + .bind(self.queue_name.clone()), + ) + .await?; + Ok(()) + } +} + +#[allow(clippy::too_many_arguments)] +pub(super) async fn insert_job( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + name: &str, + payload_json: &str, + options_json: &str, + options: &QueueJobOptions, + now: i64, +) -> Result { + let deduplication_id = options + .deduplication + .as_ref() + .map(|deduplication| deduplication.id.clone()); + let deduplication_expires_at = options + .deduplication + .as_ref() + .and_then(|deduplication| deduplication.ttl) + .map(|ttl| add_duration(now, ttl)); + let available_at = options.delay.map_or(now, |delay| add_duration(now, delay)); + let rows = execute_query( + transaction, + sql_query::<()>( + "INSERT INTO boot_queue_jobs (queue_name, job_id, job_name, payload_json, \ + options_json, state, priority, lifo, available_at_nanos, attempts_made, \ + stalled_count, deduplication_id, deduplication_expires_at_nanos, \ + created_at_nanos, updated_at_nanos) VALUES (", + ) + .bind(queue_name) + .append(", ") + .bind(job_id) + .append(", ") + .bind(name) + .append(", ") + .bind(payload_json) + .append(", ") + .bind(options_json) + .append(", 'pending', ") + .bind(i64::from(options.priority)) + .append(", ") + .bind(options.lifo) + .append(", ") + .bind(available_at) + .append(", 0, 0, ") + .bind(deduplication_id) + .append(", ") + .bind(deduplication_expires_at) + .append(", ") + .bind(now) + .append(", ") + .bind(now) + .append(") ON CONFLICT (queue_name, job_id) DO NOTHING"), + ) + .await?; + Ok(rows == 1) +} + +pub(super) async fn ensure_idempotent_job( + transaction: &PostgresTransaction, + queue_name: &str, + job_id: &str, + name: &str, + payload_json: &str, + options_json: &str, +) -> Result<()> { + let existing = fetch_optional_query( + transaction, + sql_query::<(String, String, String)>( + "SELECT job_name, payload_json, options_json FROM boot_queue_jobs WHERE queue_name = ", + ) + .bind(queue_name) + .append(" AND job_id = ") + .bind(job_id), + ) + .await?; + if existing + .as_ref() + .is_some_and(|(stored_name, stored_payload, stored_options)| { + stored_name == name && stored_payload == payload_json && stored_options == options_json + }) + { + return Ok(()); + } + Err(BootError::Conflict(format!( + "PostgreSQL queue job id {job_id} is already used by different work" + ))) +} + +fn validate_job(name: String) -> Result { + let name = name.trim().to_string(); + if name.is_empty() { + return Err(BootError::BadRequest( + "PostgreSQL queue job name cannot be empty".to_string(), + )); + } + Ok(name) +} + +fn validate_options(options: &QueueJobOptions) -> Result<()> { + if options + .job_id + .as_ref() + .is_some_and(|id| id.trim().is_empty()) + { + return Err(BootError::BadRequest( + "PostgreSQL queue job id cannot be empty".to_string(), + )); + } + if options.timeout.is_some_and(|timeout| timeout.is_zero()) { + return Err(BootError::BadRequest( + "PostgreSQL queue timeout must be greater than zero".to_string(), + )); + } + if !options.retry_policy.multiplier.is_finite() || options.retry_policy.multiplier < 0.0 { + return Err(BootError::BadRequest( + "PostgreSQL queue retry multiplier must be finite and non-negative".to_string(), + )); + } + if let Some(deduplication) = options.deduplication.as_ref() { + if deduplication.id.trim().is_empty() { + return Err(BootError::BadRequest( + "PostgreSQL queue deduplication id cannot be empty".to_string(), + )); + } + if deduplication.ttl.is_some_and(|ttl| ttl.is_zero()) { + return Err(BootError::BadRequest( + "PostgreSQL queue deduplication TTL must be greater than zero".to_string(), + )); + } + } + if let Some(retention) = options.completion_retention.as_ref() { + validate_retention("completion retention", retention)?; + } + if let Some(retention) = options.failure_retention.as_ref() { + validate_retention("failure retention", retention)?; + } + if options.repeat.is_some() + || options.ignore_dependency_on_failure + || options.remove_dependency_on_failure + || options.continue_parent_on_failure + || options.fail_parent_on_failure + { + return Err(BootError::NotImplemented( + "PostgreSQL Boot queues do not support Lane flow or repeat options".to_string(), + )); + } + Ok(()) +} + +fn validate_retention(label: &str, retention: &QueueJobRetention) -> Result<()> { + if retention.age.is_none() && retention.count.is_none() { + return Err(BootError::BadRequest(format!( + "PostgreSQL queue {label} must specify an age or count" + ))); + } + if retention.age.is_some_and(|age| age.is_zero()) { + return Err(BootError::BadRequest(format!( + "PostgreSQL queue {label} age must be greater than zero" + ))); + } + if retention.limit == Some(0) { + return Err(BootError::BadRequest(format!( + "PostgreSQL queue {label} limit must be greater than zero" + ))); + } + Ok(()) +} + +fn claimed_job_from_row(row: ClaimedJobRow) -> Result { + Ok(ClaimedJob { + id: row.id, + name: row.name, + payload: serde_json::from_str(&row.payload_json).map_err(stored_json_error)?, + options: serde_json::from_str(&row.options_json).map_err(stored_json_error)?, + lock_token: row.lock_token, + }) +} + +fn job_info_from_row(row: JobInfoRow) -> Result { + let state = match row.state.as_str() { + "pending" => QueueJobState::Pending, + "active" => QueueJobState::Active, + "completed" => QueueJobState::Completed, + "failed" => QueueJobState::Failed, + other => { + return Err(BootError::Internal(format!( + "stored PostgreSQL queue state is invalid: {other}" + ))) + } + }; + Ok(QueueJobInfo { + id: row.id, + name: row.name, + state, + data: serde_json::from_str(&row.payload_json).map_err(stored_json_error)?, + }) +} + +pub(super) fn require_fenced_row(rows: u64, job_id: &str, action: &str) -> Result<()> { + if rows == 1 { + return Ok(()); + } + Err(lease_conflict(job_id, action)) +} + +fn lease_conflict(job_id: &str, action: &str) -> BootError { + BootError::Conflict(format!( + "PostgreSQL queue job {job_id} lost its lease before {action}" + )) +} + +fn now_nanos() -> i64 { + Utc::now().timestamp_nanos_opt().unwrap_or(i64::MAX) +} + +pub(super) fn add_duration(value: i64, duration: Duration) -> i64 { + let nanos = i64::try_from(duration.as_nanos()).unwrap_or(i64::MAX); + value.saturating_add(nanos) +} + +pub(super) fn subtract_duration(value: i64, duration: Duration) -> i64 { + let nanos = i64::try_from(duration.as_nanos()).unwrap_or(i64::MAX); + value.saturating_sub(nanos) +} + +fn json_error(error: serde_json::Error) -> BootError { + BootError::BadRequest(format!("could not encode PostgreSQL queue job: {error}")) +} + +pub(super) fn stored_json_error(error: serde_json::Error) -> BootError { + BootError::Internal(format!("stored PostgreSQL queue job is invalid: {error}")) +} + +fn map_transaction( + result: std::result::Result>, +) -> Result { + match result { + Ok(value) => Ok(value), + Err(PostgresTransactionError::Operation(error)) => Err(error), + Err(error) => Err(BootError::Internal(format!( + "PostgreSQL queue transaction failed: {error}" + ))), + } +} + +pub(super) async fn execute_query(executor: &E, query: SqlQuery<()>) -> Result +where + E: Executor, +{ + let query = query.compile(&PostgresDialect).map_err(query_error)?; + Ok(executor + .execute(&query) + .await + .map_err(database_error)? + .rows_affected) +} + +async fn fetch_all_query(executor: &E, query: SqlQuery) -> Result> +where + T: FromRow + Send, + E: Executor, +{ + let query = query.compile(&PostgresDialect).map_err(query_error)?; + executor + .fetch_all(&query) + .await + .map_err(database_error)? + .rows + .iter() + .map(T::from_row) + .collect::, _>>() + .map_err(decode_error) +} + +pub(super) async fn fetch_optional_query( + executor: &E, + query: SqlQuery, +) -> Result> +where + T: FromRow + Send, + E: Executor, +{ + let mut rows = fetch_all_query(executor, query).await?; + match rows.len() { + 0 => Ok(None), + 1 => Ok(rows.pop()), + actual => Err(BootError::Internal(format!( + "PostgreSQL queue returned {actual} rows where at most one was expected" + ))), + } +} + +pub(super) async fn fetch_one_query(executor: &E, query: SqlQuery) -> Result +where + T: FromRow + Send, + E: Executor, +{ + fetch_optional_query(executor, query) + .await? + .ok_or_else(|| BootError::Internal("PostgreSQL queue returned no row".to_string())) +} + +pub(super) fn decode( + row: &impl Row, + index: usize, +) -> std::result::Result { + let value = row + .value(index) + .ok_or(DecodeError::MissingColumn { index })?; + T::from_value(value, index) +} + +fn query_error(error: a3s_orm::Error) -> BootError { + BootError::Internal(format!("PostgreSQL queue query build failed: {error}")) +} + +fn database_error(error: PostgresError) -> BootError { + BootError::Internal(format!("PostgreSQL queue storage failed: {error}")) +} + +fn decode_error(error: DecodeError) -> BootError { + BootError::Internal(format!("PostgreSQL queue row decoding failed: {error}")) +} diff --git a/tests/postgres_queue.rs b/tests/postgres_queue.rs new file mode 100644 index 0000000..df01512 --- /dev/null +++ b/tests/postgres_queue.rs @@ -0,0 +1,454 @@ +#![cfg(feature = "queue-postgres")] + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::Duration; +use std::{process::Command, process::Stdio}; + +use a3s_boot::{ + ModuleRef, PostgresQueueBackend, Queue, QueueContext, QueueJob, QueueJobOptions, QueueOptions, + QueueRetryPolicy, +}; +use serde_json::json; +use tokio::sync::Notify; +use uuid::Uuid; + +fn postgres_url() -> Option { + std::env::var("A3S_BOOT_POSTGRES_URL") + .ok() + .filter(|value| !value.trim().is_empty()) +} + +fn queue_name(label: &str) -> String { + format!("boot-{label}-{}", Uuid::new_v4()) +} + +async fn wait_until(mut predicate: impl FnMut() -> bool) { + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if predicate() { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("PostgreSQL queue condition should become true"); +} + +#[tokio::test] +async fn queued_job_survives_backend_reconstruction() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("restart"); + let options = QueueOptions::new() + .with_poll_interval(Duration::from_millis(5)) + .with_lease_duration(Duration::from_millis(200)); + let first = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect first PostgreSQL queue backend"); + Queue::new(name.clone(), first) + .enqueue("durable", &json!({"value": 42})) + .await + .expect("enqueue durable job"); + + let calls = Arc::new(AtomicUsize::new(0)); + let second = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("reconnect PostgreSQL queue backend"); + let queue = Queue::new(name, second); + let observed = Arc::clone(&calls); + queue + .process("durable", move |job: QueueJob, _context: QueueContext| { + let observed = Arc::clone(&observed); + async move { + assert_eq!(job.data, json!({"value": 42})); + observed.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + }) + .expect("register durable processor"); + queue + .start(ModuleRef::new()) + .await + .expect("start reconstructed queue"); + wait_until(|| calls.load(Ordering::SeqCst) == 1).await; + queue.shutdown().await.expect("shutdown queue"); + assert_eq!(queue.stats().expect("queue stats").completed, 1); +} + +#[tokio::test] +async fn processor_failures_follow_the_typed_retry_policy() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("retry"); + let backend = PostgresQueueBackend::connect( + &url, + &name, + QueueOptions::new().with_poll_interval(Duration::from_millis(5)), + ) + .await + .expect("connect PostgreSQL queue backend"); + let queue = Queue::new(name, backend); + let calls = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&calls); + queue + .process("retry", move |_job: QueueJob, _context: QueueContext| { + let observed = Arc::clone(&observed); + async move { + if observed.fetch_add(1, Ordering::SeqCst) == 0 { + return Err(a3s_boot::BootError::Internal("first attempt".into())); + } + Ok(()) + } + }) + .expect("register retry processor"); + queue + .enqueue_with_options( + "retry", + &json!({}), + QueueJobOptions::new() + .with_retry_policy(QueueRetryPolicy::fixed(1, Duration::from_millis(5))), + ) + .await + .expect("enqueue retried job"); + queue.start(ModuleRef::new()).await.expect("start queue"); + wait_until(|| calls.load(Ordering::SeqCst) == 2).await; + queue.shutdown().await.expect("shutdown queue"); + let stats = queue.stats().expect("queue stats"); + assert_eq!(stats.completed, 1); + assert_eq!(stats.failed, 0); +} + +#[tokio::test] +async fn active_deduplication_keeps_only_the_latest_successor() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("dedup"); + let backend = PostgresQueueBackend::connect( + &url, + &name, + QueueOptions::new().with_poll_interval(Duration::from_millis(5)), + ) + .await + .expect("connect PostgreSQL queue backend"); + let queue = Queue::new(name, backend); + let started = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let values = Arc::new(std::sync::Mutex::new(Vec::new())); + let processor_started = Arc::clone(&started); + let processor_release = Arc::clone(&release); + let observed = Arc::clone(&values); + queue + .process("dedup", move |job: QueueJob, _context: QueueContext| { + let started = Arc::clone(&processor_started); + let release = Arc::clone(&processor_release); + let observed = Arc::clone(&observed); + async move { + let value = job.data["value"].as_u64().expect("numeric value"); + observed.lock().expect("values lock").push(value); + if value == 1 { + started.notify_one(); + release.notified().await; + } + Ok(()) + } + }) + .expect("register deduplicated processor"); + let deduplicated = || { + let mut options = QueueJobOptions::new().with_deduplication_id("same-target"); + options + .deduplication + .as_mut() + .expect("deduplication options") + .keep_last_if_active = true; + options + }; + queue + .enqueue_with_options("dedup", &json!({"value": 1}), deduplicated()) + .await + .expect("enqueue active owner"); + queue.start(ModuleRef::new()).await.expect("start queue"); + tokio::time::timeout(Duration::from_secs(5), started.notified()) + .await + .expect("first job should become active"); + queue + .enqueue_with_options("dedup", &json!({"value": 2}), deduplicated()) + .await + .expect("enqueue first successor"); + queue + .enqueue_with_options("dedup", &json!({"value": 3}), deduplicated()) + .await + .expect("replace successor with latest payload"); + release.notify_one(); + wait_until(|| values.lock().expect("values lock").len() == 2).await; + queue.shutdown().await.expect("shutdown queue"); + assert_eq!(*values.lock().expect("values lock"), vec![1, 3]); +} + +#[tokio::test] +async fn processor_timeout_is_fenced_and_recorded() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("timeout"); + let backend = PostgresQueueBackend::connect( + &url, + &name, + QueueOptions::new().with_poll_interval(Duration::from_millis(5)), + ) + .await + .expect("connect timeout queue"); + let diagnostics = backend.clone(); + let queue = Queue::new(name, backend); + queue + .process( + "timeout", + |_job: QueueJob, _context: QueueContext| async move { + std::future::pending::>().await + }, + ) + .expect("register timeout processor"); + queue + .enqueue_with_options( + "timeout", + &json!({}), + QueueJobOptions::new().with_timeout(Duration::from_millis(20)), + ) + .await + .expect("enqueue timeout job"); + queue.start(ModuleRef::new()).await.expect("start queue"); + let timeout_result = tokio::time::timeout(Duration::from_secs(5), async { + loop { + if diagnostics + .stats_async() + .await + .expect("timeout stats") + .failed + == 1 + { + break; + } + assert_eq!(diagnostics.last_worker_error().expect("worker error"), None); + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await; + if timeout_result.is_err() { + let worker_error = diagnostics.last_worker_error().expect("worker error"); + let shutdown_error = queue.shutdown().await.err(); + panic!( + "timeout job did not fail; worker_error={worker_error:?}, shutdown_error={shutdown_error:?}" + ); + } + queue.shutdown().await.expect("shutdown queue"); + let failures = diagnostics + .failures_async() + .await + .expect("timeout failures"); + assert_eq!(failures.len(), 1); + assert!(failures[0].message.contains("timed out")); +} + +#[tokio::test] +async fn heartbeat_prevents_competing_workers_from_reclaiming_live_work() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("heartbeat"); + let options = QueueOptions::new() + .with_worker_count(2) + .with_poll_interval(Duration::from_millis(5)) + .with_lease_duration(Duration::from_millis(60)); + let backend = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect heartbeat queue"); + let diagnostics = backend.clone(); + let queue = Queue::new(name, backend); + let calls = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&calls); + queue + .process("slow", move |_job: QueueJob, _context: QueueContext| { + let observed = Arc::clone(&observed); + async move { + observed.fetch_add(1, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(220)).await; + Ok(()) + } + }) + .expect("register slow processor"); + queue + .enqueue("slow", &json!({})) + .await + .expect("enqueue slow job"); + queue.start(ModuleRef::new()).await.expect("start queue"); + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if diagnostics + .stats_async() + .await + .expect("heartbeat stats") + .completed + == 1 + { + break; + } + assert_eq!(diagnostics.last_worker_error().expect("worker error"), None); + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("heartbeat job should complete"); + queue.shutdown().await.expect("shutdown queue"); + assert_eq!(calls.load(Ordering::SeqCst), 1); + assert_eq!( + diagnostics + .stats_async() + .await + .expect("heartbeat stats") + .failed, + 0 + ); +} + +#[tokio::test] +async fn expired_lease_recovers_the_same_job_after_process_death() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("process-death"); + let options = QueueOptions::new() + .with_poll_interval(Duration::from_millis(5)) + .with_lease_duration(Duration::from_millis(150)); + let backend = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect process-death queue"); + let queue = Queue::new(name.clone(), backend); + let receipt = queue + .enqueue("block", &json!({"effectId": "stable-effect"})) + .await + .expect("enqueue process-death job"); + + let executable = std::env::current_exe().expect("resolve PostgreSQL test executable"); + let child = Command::new(executable) + .arg("--exact") + .arg("postgres_queue_process_death_probe") + .arg("--nocapture") + .env("A3S_BOOT_POSTGRES_URL", &url) + .env("A3S_BOOT_POSTGRES_CRASH_PROBE", "1") + .env("A3S_BOOT_POSTGRES_CRASH_QUEUE", &name) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()) + .spawn() + .expect("spawn PostgreSQL queue crash probe"); + let mut child = ChildGuard(Some(child)); + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if queue.stats().expect("process-death stats").active == 1 { + return; + } + if child + .0 + .as_mut() + .expect("crash probe child") + .try_wait() + .expect("poll crash probe") + .is_some() + { + panic!("PostgreSQL queue crash probe exited before leasing work"); + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("crash probe should lease the job"); + let mut killed = child.0.take().expect("take crash probe child"); + killed.kill().expect("kill PostgreSQL queue worker process"); + let killed_status = killed.wait().expect("wait for killed queue worker"); + assert!(!killed_status.success()); + + tokio::time::sleep(Duration::from_millis(200)).await; + let recovered = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("reconnect after process death"); + let recovered_queue = Queue::new(name, recovered); + let recovered_job_id = receipt.id.clone(); + let completed = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&completed); + recovered_queue + .process("block", move |job: QueueJob, _context: QueueContext| { + let observed = Arc::clone(&observed); + let recovered_job_id = recovered_job_id.clone(); + async move { + assert_eq!(job.id, recovered_job_id); + assert_eq!(job.data["effectId"], "stable-effect"); + observed.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + }) + .expect("register recovered processor"); + recovered_queue + .start(ModuleRef::new()) + .await + .expect("start recovered queue"); + wait_until(|| completed.load(Ordering::SeqCst) == 1).await; + recovered_queue + .shutdown() + .await + .expect("shutdown recovered queue"); + let stats = recovered_queue.stats().expect("recovered queue stats"); + assert_eq!(stats.active, 0); + assert_eq!(stats.pending, 0); + assert_eq!(stats.completed, 1); + assert_eq!(stats.failed, 0); +} + +#[tokio::test] +async fn postgres_queue_process_death_probe() { + if std::env::var("A3S_BOOT_POSTGRES_CRASH_PROBE").as_deref() != Ok("1") { + return; + } + let url = std::env::var("A3S_BOOT_POSTGRES_URL").expect("crash probe PostgreSQL URL"); + let name = std::env::var("A3S_BOOT_POSTGRES_CRASH_QUEUE").expect("crash probe queue name"); + let options = QueueOptions::new() + .with_poll_interval(Duration::from_millis(5)) + .with_lease_duration(Duration::from_millis(150)); + let backend = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect crash probe queue"); + let queue = Queue::new(name, backend); + queue + .process( + "block", + |_job: QueueJob, _context: QueueContext| async move { + std::future::pending::>().await + }, + ) + .expect("register blocking crash probe processor"); + queue + .start(ModuleRef::new()) + .await + .expect("start crash probe queue"); + std::future::pending::<()>().await; +} + +struct ChildGuard(Option); + +impl Drop for ChildGuard { + fn drop(&mut self) { + if let Some(child) = self.0.as_mut() { + let _ = child.kill(); + let _ = child.wait(); + } + } +} diff --git a/tests/postgres_queue_contract.rs b/tests/postgres_queue_contract.rs new file mode 100644 index 0000000..7d64a30 --- /dev/null +++ b/tests/postgres_queue_contract.rs @@ -0,0 +1,166 @@ +#![cfg(feature = "queue-postgres")] + +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::Arc; +use std::time::Duration; + +use a3s_boot::{ + BootError, ModuleRef, PostgresQueueBackend, Queue, QueueJobOptions, QueueJobRetention, + QueueOptions, +}; +use serde_json::json; +use uuid::Uuid; + +fn postgres_url() -> Option { + std::env::var("A3S_BOOT_POSTGRES_URL") + .ok() + .filter(|value| !value.trim().is_empty()) +} + +fn queue_name(label: &str) -> String { + format!("boot-{label}-{}", Uuid::new_v4()) +} + +async fn wait_until(mut predicate: impl FnMut() -> bool) { + tokio::time::timeout(Duration::from_secs(5), async { + loop { + if predicate() { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("PostgreSQL queue condition should become true"); +} + +#[tokio::test] +async fn caller_assigned_job_ids_are_idempotent_and_conflicts_are_explicit() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("idempotency"); + let backend = PostgresQueueBackend::connect(&url, &name, QueueOptions::new()) + .await + .expect("connect idempotency queue"); + let diagnostics = backend.clone(); + let queue = Queue::new(name, backend); + let options = QueueJobOptions::new().with_job_id("stable-job-id"); + + let first = queue + .enqueue_with_options("work", &json!({"value": 1}), options.clone()) + .await + .expect("enqueue caller-assigned job"); + let replay = queue + .enqueue_with_options("work", &json!({"value": 1}), options.clone()) + .await + .expect("replay caller-assigned job"); + assert_eq!(first, replay); + assert_eq!( + diagnostics + .stats_async() + .await + .expect("queue stats") + .pending, + 1 + ); + + let error = queue + .enqueue_with_options("work", &json!({"value": 2}), options) + .await + .expect_err("different work must not reuse a caller-assigned job id"); + assert!(matches!(error, BootError::Conflict(_))); + + diagnostics.clear_async().await.expect("clear queue"); + assert_eq!( + diagnostics + .stats_async() + .await + .expect("queue stats") + .pending, + 0 + ); +} + +#[tokio::test] +async fn completion_retention_keeps_only_the_newest_terminal_job() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("retention"); + let backend = PostgresQueueBackend::connect( + &url, + &name, + QueueOptions::new().with_poll_interval(Duration::from_millis(5)), + ) + .await + .expect("connect retention queue"); + let diagnostics = backend.clone(); + let queue = Queue::new(name, backend); + let calls = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&calls); + queue + .process("retain", move |_job, _context| { + let observed = Arc::clone(&observed); + async move { + observed.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + }) + .expect("register retention processor"); + for value in 1..=3 { + queue + .enqueue_with_options( + "retain", + &json!({"value": value}), + QueueJobOptions::new().with_completion_retention(QueueJobRetention::count(1)), + ) + .await + .expect("enqueue retained job"); + } + queue.start(ModuleRef::new()).await.expect("start queue"); + wait_until(|| calls.load(Ordering::SeqCst) == 3).await; + queue.shutdown().await.expect("shutdown queue"); + + let jobs = diagnostics.jobs_async().await.expect("retained jobs"); + assert_eq!(jobs.len(), 1); + assert_eq!(jobs[0].data, json!({"value": 3})); + assert_eq!( + diagnostics + .stats_async() + .await + .expect("queue stats") + .completed, + 1 + ); +} + +#[tokio::test] +async fn malformed_retention_is_rejected_before_storage() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = queue_name("retention-validation"); + let backend = PostgresQueueBackend::connect(&url, &name, QueueOptions::new()) + .await + .expect("connect retention validation queue"); + let queue = Queue::new(name, backend); + let malformed = QueueJobRetention { + age: None, + count: None, + limit: None, + }; + + let error = queue + .enqueue_with_options( + "work", + &json!({}), + QueueJobOptions::new().with_completion_retention(malformed), + ) + .await + .expect_err("empty retention policy must be rejected"); + assert!(matches!(error, BootError::BadRequest(_))); +} diff --git a/tests/postgres_queue_shared.rs b/tests/postgres_queue_shared.rs new file mode 100644 index 0000000..bd02149 --- /dev/null +++ b/tests/postgres_queue_shared.rs @@ -0,0 +1,102 @@ +#![cfg(feature = "queue-postgres")] + +use std::collections::BTreeSet; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::Duration; + +use a3s_boot::{ModuleRef, PostgresQueueBackend, Queue, QueueContext, QueueJob, QueueOptions}; +use serde_json::json; +use uuid::Uuid; + +fn postgres_url() -> Option { + std::env::var("A3S_BOOT_POSTGRES_URL") + .ok() + .filter(|value| !value.trim().is_empty()) +} + +#[tokio::test] +async fn independent_backends_share_leases_without_duplicate_processing() { + let Some(url) = postgres_url() else { + eprintln!("skipping PostgreSQL queue test; set A3S_BOOT_POSTGRES_URL"); + return; + }; + let name = format!("boot-shared-{}", Uuid::new_v4()); + let options = QueueOptions::new() + .with_poll_interval(Duration::from_millis(5)) + .with_lease_duration(Duration::from_millis(60)); + let first_backend = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect first shared backend"); + let diagnostics = first_backend.clone(); + let second_backend = PostgresQueueBackend::connect(&url, &name, options) + .await + .expect("connect second shared backend"); + let first = Queue::new(name.clone(), first_backend); + let second = Queue::new(name, second_backend); + let calls = Arc::new(AtomicUsize::new(0)); + let processed = Arc::new(Mutex::new(Vec::new())); + let processor = || { + let calls = Arc::clone(&calls); + let processed = Arc::clone(&processed); + move |job: QueueJob, _context: QueueContext| { + let calls = Arc::clone(&calls); + let processed = Arc::clone(&processed); + async move { + tokio::time::sleep(Duration::from_millis(90)).await; + processed.lock().expect("processed lock").push(job.id); + calls.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + } + }; + first + .process("shared", processor()) + .expect("register first processor"); + second + .process("shared", processor()) + .expect("register second processor"); + + let mut expected = BTreeSet::new(); + for value in 0..12 { + let receipt = first + .enqueue("shared", &json!({"value": value})) + .await + .expect("enqueue shared job"); + expected.insert(receipt.id); + } + first + .start(ModuleRef::new()) + .await + .expect("start first queue"); + second + .start(ModuleRef::new()) + .await + .expect("start second queue"); + tokio::time::timeout(Duration::from_secs(10), async { + while calls.load(Ordering::SeqCst) < expected.len() { + tokio::time::sleep(Duration::from_millis(10)).await; + } + }) + .await + .expect("both backends should drain shared work"); + first.shutdown().await.expect("shutdown first queue"); + second.shutdown().await.expect("shutdown second queue"); + + let actual = processed + .lock() + .expect("processed lock") + .iter() + .cloned() + .collect::>(); + assert_eq!(calls.load(Ordering::SeqCst), expected.len()); + assert_eq!(actual, expected); + assert_eq!( + diagnostics + .stats_async() + .await + .expect("queue stats") + .completed, + 12 + ); +} diff --git a/tests/queue.rs b/tests/queue.rs index 35f47f5..d457371 100644 --- a/tests/queue.rs +++ b/tests/queue.rs @@ -2,7 +2,8 @@ use a3s_boot::{ BootApplication, BootError, BoxFuture, Module, ModuleRef, ProviderDefinition, Queue, - QueueContext, QueueJob, QueueJobOptions, QueueJobState, QueueModule, QueueOptions, Result, + QueueContext, QueueJob, QueueJobOptions, QueueJobState, QueueModule, QueueOptions, + QueueRetryPolicy, Result, }; use serde::{Deserialize, Serialize}; use std::sync::atomic::{AtomicUsize, Ordering}; @@ -286,6 +287,61 @@ async fn queue_records_processor_failures_and_job_states() { assert_eq!(jobs[0].state, QueueJobState::Failed); } +#[tokio::test] +async fn queue_applies_retry_and_timeout_options_to_processors() { + let attempts = Arc::new(AtomicUsize::new(0)); + let observed = Arc::clone(&attempts); + let queue = Queue::in_process_with_options( + "retry-queue", + QueueOptions::new().with_poll_interval(Duration::from_millis(5)), + ); + queue + .process("retry", move |_job, _context| { + let observed = Arc::clone(&observed); + async move { + if observed.fetch_add(1, Ordering::SeqCst) == 0 { + return Err(BootError::Internal("transient failure".to_string())); + } + Ok(()) + } + }) + .unwrap(); + queue + .process("timeout", |_job, _context| async move { + std::future::pending::>().await + }) + .unwrap(); + queue + .enqueue_with_options( + "retry", + &EmailJob { + to: "retry@example.com".to_string(), + }, + QueueJobOptions::new() + .with_retry_policy(QueueRetryPolicy::fixed(1, Duration::from_millis(5))), + ) + .await + .unwrap(); + queue + .enqueue_with_options( + "timeout", + &EmailJob { + to: "timeout@example.com".to_string(), + }, + QueueJobOptions::new().with_timeout(Duration::from_millis(10)), + ) + .await + .unwrap(); + queue.start(ModuleRef::new()).await.unwrap(); + wait_until(|| attempts.load(Ordering::SeqCst) == 2 && queue.stats().unwrap().failed == 1).await; + queue.shutdown().await.unwrap(); + + let stats = queue.stats().unwrap(); + assert_eq!(stats.completed, 1); + assert_eq!(stats.failed, 1); + assert!(queue.failures().unwrap()[0].message.contains("timed out")); +} + #[test] fn queue_validates_job_names_and_worker_count() { let app = BootApplication::builder() From 0129e8cafe8c353349f8a64611e41ce9197f7485 Mon Sep 17 00:00:00 2001 From: RoyLin Date: Sun, 9 Aug 2026 08:56:12 +0800 Subject: [PATCH 2/2] release: prepare a3s-boot 0.1.4 --- Cargo.toml | 2 +- README.md | 8 ++++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 9824f32..03a04e6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "a3s-boot" -version = "0.1.3" +version = "0.1.4" edition = "2021" authors = ["A3S Lab"] license = "MIT" diff --git a/README.md b/README.md index d094bac..ac102db 100644 --- a/README.md +++ b/README.md @@ -144,7 +144,7 @@ not a claim of support for every production backend. ```toml [dependencies] -a3s-boot = "0.1.3" +a3s-boot = "0.1.4" tokio = { version = "1", features = ["macros", "rt-multi-thread"] } ``` @@ -152,14 +152,14 @@ For a core-only build without Axum, macros, or shutdown signal handling: ```toml [dependencies] -a3s-boot = { version = "0.1.3", default-features = false } +a3s-boot = { version = "0.1.4", default-features = false } ``` Enable only the optional modules an application uses: ```toml [dependencies] -a3s-boot = { version = "0.1.3", features = ["auth", "security", "openapi-schemas"] } +a3s-boot = { version = "0.1.4", features = ["auth", "security", "openapi-schemas"] } serde = { version = "1", features = ["derive"] } tokio = { version = "1", features = ["macros", "rt-multi-thread"] } ``` @@ -260,7 +260,7 @@ expired leases after worker or process death. ```toml [dependencies] -a3s-boot = { version = "0.1.3", features = ["queue-postgres"] } +a3s-boot = { version = "0.1.4", features = ["queue-postgres"] } serde_json = "1" tokio = { version = "1", features = ["macros", "rt-multi-thread"] } ```