#![warn(missing_docs)]
mod cipher;
mod error;
mod segment;
#[cfg(feature = "encryption")]
pub use cipher::AesGcmCipher;
pub use cipher::{CipherError, SegmentCipher};
pub use error::{Result, SegmentError};
use std::fs;
use std::path::PathBuf;
use std::time::Instant;
use parking_lot::Mutex;
use serde::de::DeserializeOwned;
use serde::Serialize;
use tracing::{debug, info};
use segment::SegmentRange;
#[non_exhaustive]
pub struct SegmentConfig {
pub max_batch_events: usize,
pub flush_interval_secs: u64,
pub max_size_bytes: u64,
pub compression_level: i32,
pub cipher: Option<Box<dyn SegmentCipher>>,
}
impl std::fmt::Debug for SegmentConfig {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SegmentConfig")
.field("max_batch_events", &self.max_batch_events)
.field("flush_interval_secs", &self.flush_interval_secs)
.field("max_size_bytes", &self.max_size_bytes)
.field("compression_level", &self.compression_level)
.field("cipher", &self.cipher.as_ref().map(|_| "[set]"))
.finish()
}
}
impl Default for SegmentConfig {
fn default() -> Self {
Self {
max_batch_events: 256,
flush_interval_secs: 5,
max_size_bytes: 10 * 1024 * 1024 * 1024,
compression_level: 3,
cipher: None,
}
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct BufferStats {
pub pending_count: u64,
pub latest_sequence: u64,
pub head_sequence: u64,
pub next_sequence: u64,
pub approx_disk_bytes: u64,
pub max_size_bytes: u64,
pub store_pressure: f32,
}
struct BufferInner<T> {
unflushed: Vec<T>,
next_seq: u64,
head_seq: u64,
last_flush: Instant,
approx_disk_bytes: u64,
}
pub struct SegmentBuffer<T> {
dir: PathBuf,
config: SegmentConfig,
inner: Mutex<BufferInner<T>>,
}
impl<T> std::fmt::Debug for SegmentBuffer<T>
where
T: Serialize + DeserializeOwned + Clone + Send + 'static,
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let stats = self.stats();
f.debug_struct("SegmentBuffer")
.field("dir", &self.dir)
.field("pending_count", &stats.pending_count)
.field("latest_sequence", &stats.latest_sequence)
.field("head_sequence", &stats.head_sequence)
.field("next_sequence", &stats.next_sequence)
.field("approx_disk_bytes", &stats.approx_disk_bytes)
.field("max_size_bytes", &stats.max_size_bytes)
.field("store_pressure", &stats.store_pressure)
.finish()
}
}
impl<T> SegmentBuffer<T>
where
T: Serialize + DeserializeOwned + Clone + Send + 'static,
{
pub fn open(dir: impl Into<PathBuf>, config: SegmentConfig) -> Result<Self> {
let dir = dir.into();
fs::create_dir_all(&dir)?;
let buffer = Self {
dir,
config,
inner: Mutex::new(BufferInner {
unflushed: Vec::new(),
next_seq: 0,
head_seq: 0,
last_flush: Instant::now(),
approx_disk_bytes: 0,
}),
};
buffer.recover()?;
Ok(buffer)
}
pub fn append(&self, event: T) -> Result<u64> {
let (should_flush, seq) = {
let mut inner = self.inner.lock();
inner.unflushed.push(event);
inner.next_seq += 1;
let seq = inner.next_seq - 1;
let batch_full = inner.unflushed.len() >= self.config.max_batch_events;
let interval_elapsed =
inner.last_flush.elapsed().as_secs() >= self.config.flush_interval_secs;
(batch_full || interval_elapsed, seq)
};
if should_flush {
self.flush()?;
}
Ok(seq)
}
pub fn flush(&self) -> Result<()> {
let (events, start_seq, end_seq) = {
let mut inner = self.inner.lock();
inner.last_flush = Instant::now();
if inner.unflushed.is_empty() {
return Ok(());
}
let events = std::mem::take(&mut inner.unflushed);
let count = events.len() as u64;
let end_seq = inner.next_seq - 1;
let start_seq = end_seq + 1 - count;
(events, start_seq, end_seq)
};
let compressed_len = self.write_segment(start_seq, end_seq, &events)?;
{
let mut inner = self.inner.lock();
inner.approx_disk_bytes += compressed_len;
}
debug!(start_seq, end_seq, count = events.len(), "Flushed segment");
Ok(())
}
pub fn read_from(&self, start_seq: u64, limit: usize) -> Result<Vec<T>> {
if limit == 0 {
return Ok(Vec::new());
}
let mut result: Vec<T> = Vec::with_capacity(limit.min(1024));
let segments = self.scan_segments()?;
for seg in &segments {
if result.len() >= limit {
break;
}
if seg.end < start_seq {
continue;
}
let events = self.read_segment(*seg)?;
let skip = if seg.start < start_seq {
(start_seq - seg.start) as usize
} else {
0
};
for event in events.into_iter().skip(skip) {
if result.len() >= limit {
break;
}
result.push(event);
}
}
if result.len() < limit {
let inner = self.inner.lock();
let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
for (i, event) in inner.unflushed.iter().enumerate() {
let seq = pending_start + i as u64;
if seq < start_seq {
continue;
}
if result.len() >= limit {
break;
}
result.push(event.clone());
}
}
Ok(result)
}
pub fn delete_acked(&self, acked_seq: u64) -> Result<usize> {
let segments = self.scan_segments()?;
let mut deleted = 0;
let mut freed_bytes: u64 = 0;
let mut new_head = None;
for seg in &segments {
if seg.end <= acked_seq {
let path = self.segment_path(seg.start, seg.end);
if let Ok(meta) = fs::metadata(&path) {
freed_bytes += meta.len();
}
match fs::remove_file(&path) {
Ok(()) => {
deleted += 1;
debug!(start = seg.start, end = seg.end, "Deleted acked segment");
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {}
Err(e) => return Err(e.into()),
}
} else if new_head.is_none() {
new_head = Some(seg.start);
}
}
{
let mut inner = self.inner.lock();
inner.approx_disk_bytes = inner.approx_disk_bytes.saturating_sub(freed_bytes);
let pending_start = inner.next_seq.saturating_sub(inner.unflushed.len() as u64);
inner.head_seq = new_head.unwrap_or(inner.next_seq).min(pending_start);
}
if deleted > 0 {
info!(deleted, freed_bytes, acked_seq, "Deleted acked segments");
}
Ok(deleted)
}
#[must_use = "the sequence number is meaningless if discarded"]
pub fn latest_sequence(&self) -> u64 {
let inner = self.inner.lock();
if inner.next_seq == 0 {
0
} else {
inner.next_seq - 1
}
}
#[must_use = "the backlog size is meaningless if discarded"]
pub fn pending_count(&self) -> u64 {
let inner = self.inner.lock();
inner.next_seq.saturating_sub(inner.head_seq)
}
#[must_use = "the backlog size is meaningless if discarded"]
pub fn len(&self) -> u64 {
self.pending_count()
}
#[must_use = "the emptiness flag is meaningless if discarded"]
pub fn is_empty(&self) -> bool {
self.pending_count() == 0
}
#[must_use = "the pressure value is meaningless if discarded"]
pub fn store_pressure(&self) -> f32 {
let inner = self.inner.lock();
if self.config.max_size_bytes == 0 {
return 0.0;
}
(inner.approx_disk_bytes as f32 / self.config.max_size_bytes as f32).min(1.0)
}
#[must_use = "the overload flag is meaningless if discarded"]
pub fn is_overloaded(&self) -> bool {
self.store_pressure() > 0.9
}
#[must_use = "the snapshot is meaningless if discarded"]
pub fn stats(&self) -> BufferStats {
let inner = self.inner.lock();
let pending_count = inner.next_seq.saturating_sub(inner.head_seq);
let latest_sequence = if inner.next_seq == 0 {
0
} else {
inner.next_seq - 1
};
let store_pressure = if self.config.max_size_bytes == 0 {
0.0
} else {
(inner.approx_disk_bytes as f32 / self.config.max_size_bytes as f32).min(1.0)
};
BufferStats {
pending_count,
latest_sequence,
head_sequence: inner.head_seq,
next_sequence: inner.next_seq,
approx_disk_bytes: inner.approx_disk_bytes,
max_size_bytes: self.config.max_size_bytes,
store_pressure,
}
}
fn recover(&self) -> Result<()> {
segment::clean_tmp(&self.dir)?;
let segments = self.scan_segments()?;
let total_bytes: u64 = segments
.iter()
.filter_map(|s| fs::metadata(self.segment_path(s.start, s.end)).ok())
.map(|m| m.len())
.sum();
let (head_seq, next_seq) = match (segments.first(), segments.last()) {
(Some(first), Some(last)) => (first.start, last.end + 1),
_ => (0, 0),
};
let segment_count = segments.len();
{
let mut inner = self.inner.lock();
inner.head_seq = head_seq;
inner.next_seq = next_seq;
inner.approx_disk_bytes = total_bytes;
}
info!(
segments = segment_count,
head_seq,
next_seq,
disk_bytes = total_bytes,
"Segment buffer recovered"
);
Ok(())
}
fn write_segment(&self, start: u64, end: u64, events: &[T]) -> Result<u64> {
segment::write(
&self.dir,
self.config.cipher.as_deref(),
self.config.compression_level,
SegmentRange::new(start, end),
events,
)
}
fn read_segment(&self, seg: SegmentRange) -> Result<Vec<T>> {
segment::read(&self.dir, self.config.cipher.as_deref(), seg)
}
fn scan_segments(&self) -> Result<Vec<SegmentRange>> {
segment::scan(&self.dir)
}
fn segment_path(&self, start: u64, end: u64) -> PathBuf {
self.dir.join(segment::filename(start, end))
}
}
const _: () = {
const fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<SegmentBuffer<()>>();
};
#[cfg(test)]
mod tests;
#[cfg(test)]
mod property_tests;