Skip to main content

SegmentBuffer

Struct SegmentBuffer 

Source
pub struct SegmentBuffer<T> { /* private fields */ }
Expand description

Durable bounded queue of T backed by compressed segment files.

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>
where T: Serialize + DeserializeOwned + Clone + Send + 'static,

Source

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.

Source

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.

Source

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);
Source

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);
Source

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]);
Source

pub fn for_each_from<F>( &self, start_seq: u64, limit: usize, f: F, ) -> Result<usize>
where F: FnMut(u64, &T),

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):

Itemsread_fromfor_each_fromSpeedup
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);
Source

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.

Source

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);
Source

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);
Source

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());
Source

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());
Source

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);
Source

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());
Source

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):

OperationMeasured 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);
Source

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());
Source

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"),
    _ => {}
}
Source

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.

Source

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.

Trait Implementations§

Source§

impl<T> Debug for SegmentBuffer<T>
where T: Serialize + DeserializeOwned + Clone + Send + 'static,

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§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl<T> !Freeze for SegmentBuffer<T>

§

impl<T> !RefUnwindSafe for SegmentBuffer<T>

§

impl<T> !UnwindSafe for SegmentBuffer<T>

§

impl<T> Send for SegmentBuffer<T>
where T: Send,

§

impl<T> Sync for SegmentBuffer<T>
where T: Send,

§

impl<T> Unpin for SegmentBuffer<T>
where T: Unpin,

§

impl<T> UnsafeUnpin for SegmentBuffer<T>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more