segment-buffer
High-throughput local buffer for cloud sync. Single-process by design, durability-configurable, optional performant encryption, at-least-once delivery.
Extracted from monitor365 (private), proven on 597M+ events.
Why?
The cloud is the durable layer. This crate is the local throughput buffer in front of it — spool items to disk fast, drain them to your cloud endpoint at your own pace, and delete them only after the server acknowledges. When the cloud is unreachable (offline, partitioned, rate-limited), the buffer holds the backlog; when it comes back, drain resumes from where it left off.
There are many disk-backed queues in the Rust ecosystem, but none target this shape:
- Single-process by design — one owner per buffer directory; no cross-process or distributed coordination tax. An exclusive
flock-based lock atopen()enforces this (since v0.5.0). - Throughput-first, durability-configurable — pick your crash-resilience level (see Crash behavior). When the cloud is the durable copy, skip fsync; when this buffer is the last copy, fsync file + directory.
- At-least-once delivery built in —
append()returns a stable sequence number;delete_acked(seq)is the commit point. Crash before the ack and items are re-delivered on recovery. Your server-side handler MUST be idempotent on(producer_id, seq); the library delivers at-least-once, never exactly-once. - Optional performant encryption at rest — AES-256-GCM (byte-compatible with monitor365) and XChaCha20-Poly1305 (extended 24-byte nonce, no 2³²-message limit per key; constant-time in software).
SegmentConfigBuilder::recommended_cipher(key)installs XChaCha20-Poly1305 for new buffers; legacy AES-GCM segments still decrypt. Streaming/incremental cipher is a long-term direction. - zstd + CBOR segment files — efficient storage, filename-based crash recovery.
lsthe directory and you see the state; no WAL, no metadata DB.
For cross-machine replicated queues or server-side fanout, use a different tool — this crate is the producer-side local buffer.
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 last_acked_seq = 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.
#
# // Requires --features encryption. Without it, AesGcmCipher does not exist.
#
#
#
#
See examples/encrypted.rs for a runnable end-to-end example.
Cloud sync: the at-least-once drain loop
# use ;
#
#
use ;
let buf = open?;
// Producer side: append as fast as you can; the buffer handles batching.
for i in 0..10_000
buf.flush?; // ensure everything is on disk before draining
// Drain side: read a batch, send to cloud, ack what the server confirmed.
// The starting cursor comes from stats(); from then on, track it yourself.
let mut next = buf.stats.head_sequence;
loop
# Ok::
Crash semantics. If the process dies between cloud_upload and delete_acked, the batch is still on disk. On restart, read_from(next, ...) returns it again — your cloud_upload will see the same (producer, seq) pairs a second time and must treat them as no-ops. Only the unflushed in-memory tail is at risk of loss; call flush() to drain it before crash-sensitive boundaries. The library provides the sequence-number substrate; idempotency lives in your server.
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]
│
▼
prepend 8-byte SBF1 envelope ─► write seg_*.zst.tmp ─► fsync ─► atomic rename
│
▼ (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.
Crash behavior (configurable)
Shipped in v0.5.0. The default remains
Segment(today's behavior) for one release for backward compatibility; cloud-sync deployments should switch toThroughputonce the cloud endpoint holds the durable copy. Pick the policy at construction:SegmentConfig::builder().durability(DurabilityPolicy::Throughput).build().
DurabilityPolicy |
Fsync file | Fsync dir after rename | Worst-case crash loss | Use case |
|---|---|---|---|---|
Maximal |
yes | yes | last in-flight flush only | standalone queue — this buffer is the last copy |
Segment (today's default) |
yes | no | the rename window (~5–30s of flushes, kernel-dependent) | backwards-compatible default |
Throughput (for cloud sync) |
no | no | entire OS dirty window (~30s); the cloud is the durable layer | high-throughput producer with cloud ack |
Note: today's code fsyncs the segment file's data but not the directory inode after rename, so Segment already has a real (small) crash window. Maximal closes it at the cost of one extra dir.sync_all() per flush. Throughput removes the per-flush fsync entirely.
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. In cloud-sync
deployments the typical policy is: when store_pressure() exceeds your
threshold, apply backpressure to the producer (slow down, sample, drop) — the
buffer is holding the backlog until the cloud endpoint recovers.
Comparison
Comparison tables rot. This one was written against the versions current as of 2026-07; verify against the upstream crates before making a storage decision.
Reframed for the cloud-sync producer-side buffer target. Comparison tables rot; verify upstream before deciding.
| Feature | segment-buffer | yaque | disk_backed_queue |
|---|---|---|---|
| Target shape | local spool for cloud sync | general queue | general queue |
| Process model | single-process (locked at open) | SPSC | multi-process |
| Segment files | zstd+CBOR | raw bytes | SQLite |
| Ack/delete | delete_acked() (at-least-once) |
RecvGuard commit/revert |
partial |
| Crash recovery | filename-based | replay or loss | SQLite WAL |
| Compression | zstd | none | none |
| Durability knob | 3 policies (Maximal/Segment/Throughput) | write-through | SQLite full/normal/off |
| Encryption | optional (AES-GCM + XChaCha20-Poly1305) | no | no |
Status
v0.5.0 — the cloud-sync throughput batch ships the 2026-07-20 reframing: single-process throughput buffer for cloud sync, durability-configurable, XChaCha20-Poly1305 recommended cipher, at-least-once delivery. Breaking changes are batched so users upgrade once — see CHANGELOG.md for the full per-item detail and migration notes. See FEATURES.md, ROADMAP.md.
v0.4.2 — the "process debt + semver-leak closure" release. Gates
fuzz_hooks behind a #[cfg] feature (closes the v0.4.1 semver leak), adds
a CI loom job (prevents the v0.4.0-v0.4.1 silent rot of the loom test),
adds 1 new fuzz target (fuzz_append_all) and 2 new property tests, and
ships docs/DOMAIN_LANGUAGE.md + docs/CIPHERS.md. Non-breaking; drop-in
upgrade from v0.4.1.
See CHANGELOG.md for details.
See FEATURES.md, ROADMAP.md.
Performance: the 2026-07-20 PGO session
(docs/perf/2026-07-20_hot-path-flamegraph.md)
found that ~66% of flush CPU was in zstd re-initialising its ~200 KB CCtx
on every encode_all call. Pooling a zstd::bulk::Compressor on SegmentBuffer
made append/batch_1 2.07× faster (15.09 µs → 7.75 µs), with smaller wins
at larger batches. The crate is now substantially faster than v0.1.0 on small
batches; the previous "30–65% regression" framing is obsolete. A Throughput
durability policy (above) would remove the per-flush fsync and open a further
large gain on cloud-sync workloads.
Methodology caveat: single-run, single-machine criterion medians without
statistical noise bars — indicative of direction, not publication-grade. See
docs/PERFORMANCE.md for methodology and reproduction.
License
Licensed under the Apache License, Version 2.0.