use segment_buffer::{FlushPolicy, SegmentBuffer, SegmentConfig};
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::Duration;
#[derive(Serialize, Deserialize, Clone, Debug)]
struct Event {
id: u64,
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let tmp = tempfile::tempdir()?;
let config = SegmentConfig::builder()
.flush_policy(FlushPolicy::Manual)
.max_size_bytes(1024 * 1024)
.build();
let buffer = Arc::new(SegmentBuffer::<Event>::open(tmp.path(), config)?);
let shutdown = Arc::new(AtomicBool::new(false));
let flusher_buf = Arc::clone(&buffer);
let flusher_shutdown = Arc::clone(&shutdown);
let flusher = thread::Builder::new()
.name("segment-buffer-flusher".into())
.spawn(move || {
while !flusher_shutdown.load(Ordering::Relaxed) {
thread::sleep(Duration::from_millis(50));
if let Err(e) = flusher_buf.flush() {
eprintln!("background flush failed: {e}");
break;
}
}
})?;
for i in 0..10_000 {
buffer.append(Event { id: i })?;
}
shutdown.store(true, Ordering::Relaxed);
flusher.join().expect("flusher panicked");
buffer.flush()?;
let items = buffer.read_from(0, 20_000)?;
println!("Recovered {} items", items.len());
assert_eq!(items.len(), 10_000);
Ok(())
}