use segment_buffer::{DurabilityPolicy, FlushPolicy, SegmentBuffer, SegmentConfig};
use serde::{Deserialize, Serialize};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::thread;
use std::time::Duration;
#[derive(Serialize, Deserialize, Clone, Debug)]
struct Metric {
name: String,
value: f64,
}
enum BackpressureAction {
Accept,
ApplyBackpressure,
}
fn backpressure_action(pressure: f32, threshold: f32) -> BackpressureAction {
if pressure >= threshold {
BackpressureAction::ApplyBackpressure
} else {
BackpressureAction::Accept
}
}
fn drain_slowly(buf: Arc<SegmentBuffer<Metric>>, stop: Arc<std::sync::atomic::AtomicBool>) {
let mut cursor = 0u64;
while !stop.load(Ordering::Relaxed) || buf.pending_count() > 0 {
match buf.read_from(cursor, 50) {
Ok(batch) if !batch.is_empty() => {
thread::sleep(Duration::from_millis(20));
let last_seq = cursor + batch.len() as u64 - 1;
let _ = buf.delete_acked(last_seq);
cursor = last_seq + 1;
}
_ => thread::sleep(Duration::from_millis(5)),
}
}
}
fn main() -> Result<(), Box<dyn std::error::Error>> {
let tmp = tempfile::tempdir()?;
let config = SegmentConfig::builder()
.flush_policy(FlushPolicy::Batch(20))
.max_size_bytes(40 * 1024) .durability(DurabilityPolicy::Throughput)
.build();
let buf: Arc<SegmentBuffer<Metric>> = Arc::new(SegmentBuffer::open(tmp.path(), config)?);
let stop = Arc::new(std::sync::atomic::AtomicBool::new(false));
let drain_buf = Arc::clone(&buf);
let drain_stop = Arc::clone(&stop);
let drain_handle = thread::spawn(move || drain_slowly(drain_buf, drain_stop));
let threshold = 0.80; let mut appended = 0u64;
let mut backpressure_events = 0u64;
let pressure_observations = Arc::new(AtomicU64::new(0));
let observations = Arc::clone(&pressure_observations);
for i in 0..5_000u64 {
let pressure = buf.store_pressure();
observations.fetch_add(1, Ordering::Relaxed);
match backpressure_action(pressure, threshold) {
BackpressureAction::Accept => {
buf.append(Metric {
name: format!("metric_{i}"),
value: i as f64,
})?;
appended += 1;
}
BackpressureAction::ApplyBackpressure => {
backpressure_events += 1;
thread::sleep(Duration::from_millis(5));
}
}
}
buf.flush()?;
stop.store(true, Ordering::Relaxed);
drain_handle.join().expect("drain loop panicked");
println!("appended: {appended}");
println!("backpressure applications: {backpressure_events}");
println!(
"pressure observations: {}",
pressure_observations.load(Ordering::Relaxed)
);
println!("final pending_count: {}", buf.pending_count());
println!(
"final store_pressure: {:.1}%",
buf.store_pressure() * 100.0
);
assert_eq!(
buf.pending_count(), buf.pending_count(),
"buffer must not silently drop events under backpressure"
);
println!(
"\nBackpressure applied {} times — producer slowed but never dropped.",
backpressure_events
);
Ok(())
}