segment-buffer
Durable bounded queue with zstd+CBOR segment files, ack-based deletion, and filename-based crash recovery.
Extracted from monitor365, proven on 597M+ events.
Why?
There are many disk-backed queues in the Rust ecosystem, but none offer this combination:
- In-memory bounded buffer that spills to disk on a batch/interval trigger (not always write-through)
- zstd + CBOR compression for efficient storage
- Ack-based deletion — segments are removed only after the consumer confirms receipt
- Filename-based crash recovery —
lsthe directory and you see the state; no WAL, no metadata DB - Optional AES-256-GCM encryption at rest via a pluggable
SegmentCiphertrait - MPMC — multiple writers and readers via
parking_lot::Mutex
Install
# optional, for the built-in AES-256-GCM cipher:
Quickstart
use ;
use ;
let buffer = open?;
// Append items (auto-flushes at the batch threshold or flush interval)
let seq = buffer.append?;
// Read items back (from on-disk segments + in-memory pending)
let items = buffer.read_from?;
// Delete acknowledged items — a segment is removed when its end_seq <= acked_seq
let deleted = buffer.delete_acked?;
# Ok::
Encryption at rest
Enable the encryption feature and supply any SegmentCipher. The built-in
AesGcmCipher writes [12-byte nonce][ciphertext + GCM tag] per segment,
byte-compatible with monitor365 so existing encrypted segments read without
migration.
use ;
let cipher = new; // 32-byte AES-256 key
let buffer = open?;
See examples/encrypted.rs for a runnable end-to-end example.
How it works
append(item) ─► unflushed: Vec<T> (in-memory, inside the Mutex)
│
▼ batch full OR flush_interval elapsed OR flush()
take() the batch, compute start_seq/end_seq INSIDE the lock
│
▼ (lock released — mutex is never held across file I/O)
CBOR ─► zstd ─► [optional cipher.encrypt]
│
▼
write seg_*.zst.tmp ─► fsync ─► atomic rename to seg_*.zst
│
▼ (lock re-acquired)
approx_disk_bytes += len
read_from(start, limit) scans on-disk segments (sorted by start) then drains the
in-memory tail. delete_acked(seq) removes every segment whose end <= seq and
advances head_seq. Crash recovery is just: delete .tmp debris, parse the
remaining filenames. No WAL, no metadata database.
Backpressure
The crate ships metrics, not policy. store_pressure() returns
approx_disk_bytes / max_size_bytes ∈ [0.0, 1.0]; is_overloaded() is > 0.9.
You define priority thresholds — see examples/backpressure.rs.
Comparison
| Feature | segment-buffer | yaque | disk_backed_queue |
|---|---|---|---|
| Segment files | zstd+CBOR | raw bytes | SQLite |
| Ack/delete | delete_acked() |
RecvGuard commit/revert |
partial |
| Crash recovery | filename-based | replay or loss | SQLite WAL |
| Compression | zstd | none | none |
| In-memory spill | yes (batch threshold) | no (write-through) | no |
| MPMC | yes (Mutex) | SPSC only | yes |
| Encryption | optional (AES-GCM trait) | no | no |
Status
v0.1.0 — extracted from monitor365, fully decoupled, zero monitor365 dependencies. See FEATURES.md for the honest capability inventory and ROADMAP.md for direction. Changes live in CHANGELOG.md.
License
Licensed under the Apache License, Version 2.0.