pub struct SegmentBuffer<T> { /* private fields */ }Expand description
High-throughput local buffer for cloud sync, holding items of T in
memory and spilling them to compressed segment files for at-least-once
delivery to a cloud endpoint.
Thread-safe via parking_lot::Mutex. All file I/O is synchronous. The mutex
is never held across an async boundary because there are no await points.
Create with SegmentBuffer::open, supplying the directory and config.
Implementations§
Source§impl<T> SegmentBuffer<T>
impl<T> SegmentBuffer<T>
Sourcepub fn open(dir: impl Into<PathBuf>, config: SegmentConfig) -> Result<Self>
pub fn open(dir: impl Into<PathBuf>, config: SegmentConfig) -> Result<Self>
Open (or create) a buffer at dir, recovering from any existing
segment files.
Recovery is filename-based: it scans the directory to rebuild
head_seq / next_seq and deletes leftover .tmp debris. Segment
contents are not read until read_from, so a
corrupted segment does not fail here — it fails when read.
If you need the recovery summary (segments found, bytes, head/next seq)
programmatically, use SegmentBuffer::open_with_report instead. The
same data is logged via tracing::info! from this call.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;§Errors
Returns SegmentError::Io if the directory cannot be created or read.
Sourcepub fn open_with_report(
dir: impl Into<PathBuf>,
config: SegmentConfig,
) -> Result<(Self, RecoveryReport)>
pub fn open_with_report( dir: impl Into<PathBuf>, config: SegmentConfig, ) -> Result<(Self, RecoveryReport)>
Like SegmentBuffer::open, but also returns a RecoveryReport
describing what the recovery scan found on disk.
Useful for operational dashboards or migration tools that need to know the on-disk state without re-scanning.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let (buf, report) =
SegmentBuffer::<u64>::open_with_report(dir.path(), SegmentConfig::default())?;
assert_eq!(report.segment_count, 0); // fresh dir
assert_eq!(report.head_seq, 0);
assert_eq!(report.next_seq, 0);§Errors
Returns SegmentError::Io if the directory cannot be created or read.
Returns SegmentError::Locked if another process holds the
exclusive single-process lock on <dir>/.segment-buffer.lock.
Sourcepub fn append(&self, event: T) -> Result<u64>
pub fn append(&self, event: T) -> Result<u64>
Append an item to the buffer. Assigns the next sequence number and auto-flushes if the batch threshold or interval is reached.
Returns the assigned sequence number. The first append returns 0,
and the number increments by 1 for each subsequent append.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert_eq!(buf.append(1)?, 0);
assert_eq!(buf.append(2)?, 1);
assert_eq!(buf.append(3)?, 2);Sourcepub fn flush(&self) -> Result<()>
pub fn flush(&self) -> Result<()>
Flush buffered items to a segment file. No-op if nothing is buffered.
Flushing is also triggered automatically by append
according to the configured FlushPolicy (batch threshold, interval,
both, or manual). Call this explicitly when you need durability before
a known threshold, or when using FlushPolicy::Manual.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
buf.append(1)?;
buf.append(2)?;
buf.flush()?; // items now durable on disk
assert_eq!(buf.pending_count(), 2);Sourcepub fn read_from(&self, start_seq: u64, limit: usize) -> Result<Vec<T>>
pub fn read_from(&self, start_seq: u64, limit: usize) -> Result<Vec<T>>
Read up to limit items starting from start_seq (inclusive).
Reads from both on-disk segment files and in-memory pending items. Items are returned in ascending sequence order.
Passing limit = 0 returns an empty Vec without scanning.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
buf.append(10)?;
buf.append(20)?;
buf.append(30)?;
buf.flush()?;
let items = buf.read_from(0, 100)?;
assert_eq!(items, vec![10, 20, 30]);
// start_seq skips already-read items:
let tail = buf.read_from(2, 100)?;
assert_eq!(tail, vec![30]);Sourcepub fn for_each_from<F>(
&self,
start_seq: u64,
limit: usize,
f: F,
) -> Result<usize>
pub fn for_each_from<F>( &self, start_seq: u64, limit: usize, f: F, ) -> Result<usize>
Lending-iterator counterpart to read_from: invoke
f(seq, item) for up to limit items starting at start_seq, without
materialising them into a Vec<T>.
This avoids the per-item Clone that read_from
pays for in-memory pending items. On-disk segments still deserialize
into a temporary Vec<T> per segment (the on-disk format is bytes, not
T), but items are passed to f by reference rather than being
re-collected.
Returns the number of items the callback was invoked for.
§Performance
Micro-benchmarked in benches/bench_read_vs_for_each.rs against
in-memory pending items (no segment files):
| Items | read_from | for_each_from | Speedup |
|---|---|---|---|
| 1,000 | ~26 µs | ~1.2 µs | ~21× |
| 10,000 | ~200 µs | ~10 µs | ~20× |
The speedup shrinks toward zero once on-disk segments dominate, because both paths pay the same CBOR+zstd+cipher decode cost per segment — the clone saving only applies to the in-memory tail.
§Re-entrancy contract
The buffer mutex is held across f while iterating the in-memory
pending items. Calling any other &self method on SegmentBuffer
from inside f would deadlock (parking_lot::Mutex is not reentrant).
To make this footgun impossible to hit silently, every other method
asserts it is not being re-entered from inside a for_each_from
callback and panics with a clear message if it is. The callback
receives only (seq, &T), which gives no way to reach the buffer, but
a closure that captures a clone of the Arc<SegmentBuffer<T>> can
still attempt re-entry — and will now get an immediate, diagnosable
panic instead of a silent hang.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
for i in 0..5u64 {
buf.append(i * 10)?;
}
buf.flush()?;
let mut sum = 0u64;
let count = buf.for_each_from(0, 100, |_seq, item| { sum += *item; })?;
assert_eq!(count, 5);
assert_eq!(sum, 0 + 10 + 20 + 30 + 40);Sourcepub fn delete_acked(&self, acked_seq: u64) -> Result<usize>
pub fn delete_acked(&self, acked_seq: u64) -> Result<usize>
Delete all on-disk segment files whose items are fully covered by
acked_seq.
A segment is deleted when its end_seq <= acked_seq. Returns the number
of segment files removed.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
for i in 0..5u64 {
buf.append(i)?;
}
buf.flush()?;
// Consumer has processed sequence 0..=4; acknowledge them:
let removed = buf.delete_acked(4)?;
assert_eq!(removed, 1); // one segment file deleted
assert_eq!(buf.pending_count(), 0);§Limitation
Acknowledgement only removes flushed segment files. Items still held
in the in-memory pending batch have no segment file to delete, so they
remain readable (and counted by SegmentBuffer::pending_count) until
they are flushed and acknowledged in a later call. head_seq is clamped
so it never advances past the pending window, keeping the backlog count
honest.
Sourcepub fn latest_sequence(&self) -> u64
pub fn latest_sequence(&self) -> u64
The highest sequence number assigned (or 0 if buffer is empty).
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert_eq!(buf.latest_sequence(), 0);
buf.append(7)?;
assert_eq!(buf.latest_sequence(), 0);
buf.append(8)?;
assert_eq!(buf.latest_sequence(), 1);Sourcepub fn pending_count(&self) -> u64
pub fn pending_count(&self) -> u64
Total items waiting in the buffer (on-disk + in-memory pending).
Equivalent to latest_sequence() - head_seq + 1 when non-empty, 0 when
empty. Decreases as delete_acked removes files.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert_eq!(buf.pending_count(), 0);
buf.append(1)?;
buf.append(2)?;
assert_eq!(buf.pending_count(), 2);
buf.flush()?;
assert_eq!(buf.pending_count(), 2); // still pending until acked
buf.delete_acked(1)?;
assert_eq!(buf.pending_count(), 0);Sourcepub fn len(&self) -> u64
pub fn len(&self) -> u64
Standard len alias for pending_count.
Provided so SegmentBuffer reads like a normal collection at the call
site (buf.len(), buf.is_empty()). Same value as pending_count(),
kept as u64 because the buffer is proven beyond usize::MAX on
32-bit targets (597M+ events in monitor365).
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert!(buf.is_empty());
buf.append(7)?;
assert_eq!(buf.len(), 1);
assert!(!buf.is_empty());Sourcepub fn is_empty(&self) -> bool
pub fn is_empty(&self) -> bool
true when there are no items waiting in the buffer (on-disk or
in-memory). Equivalent to pending_count() == 0.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert!(buf.is_empty());Sourcepub fn store_pressure(&self) -> f32
pub fn store_pressure(&self) -> f32
Disk usage pressure as a value between 0.0 and 1.0.
Use this to implement your own admission/backpressure policy (e.g.
reject low-priority items above 0.90, reject standard items above 0.95).
Returns 0.0 when max_size_bytes == 0 (limit disabled).
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let mut cfg = SegmentConfig::default();
cfg.max_size_bytes = 1000; // tiny limit so pressure is observable
let buf: SegmentBuffer<u64> = SegmentBuffer::open(dir.path(), cfg)?;
assert!(buf.store_pressure() < 0.1);Sourcepub fn is_overloaded(&self) -> bool
pub fn is_overloaded(&self) -> bool
True when disk usage exceeds 90% of the configured limit.
Convenience wrapper around store_pressure() > 0.9.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert!(!buf.is_overloaded());Sourcepub fn stats(&self) -> BufferStats
pub fn stats(&self) -> BufferStats
Capture a consistent snapshot of buffer state under a single lock.
Cheaper and more consistent than calling
pending_count,
latest_sequence,
store_pressure etc. individually (which each
take the mutex and could observe a flush/delete between calls).
§Performance
Micro-benchmarked in benches/bench_stats.rs (run with
cargo bench --bench bench_stats --features encryption):
| Operation | Measured time (median, typical run) |
|---|---|
stats() (single lock, 7-field snapshot) | ~12 ns |
3 individual accessors (pending_count + latest_sequence + store_pressure) | ~31 ns |
So stats() is roughly 2.5× cheaper than 3 individual accessors
while also being atomic — torn reads between calls are impossible.
Numbers are from the benchmark machine and fluctuate with hardware;
the relative ratio is the durable claim.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
buf.append(1)?;
buf.append(2)?;
let snapshot = buf.stats();
assert_eq!(snapshot.pending_count, 2);
assert_eq!(snapshot.next_sequence, 2);
assert!(snapshot.store_pressure < 0.01);Sourcepub fn path(&self) -> &Path
pub fn path(&self) -> &Path
The directory this buffer reads from and writes segment files to.
Useful for operators that need to inspect, archive, or quarantine the
segment directory without parsing it out of Debug.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
assert_eq!(buf.path(), dir.path());Sourcepub fn config(&self) -> &SegmentConfig
pub fn config(&self) -> &SegmentConfig
The SegmentConfig this buffer was opened with.
Returned by reference so callers can inspect the flush policy, disk ceiling, compression level, and cipher presence without re-deriving them. The config is immutable for the lifetime of the buffer.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig, FlushPolicy};
use tempfile::tempdir;
let dir = tempdir()?;
let config = SegmentConfig::builder()
.flush_at_batch_size(128)
.build();
let buf: SegmentBuffer<u64> = SegmentBuffer::open(dir.path(), config)?;
match &buf.config().flush_policy {
FlushPolicy::Batch(n) => println!("flushing at {n} items"),
_ => {}
}Sourcepub fn sync_disk_bytes(&self) -> Result<u64>
pub fn sync_disk_bytes(&self) -> Result<u64>
Re-stat the segment directory and store the authoritative total as
BufferStats::approx_disk_bytes.
BufferStats::approx_disk_bytes is updated incrementally on every
flush/delete/recover, so it is accurate as long as only this buffer
touches the directory. If an external process (backup, compaction,
manual cleanup) adds or removes segment files, the cached value drifts.
This method recomputes it from a directory scan.
Returns the new total so callers can observe the delta without a
second call to stats.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
buf.append(1)?;
buf.flush()?;
// Simulate an external process truncating a segment file to zero bytes.
for entry in std::fs::read_dir(dir.path())? {
let _ = std::fs::write(entry?.path(), b"");
}
let synced = buf.sync_disk_bytes()?;
assert_eq!(synced, 0, "external truncation should be reflected");§Errors
Returns SegmentError::Io if the directory cannot be read.
Sourcepub fn append_all<I>(&self, items: I) -> Result<u64>where
I: IntoIterator<Item = T>,
pub fn append_all<I>(&self, items: I) -> Result<u64>where
I: IntoIterator<Item = T>,
Append a batch of items under a single lock acquisition.
Each item receives the next contiguous sequence number. Returns the
last sequence number assigned (matching the contract of
append); the full range is
[last - count + 1, last] where count is the number of items the
iterator yielded.
§Batch vs streaming semantics
All items are accumulated under a single lock acquisition, then the
flush policy is checked once at the end. This gives true atomic
batch semantics: either the entire batch lands in the buffer or the
error propagates. Callers who want per-item auto-flush semantics
(flush at every batch_size threshold) should call
append in a loop instead — append_all is
optimized for the “load this batch atomically” use case and avoids
paying the lock-acquisition cost per item.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig, FlushPolicy};
use tempfile::tempdir;
let dir = tempdir()?;
let config = SegmentConfig::builder()
.flush_policy(FlushPolicy::Manual)
.build();
let buf: SegmentBuffer<u64> = SegmentBuffer::open(dir.path(), config)?;
let last = buf.append_all([10u64, 20, 30, 40])?;
assert_eq!(last, 3); // 0-based: items got seqs 0, 1, 2, 3
assert_eq!(buf.pending_count(), 4);§Errors
Returns SegmentError::Io if a flush triggered by the batch fails.
Sourcepub fn iter_from(
&self,
start_seq: u64,
limit: usize,
) -> Result<SegmentIter<'_, T>>
pub fn iter_from( &self, start_seq: u64, limit: usize, ) -> Result<SegmentIter<'_, T>>
Owned-item iterator over buffer contents starting at start_seq.
Equivalent to read_from but yields (seq, item)
pairs one at a time so callers can write for (seq, item) in buf.iter_from(start, limit)? and chain standard
Iterator combinators (.take, .filter, .map, …).
This is a materialising iterator: items are loaded eagerly up to
limit (memory cost O(limit)). For a lending iterator that
passes in-memory items by reference without cloning, use
for_each_from — that variant is ~20× faster
on the in-memory tail but takes a closure instead of returning an
Iterator. The two coexist because no stable-Rust Iterator
trait can currently express “yield &T from &mut self” without
pre-collecting.
§Re-entrancy contract
The iterator borrows the buffer for 'a. Drop the iterator before
calling any other &self method on the same buffer; if the
iterator is alive across a flush / append / delete_acked
call, that call will panic with a clear message (same contract as
for_each_from). The simplest pattern is
for item in buf.iter_from(..)? { ... } — the for loop drops the
iterator at the end of the block.
§Example
use segment_buffer::{SegmentBuffer, SegmentConfig};
use tempfile::tempdir;
let dir = tempdir()?;
let buf: SegmentBuffer<u64> =
SegmentBuffer::open(dir.path(), SegmentConfig::default())?;
for i in 0..5u64 { buf.append(i * 10)?; }
buf.flush()?;
// `for` loop with owned items + seq numbers:
let mut seen = Vec::new();
for (seq, item) in buf.iter_from(0, 100)? {
seen.push((seq, item));
}
assert_eq!(seen, vec![
(0, 0), (1, 10), (2, 20), (3, 30), (4, 40),
]);§Errors
Returns SegmentError if the directory scan or any segment decode
fails.
Trait Implementations§
Source§impl<T> Debug for SegmentBuffer<T>
Debug mirrors the field set of BufferStats plus the directory path.
It does NOT print the in-memory unflushed items (which could be large or
sensitive), so T itself is not required to be Debug.
impl<T> Debug for SegmentBuffer<T>
Debug mirrors the field set of BufferStats plus the directory path.
It does NOT print the in-memory unflushed items (which could be large or
sensitive), so T itself is not required to be Debug.
Source§impl<T> Drop for SegmentBuffer<T>
impl<T> Drop for SegmentBuffer<T>
Source§fn drop(&mut self)
fn drop(&mut self)
Releases the single-process flock by explicitly calling unlock and
then dropping the lock file handle. The kernel would release the
advisory lock on fd close anyway, but the explicit call makes the
release point diagnosable in a flamegraph (vs. waiting for File’s
own Drop to run somewhere in the field-tear-down sequence).
Deliberately no T: Serialize + ... bound: Drop impls must match
the struct’s bounds (Rust rule E0367), and the struct itself has no
bounds — the bound lives on the API-impl block. The lock-release
logic doesn’t touch T at all, so no bound is needed here.