Skip to main content

segment_buffer/
lib.rs

1//! Durable bounded queue backed by zstd-compressed CBOR segment files.
2//!
3//! Items are accumulated in memory, flushed as zstd-compressed CBOR batches
4//! to `seg_{start:012}_{end:012}.zst` files, and deleted once the consumer
5//! acknowledges receipt via [`SegmentBuffer::delete_acked`].
6//!
7//! The buffer is generic over any `T: Serialize + DeserializeOwned + Clone + Send + 'static`.
8//! Crash recovery is filename-based: scanning the directory rebuilds `head_seq`
9//! and `next_seq` without any WAL or metadata database.
10//!
11//! # Example
12//!
13//! ```no_run
14//! use segment_buffer::{SegmentBuffer, SegmentConfig};
15//! use serde::{Serialize, Deserialize};
16//!
17//! #[derive(Serialize, Deserialize, Clone)]
18//! struct MyItem { id: u64 }
19//!
20//! let buffer = SegmentBuffer::<MyItem>::open("/tmp/my-queue", SegmentConfig::default())?;
21//! let seq = buffer.append(MyItem { id: 1 })?;
22//! let items = buffer.read_from(0, 100)?;
23//! # Ok::<(), Box<dyn std::error::Error>>(())
24//! ```
25
26#![warn(missing_docs)]
27
28mod cipher;
29mod error;
30mod segment;
31
32#[cfg(feature = "encryption")]
33pub use cipher::AesGcmCipher;
34pub use cipher::SegmentCipher;
35pub use error::{Result, SegmentError};
36
37use std::fs;
38use std::path::PathBuf;
39use std::time::Instant;
40
41use parking_lot::Mutex;
42use serde::de::DeserializeOwned;
43use serde::Serialize;
44use tracing::{debug, info};
45
46use segment::SegmentRange;
47
48/// Configuration knobs for [`SegmentBuffer`].
49pub struct SegmentConfig {
50    /// Max events accumulated in RAM before auto-flush (default: 256).
51    pub max_batch_events: usize,
52    /// Max seconds between flushes. An append after this interval triggers a
53    /// flush even if the batch threshold hasn't been reached (default: 5s).
54    pub flush_interval_secs: u64,
55    /// Max total disk usage before the buffer reports overload pressure (default: 10 GB).
56    pub max_size_bytes: u64,
57    /// zstd compression level (1-22; 3 is fast with a good ratio).
58    pub compression_level: i32,
59    /// Optional cipher for encrypting segment files at rest. When `None`,
60    /// segments are written as plaintext zstd+CBOR.
61    pub cipher: Option<Box<dyn SegmentCipher>>,
62}
63
64impl std::fmt::Debug for SegmentConfig {
65    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
66        f.debug_struct("SegmentConfig")
67            .field("max_batch_events", &self.max_batch_events)
68            .field("flush_interval_secs", &self.flush_interval_secs)
69            .field("max_size_bytes", &self.max_size_bytes)
70            .field("compression_level", &self.compression_level)
71            .field("cipher", &self.cipher.as_ref().map(|_| "[set]"))
72            .finish()
73    }
74}
75
76impl Default for SegmentConfig {
77    fn default() -> Self {
78        Self {
79            max_batch_events: 256,
80            flush_interval_secs: 5,
81            max_size_bytes: 10 * 1024 * 1024 * 1024,
82            compression_level: 3,
83            cipher: None,
84        }
85    }
86}
87
88struct BufferInner<T> {
89    /// Items buffered in memory, not yet written to a segment file. Drained by
90    /// [`SegmentBuffer::flush`] and rebuilt empty on crash recovery (unflushed
91    /// items do not survive a crash by design).
92    unflushed: Vec<T>,
93    next_seq: u64,
94    head_seq: u64,
95    last_flush: Instant,
96    approx_disk_bytes: u64,
97}
98
99/// Durable bounded queue of `T` backed by compressed segment files.
100///
101/// Thread-safe via `parking_lot::Mutex`. All file I/O is synchronous. The mutex
102/// is never held across an async boundary because there are no await points.
103///
104/// Create with [`SegmentBuffer::open`], supplying the directory and config.
105pub struct SegmentBuffer<T> {
106    dir: PathBuf,
107    config: SegmentConfig,
108    inner: Mutex<BufferInner<T>>,
109}
110
111impl<T> SegmentBuffer<T>
112where
113    T: Serialize + DeserializeOwned + Clone + Send + 'static,
114{
115    /// Open (or create) a buffer at `dir`, recovering from any existing
116    /// segment files.
117    ///
118    /// Recovery is **filename-based**: it scans the directory to rebuild
119    /// `head_seq` / `next_seq` and deletes leftover `.tmp` debris. Segment
120    /// *contents* are not read until [`read_from`](Self::read_from), so a
121    /// corrupted segment does not fail here — it fails when read.
122    ///
123    /// # Errors
124    ///
125    /// Returns [`SegmentError::Io`] if the directory cannot be created or read.
126    pub fn open(dir: impl Into<PathBuf>, config: SegmentConfig) -> Result<Self> {
127        let dir = dir.into();
128        fs::create_dir_all(&dir)?;
129
130        let buffer = Self {
131            dir,
132            config,
133            inner: Mutex::new(BufferInner {
134                unflushed: Vec::new(),
135                next_seq: 0,
136                head_seq: 0,
137                last_flush: Instant::now(),
138                approx_disk_bytes: 0,
139            }),
140        };
141
142        buffer.recover()?;
143        Ok(buffer)
144    }
145
146    // -----------------------------------------------------------------------
147    // Public API
148    // -----------------------------------------------------------------------
149
150    /// Append an item to the buffer. Assigns the next sequence number and
151    /// auto-flushes if the batch threshold or interval is reached.
152    ///
153    /// Returns the assigned sequence number.
154    pub fn append(&self, event: T) -> Result<u64> {
155        let (should_flush, seq) = {
156            let mut inner = self.inner.lock();
157            inner.unflushed.push(event);
158            inner.next_seq += 1;
159            let seq = inner.next_seq - 1;
160
161            let batch_full = inner.unflushed.len() >= self.config.max_batch_events;
162            let interval_elapsed =
163                inner.last_flush.elapsed().as_secs() >= self.config.flush_interval_secs;
164            (batch_full || interval_elapsed, seq)
165        };
166
167        if should_flush {
168            self.flush()?;
169        }
170
171        Ok(seq)
172    }
173
174    /// Flush buffered items to a segment file. No-op if nothing is buffered.
175    pub fn flush(&self) -> Result<()> {
176        let (events, start_seq, end_seq) = {
177            let mut inner = self.inner.lock();
178            inner.last_flush = Instant::now();
179            if inner.unflushed.is_empty() {
180                return Ok(());
181            }
182            let events = std::mem::take(&mut inner.unflushed);
183            let count = events.len() as u64;
184            let end_seq = inner.next_seq - 1;
185            let start_seq = end_seq + 1 - count;
186            (events, start_seq, end_seq)
187        };
188
189        let compressed_len = self.write_segment(start_seq, end_seq, &events)?;
190
191        {
192            let mut inner = self.inner.lock();
193            inner.approx_disk_bytes += compressed_len;
194        }
195
196        debug!(start_seq, end_seq, count = events.len(), "Flushed segment");
197        Ok(())
198    }
199
200    /// Read up to `limit` items starting from `start_seq` (inclusive).
201    ///
202    /// Reads from both on-disk segment files and in-memory pending items.
203    /// Items are returned in ascending sequence order.
204    pub fn read_from(&self, start_seq: u64, limit: usize) -> Result<Vec<T>> {
205        if limit == 0 {
206            return Ok(Vec::new());
207        }
208
209        let mut result: Vec<T> = Vec::with_capacity(limit.min(1024));
210
211        // Phase 1: read from on-disk segments.
212        let segments = self.scan_segments()?;
213        for seg in &segments {
214            if result.len() >= limit {
215                break;
216            }
217            if seg.end < start_seq {
218                continue;
219            }
220
221            let events = self.read_segment(*seg)?;
222            let skip = if seg.start < start_seq {
223                (start_seq - seg.start) as usize
224            } else {
225                0
226            };
227
228            for event in events.into_iter().skip(skip) {
229                if result.len() >= limit {
230                    break;
231                }
232                result.push(event);
233            }
234        }
235
236        // Phase 2: read from in-memory pending events.
237        if result.len() < limit {
238            let inner = self.inner.lock();
239            let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
240            for (i, event) in inner.unflushed.iter().enumerate() {
241                let seq = pending_start + i as u64;
242                if seq < start_seq {
243                    continue;
244                }
245                if result.len() >= limit {
246                    break;
247                }
248                result.push(event.clone());
249            }
250        }
251
252        Ok(result)
253    }
254
255    /// Delete all on-disk segment files whose items are fully covered by
256    /// `acked_seq`.
257    ///
258    /// A segment is deleted when its `end_seq <= acked_seq`. Returns the number
259    /// of segment files removed.
260    ///
261    /// # Limitation
262    ///
263    /// Acknowledgement only removes **flushed** segment files. Items still held
264    /// in the in-memory pending batch have no segment file to delete, so they
265    /// remain readable (and counted by [`SegmentBuffer::pending_count`]) until
266    /// they are flushed and acknowledged in a later call. `head_seq` is clamped
267    /// so it never advances past the pending window, keeping the backlog count
268    /// honest.
269    pub fn delete_acked(&self, acked_seq: u64) -> Result<usize> {
270        let segments = self.scan_segments()?;
271        let mut deleted = 0;
272        let mut freed_bytes: u64 = 0;
273        let mut new_head = None;
274
275        for seg in &segments {
276            if seg.end <= acked_seq {
277                let path = self.segment_path(seg.start, seg.end);
278                if let Ok(meta) = fs::metadata(&path) {
279                    freed_bytes += meta.len();
280                }
281                match fs::remove_file(&path) {
282                    Ok(()) => {
283                        deleted += 1;
284                        debug!(start = seg.start, end = seg.end, "Deleted acked segment");
285                    }
286                    Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
287                    Err(e) => return Err(e.into()),
288                }
289            } else if new_head.is_none() {
290                new_head = Some(seg.start);
291            }
292        }
293
294        {
295            let mut inner = self.inner.lock();
296            inner.approx_disk_bytes = inner.approx_disk_bytes.saturating_sub(freed_bytes);
297            // `head_seq` tracks the oldest unacked sequence. Clamp it to the
298            // start of the in-memory pending window: items still waiting to be
299            // flushed cannot be acknowledged (there is no segment file to
300            // delete), so head_seq must not advance past them. Without this
301            // clamp, acknowledging past a buffer that still holds unflushed
302            // items would make `pending_count` under-report the real backlog.
303            let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
304            inner.head_seq = new_head.unwrap_or(inner.next_seq).min(pending_start);
305        }
306
307        if deleted > 0 {
308            info!(deleted, freed_bytes, acked_seq, "Deleted acked segments");
309        }
310
311        Ok(deleted)
312    }
313
314    /// The highest sequence number assigned (or 0 if buffer is empty).
315    pub fn latest_sequence(&self) -> u64 {
316        let inner = self.inner.lock();
317        if inner.next_seq == 0 {
318            0
319        } else {
320            inner.next_seq - 1
321        }
322    }
323
324    /// Total items waiting in the buffer (on-disk + in-memory pending).
325    pub fn pending_count(&self) -> u64 {
326        let inner = self.inner.lock();
327        inner.next_seq.saturating_sub(inner.head_seq)
328    }
329
330    /// Disk usage pressure as a value between 0.0 and 1.0.
331    ///
332    /// Use this to implement your own admission/backpressure policy (e.g.
333    /// reject low-priority items above 0.90, reject standard items above 0.95).
334    pub fn store_pressure(&self) -> f32 {
335        let inner = self.inner.lock();
336        if self.config.max_size_bytes == 0 {
337            return 0.0;
338        }
339        (inner.approx_disk_bytes as f32 / self.config.max_size_bytes as f32).min(1.0)
340    }
341
342    /// True when disk usage exceeds 90% of the configured limit.
343    pub fn is_overloaded(&self) -> bool {
344        self.store_pressure() > 0.9
345    }
346
347    // -----------------------------------------------------------------------
348    // Internal helpers
349    // -----------------------------------------------------------------------
350
351    fn recover(&self) -> Result<()> {
352        segment::clean_tmp(&self.dir)?;
353
354        let segments = self.scan_segments()?;
355
356        let mut inner = self.inner.lock();
357        let total_bytes: u64 = segments
358            .iter()
359            .filter_map(|s| fs::metadata(self.segment_path(s.start, s.end)).ok())
360            .map(|m| m.len())
361            .sum();
362
363        match (segments.first(), segments.last()) {
364            (Some(first), Some(last)) => {
365                inner.head_seq = first.start;
366                inner.next_seq = last.end + 1;
367            }
368            _ => {
369                inner.next_seq = 0;
370                inner.head_seq = 0;
371            }
372        }
373        inner.approx_disk_bytes = total_bytes;
374
375        info!(
376            segments = segments.len(),
377            head_seq = inner.head_seq,
378            next_seq = inner.next_seq,
379            disk_bytes = total_bytes,
380            "Segment buffer recovered"
381        );
382
383        Ok(())
384    }
385
386    fn write_segment(&self, start: u64, end: u64, events: &[T]) -> Result<u64> {
387        segment::write(
388            &self.dir,
389            self.config.cipher.as_deref(),
390            self.config.compression_level,
391            SegmentRange { start, end },
392            events,
393        )
394    }
395
396    fn read_segment(&self, seg: SegmentRange) -> Result<Vec<T>> {
397        segment::read(&self.dir, self.config.cipher.as_deref(), seg)
398    }
399
400    fn scan_segments(&self) -> Result<Vec<SegmentRange>> {
401        segment::scan(&self.dir)
402    }
403
404    fn segment_path(&self, start: u64, end: u64) -> PathBuf {
405        self.dir.join(segment::filename(start, end))
406    }
407}
408
409#[cfg(test)]
410mod tests;