Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "a3s-boot"
version = "0.1.3"
version = "0.1.4"
edition = "2021"
authors = ["A3S Lab"]
license = "MIT"
Expand Down Expand Up @@ -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"]
Expand All @@ -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 }
Expand Down Expand Up @@ -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]
Expand Down
66 changes: 60 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand All @@ -133,31 +134,32 @@ 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

### Installation

```toml
[dependencies]
a3s-boot = "0.1.3"
a3s-boot = "0.1.4"
tokio = { version = "1", features = ["macros", "rt-multi-thread"] }
```

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"] }
```
Expand Down Expand Up @@ -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.4", 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:
Expand Down
14 changes: 10 additions & 4 deletions ROADMAP.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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`,
Expand Down
9 changes: 6 additions & 3 deletions src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading