# persistent-queue
[](https://github.com/rolandjitsu/persistent-queue/actions/workflows/ci.yml)
[](https://codecov.io/gh/rolandjitsu/persistent-queue)
[](https://crates.io/crates/persistent-queue)
[](https://docs.rs/persistent-queue)
[](./LICENSE)
A durable, at-least-once MPSC queue backed by in-memory and durable backends.
Many producers push byte payloads; a single consumer reserves each one and holds it
in flight until it acks (removes) or nacks (returns) it. A dropped reservation, a
panic, or a crash all put the item back, so nothing is lost. Storage is a pluggable
`Store` - in-memory by default, `sled` or `redb` behind features, or your own - and
the core is synchronous and runtime-agnostic, with an optional `tokio` facade. See
`DESIGN.md` for the on-disk layout,
cursors, and crash recovery.
## Usage
```rust
use persistent_queue::{Builder, MemStore};
let (tx, rx) = Builder::new(MemStore::new()).capacity(1024).open().unwrap();
tx.push(b"job").unwrap(); // waits if the queue is at capacity
if let Some(item) = rx.reserve().unwrap() {
assert_eq!(&*item, b"job"); // derefs to the bytes
item.ack().unwrap(); // remove it; or item.nack() to retry later
}
```
For durability, use an on-disk backend:
```rust,ignore
use persistent_queue::{Builder, SledStore};
let store = SledStore::open("/var/lib/myapp/queue").unwrap();
let (tx, rx) = Builder::new(store).capacity(1024).open().unwrap();
```
`RedbStore::open` works the same way behind the `redb` feature.
For typed messages, open a typed queue with a `Codec` - serde and bincode behind the
`serde` feature, or your own:
```rust
use persistent_queue::{Bincode, Builder, MemStore};
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize)]
struct Job {
id: u64,
name: String,
}
let (tx, rx) = Builder::new(MemStore::new()).open_typed(Bincode).unwrap();
tx.push(&Job { id: 1, name: "build".into() }).unwrap();
if let Some(item) = rx.reserve().unwrap() {
println!("{}", item.name); // derefs to the decoded value
item.ack().unwrap();
}
```
With the `tokio` feature, `open_async` gives async handles: store I/O runs on tokio's
blocking pool, and waiting (for capacity, or the next item) is async - a blocked
`push` or `reserve` costs a task, not a thread:
```rust,ignore
use persistent_queue::{Builder, MemStore};
let (tx, rx) = Builder::new(MemStore::new()).open_async().await.unwrap();
tx.push(b"job".to_vec()).await.unwrap();
if let Some(item) = rx.reserve().await.unwrap() {
item.ack().await.unwrap();
}
```
## Backends
| `MemStore` | (default) | No | Zero dependencies; tests and baseline. |
| `SledStore` | `sled` | Yes | On-disk, backed by sled. |
| `RedbStore` | `redb` | Yes | On-disk, backed by redb. |
| `RocksStore` | `rocksdb` | Yes | On-disk, backed by RocksDB (bundled C++). |
Implement `Store` yourself for any other key/value store.
- **One process per store.** `sled`, `redb`, and `rocksdb` take an exclusive lock on
their files, so only one process can open a given database at a time.
- **Corruption is surfaced, not repaired.** A backend error on open (including a
corrupt store) is returned as `OpenError::Store`; the queue does not auto-repair or
discard data. A store written by a newer on-disk format is rejected with
`OpenError::UnsupportedVersion`.
- **RocksDB is a bundled C++ library.** The `rocksdb` feature compiles it from source
and statically links it into your binary (nothing extra to ship at runtime), but it
needs a C++ toolchain to build, adds compile time and binary size, and links the C++
runtime - so a fully static (musl) build is difficult. For a small, pure-Rust, fully
static binary, use `sled` or `redb`.
## Guarantees
- **Durability.** Once `push` returns under a durable policy, the item survives a
crash.
- **At-least-once.** Every pushed item is delivered at least once. A crash between
handling an item and its durable ack redelivers it. True exactly-once is not
possible from the queue alone; `Reserved::seq()` is a stable id you can dedupe on
for effectively-once.
- **Order.** FIFO by sequence number; `reserve` hands items out oldest first.
Durability is a policy on the builder: `Sync` (fsync every push and ack), `Group`
(batch concurrent pushes behind one fsync), or `None` (no fsync - fastest, but recent
items can be lost on a crash). `MemStore` never persists, regardless of policy.
`durable_acks(false)` drops the ack fsync - at-least-once still holds, a lost ack just
redelivers - roughly halving the per-message fsync cost under `Sync`. Capacity bounds
the unacked item count (`capacity`) and, optionally, their total bytes (`max_bytes`);
`push` blocks when either is reached.
What survives a crash (process exit or power loss):
| `MemStore` | nothing (in-memory only) | nothing (in-memory only) |
| `sled`, `redb`, `rocksdb` | not guaranteed (no fsync) | survives (fsync'd) |
## Comparison
How persistent-queue relates to some common alternatives - it is durable,
at-least-once, multi-producer, and backend-agnostic:
| `persistent-queue` | Yes | MPSC | in-memory, sled, redb, or your own |
| `yaque` | Yes | SPSC | files |
| `sled`, `redb` | Yes | not a queue | own file (key/value store) |
| `flume`, `crossbeam` | No | MPMC | in-memory |
`yaque` is a persistent queue but single-producer / single-consumer (its docs: "an
SPSC channel using your OS' filesystem"). `sled` and `redb` are key/value stores,
not queues - persistent-queue is built on top of them. `flume` and `crossbeam` are
fast in-memory channels with no durability. Backends and semantics change over time;
check each crate's current docs.
## Benchmarks
Throughput across backends, durability policies, and producer counts is in
[BENCHMARKS.md](./BENCHMARKS.md). Short version: the crate's own overhead is
in-memory-fast, and durability cost is the backend's `fsync`. Reproduce with
`cargo bench --all-features`, and measure on your own hardware.
## Roadmap
- **Richer delivery**: multiple / competing consumers with visibility timeouts,
dead-letter handling after N redeliveries, priorities and delayed delivery.
- **An rkyv codec with zero-copy reads**: add an `rkyv` codec and return borrowed or
`Bytes` data from the store, so a reserve reads the archived value without a copy.
## Status
Early, single-maintainer software. The surface is intentionally small: a producer,
a consumer, the `Reserved` guard, and the `Store` trait. Contributions welcome.
## License
[Apache-2.0](./LICENSE).