use segment_buffer::{DurabilityPolicy, FlushPolicy, SegmentBuffer, SegmentConfig};
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Mutex;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
struct Event {
id: u64,
payload: String,
}
trait CloudUploader: Send + Sync {
fn upload(&self, start_seq: u64, batch: &[Event]) -> Result<usize, String>;
}
struct ReliableUploader {
received: AtomicU64,
}
impl CloudUploader for ReliableUploader {
fn upload(&self, _start_seq: u64, batch: &[Event]) -> Result<usize, String> {
let bytes: usize = batch.iter().map(|e| e.payload.len()).sum();
self.received
.fetch_add(batch.len() as u64, Ordering::Relaxed);
Ok(bytes)
}
}
struct FlakyUploader {
received: AtomicU64,
failed_once_for: Mutex<std::collections::HashSet<u64>>,
}
impl CloudUploader for FlakyUploader {
fn upload(&self, start_seq: u64, batch: &[Event]) -> Result<usize, String> {
let mut seen = self.failed_once_for.lock().unwrap();
if !seen.contains(&start_seq) {
seen.insert(start_seq);
return Err(format!(
"transient cloud outage on batch starting at seq {start_seq}"
));
}
let bytes: usize = batch.iter().map(|e| e.payload.len()).sum();
self.received
.fetch_add(batch.len() as u64, Ordering::Relaxed);
Ok(bytes)
}
}
fn drain_loop(
buf: &SegmentBuffer<Event>,
uploader: &dyn CloudUploader,
cursor: &mut u64,
max_retries: u32,
) -> Result<(u64, u32), String> {
let batch_size = 500usize;
let mut uploaded = 0u64;
let mut retries = 0u32;
loop {
let batch = buf
.read_from(*cursor, batch_size)
.map_err(|e| format!("read_from failed: {e}"))?;
if batch.is_empty() {
break; }
let last_seq = *cursor + batch.len() as u64 - 1;
let mut succeeded = false;
for attempt in 0..=max_retries {
match uploader.upload(*cursor, &batch) {
Ok(_bytes) => {
succeeded = true;
uploaded += batch.len() as u64;
break;
}
Err(e) => {
retries += 1;
if attempt == max_retries {
return Err(format!(
"upload failed after {max_retries} retries at seq {}: {e}",
*cursor
));
}
eprintln!(" retry {attempt} at seq {}: {e}", *cursor);
std::thread::sleep(std::time::Duration::from_millis(10));
}
}
}
if !succeeded {
unreachable!("either succeeded or returned Err in the loop above");
}
let _ = buf
.delete_acked(last_seq)
.map_err(|e| format!("delete_acked failed: {e}"))?;
*cursor = last_seq + 1;
}
Ok((uploaded, retries))
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let tmp = tempfile::tempdir()?;
let config = SegmentConfig::builder()
.flush_policy(FlushPolicy::Batch(500))
.durability(DurabilityPolicy::Throughput)
.build();
let buf: SegmentBuffer<Event> = SegmentBuffer::open(tmp.path(), config)?;
const TOTAL: u64 = 10_000;
for id in 0..TOTAL {
buf.append(Event {
id,
payload: format!("event-{id}"),
})?;
}
buf.flush()?;
println!(
"produced {TOTAL} events; pending_count = {}",
buf.pending_count()
);
let mut cursor = buf.stats().head_sequence;
println!("\n--- drain with ReliableUploader ---");
let (uploaded, retries) = drain_loop(
&buf,
&ReliableUploader {
received: AtomicU64::new(0),
},
&mut cursor,
3,
)?;
println!(
"uploaded {uploaded} events, {retries} retries; pending_count = {}",
buf.pending_count()
);
assert_eq!(uploaded, TOTAL);
assert_eq!(retries, 0);
for id in 0..TOTAL {
buf.append(Event {
id: id + TOTAL, payload: format!("event2-{id}"),
})?;
}
buf.flush()?;
let mut cursor = buf.stats().head_sequence;
println!("\n--- drain with FlakyUploader (transient failures) ---");
let flaky = FlakyUploader {
received: AtomicU64::new(0),
failed_once_for: Mutex::new(std::collections::HashSet::new()),
};
let (uploaded, retries) = drain_loop(&buf, &flaky, &mut cursor, 3)?;
let received = flaky.received.load(Ordering::Relaxed);
println!(
"uploaded {uploaded} events (cloud-side: {received} received), {retries} retries; \
pending_count = {}",
buf.pending_count()
);
assert_eq!(uploaded, TOTAL);
assert_eq!(received, TOTAL);
assert!(
retries >= 20,
"flaky uploader should have forced at least one retry per batch"
);
println!("\nAll drained. The buffer directory now contains only the lock sidecar.");
Ok(())
}