use std::collections::VecDeque;
use std::sync::Arc;
use std::time::Duration;
use broadcast_common::Timestamp;
use crate::trunk::{ArchiveOverrun, SegmentCursor, SegmentCursorItem, SegmentEntry, Trunk};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum SinkOutcome {
Taken,
Busy,
}
pub trait SegmentSink: Send {
fn offer(&mut self, entry: &SegmentEntry) -> SinkOutcome;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum Retention {
HotOnly,
Tiered {
on_overrun: ArchiveOverrun,
cold_window: Duration,
},
}
impl Default for Retention {
fn default() -> Self {
Retention::HotOnly
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum SegmentLocation {
Hot,
Cold,
Evicted,
}
struct ColdEntry {
sequence_number: u32,
handed_off_at: Timestamp,
}
pub struct RetentionDriver<S> {
trunk: Arc<Trunk>,
cursor: SegmentCursor,
sink: S,
cold_window: Duration,
in_flight: Option<SegmentEntry>,
cold: VecDeque<ColdEntry>,
terminated: bool,
}
impl<S: SegmentSink> RetentionDriver<S> {
pub fn new(trunk: &Arc<Trunk>, retention: Retention, sink: S) -> Option<Self> {
match retention {
Retention::HotOnly => None,
Retention::Tiered {
on_overrun,
cold_window,
} => Some(RetentionDriver {
trunk: Arc::clone(trunk),
cursor: trunk.pin_segments(on_overrun),
sink,
cold_window,
in_flight: None,
cold: VecDeque::new(),
terminated: false,
}),
}
}
pub fn drive(&mut self, now: Timestamp) {
self.expire_cold(now);
loop {
if let Some(entry) = self.in_flight.take() {
match self.sink.offer(&entry) {
SinkOutcome::Taken => {
self.cold.push_back(ColdEntry {
sequence_number: entry.sequence_number,
handed_off_at: now,
});
}
SinkOutcome::Busy => {
self.in_flight = Some(entry);
return;
}
}
}
match self.cursor.poll() {
Some(SegmentCursorItem::Segment(entry)) => {
self.in_flight = Some(entry);
}
Some(SegmentCursorItem::Gap { .. }) | Some(SegmentCursorItem::Lagged { .. }) => {
}
Some(SegmentCursorItem::Terminated) => {
self.terminated = true;
return;
}
None => return,
}
}
}
pub fn pending_len(&self) -> usize {
usize::from(self.in_flight.is_some())
}
pub fn cold_len(&self) -> usize {
self.cold.len()
}
pub fn locate(&mut self, sequence_number: u32, now: Timestamp) -> SegmentLocation {
self.expire_cold(now);
if self
.cold
.iter()
.any(|c| c.sequence_number == sequence_number)
{
return SegmentLocation::Cold;
}
if self.terminated {
return SegmentLocation::Evicted;
}
match self.trunk.last_closed_segment() {
Some(last) if sequence_number <= last => SegmentLocation::Evicted,
_ => SegmentLocation::Hot,
}
}
fn expire_cold(&mut self, now: Timestamp) {
while let Some(front) = self.cold.front() {
if front.handed_off_at.saturating_add(self.cold_window) <= now {
self.cold.pop_front();
} else {
break;
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::trunk::TrunkConfig;
use std::num::NonZeroUsize;
use std::sync::mpsc;
use std::thread;
use transmux::SegmentMeta;
fn nz(n: usize) -> NonZeroUsize {
NonZeroUsize::new(n).expect("test capacity must be non-zero")
}
fn segment_entry(seq: u32) -> SegmentEntry {
SegmentEntry::new(
bytes::Bytes::from(vec![seq as u8; 4]),
seq,
Duration::from_secs(2),
Timestamp::from_nanos(u64::from(seq) * 2_000_000_000),
SegmentMeta {
discontinuous: false,
},
)
}
struct ScriptedSink {
always_busy: bool,
offered: Vec<u32>,
taken: Vec<u32>,
}
impl ScriptedSink {
fn new(always_busy: bool) -> Self {
ScriptedSink {
always_busy,
offered: Vec::new(),
taken: Vec::new(),
}
}
}
impl SegmentSink for ScriptedSink {
fn offer(&mut self, entry: &SegmentEntry) -> SinkOutcome {
self.offered.push(entry.sequence_number);
if self.always_busy {
SinkOutcome::Busy
} else {
self.taken.push(entry.sequence_number);
SinkOutcome::Taken
}
}
}
#[test]
fn segment_is_cold_for_the_window_then_evicted() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(4), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::Gap,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(false))
.expect("Tiered must build a driver");
let handed_off_at = Timestamp::from_nanos(1_000_000_000);
writer.publish_segment(segment_entry(1));
driver.drive(handed_off_at);
assert_eq!(driver.sink.offered, vec![1]);
assert_eq!(driver.sink.taken, vec![1]);
assert_eq!(
driver.locate(1, handed_off_at),
SegmentLocation::Cold,
"just handed off: must be Cold, not Hot or Evicted"
);
let before_expiry = handed_off_at.saturating_add(Duration::from_secs(9));
assert_eq!(driver.locate(1, before_expiry), SegmentLocation::Cold);
let at_expiry = handed_off_at.saturating_add(Duration::from_secs(10));
assert_eq!(driver.locate(1, at_expiry), SegmentLocation::Evicted);
}
#[test]
fn publish_segment_never_touches_the_sink() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(4), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::Gap,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(true))
.expect("Tiered must build a driver");
for seq in 1..=8u32 {
writer.publish_segment(segment_entry(seq));
}
assert_eq!(trunk.segment_len(), 4, "hot ring still obeys its own cap");
driver.drive(Timestamp::from_nanos(0));
assert_eq!(driver.sink.offered, vec![5]);
assert!(driver.sink.taken.is_empty());
}
#[test]
fn slow_sink_pending_hand_off_queue_is_bounded_to_one() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(64), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::Gap,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(true))
.expect("Tiered must build a driver");
for seq in 1..=10u32 {
writer.publish_segment(segment_entry(seq));
}
driver.drive(Timestamp::from_nanos(0));
assert_eq!(
driver.pending_len(),
1,
"exactly one entry may be in flight, never more, never fewer once a segment exists"
);
assert_eq!(
driver.sink.offered,
vec![1],
"only the single in-flight entry is ever offered while it is stuck; \
a flood behind it must not grow the hand-off queue"
);
driver.drive(Timestamp::from_nanos(1));
driver.drive(Timestamp::from_nanos(2));
assert_eq!(driver.pending_len(), 1);
assert!(
driver.sink.offered.iter().all(|&seq| seq == 1),
"only ever retries the single stuck entry, never advances to \
the next queued segment while it is stuck: {:?}",
driver.sink.offered
);
}
#[test]
fn archive_overrun_gap_reports_loss_and_locate_reflects_it() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(2), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::Gap,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(false))
.expect("Tiered must build a driver");
writer.publish_segment(segment_entry(1));
writer.publish_segment(segment_entry(2));
writer.publish_segment(segment_entry(3));
driver.drive(Timestamp::from_nanos(0));
assert_eq!(driver.sink.taken, vec![2, 3]);
assert_eq!(
driver.locate(1, Timestamp::from_nanos(0)),
SegmentLocation::Evicted,
"gapped before hand-off: genuinely gone, not Cold"
);
assert_eq!(
driver.locate(2, Timestamp::from_nanos(0)),
SegmentLocation::Cold
);
}
#[test]
fn archive_overrun_stall_ingest_blocks_writer_until_driver_advances() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(1), nz(8), nz(8)));
let writer = Arc::new(trunk.segment_writer().unwrap());
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::StallIngest,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(false))
.expect("Tiered must build a driver");
writer.publish_segment(segment_entry(1));
let (done_tx, done_rx) = mpsc::channel();
let blocked_writer = Arc::clone(&writer);
let handle = thread::spawn(move || {
blocked_writer.publish_segment(segment_entry(2));
done_tx.send(()).unwrap();
});
assert!(
done_rx.recv_timeout(Duration::from_millis(200)).is_err(),
"publish_segment must block: the driver has not consumed seq 1 yet"
);
driver.drive(Timestamp::from_nanos(0));
done_rx
.recv_timeout(Duration::from_secs(60))
.expect("publish_segment must unblock once the driver drains its pin");
handle.join().unwrap();
assert_eq!(
driver.sink.taken.first().copied(),
Some(1),
"seq 1 must reach the sink first — it is what the driver drained \
to release the pin"
);
driver.drive(Timestamp::from_nanos(1));
assert_eq!(
driver.sink.taken,
vec![1, 2],
"both segments must reach the sink in order once the stall clears"
);
}
#[test]
fn archive_overrun_terminate_drops_pin_and_locate_stays_honest() {
let trunk = Trunk::new(TrunkConfig::new(nz(10), nz(10), nz(1), nz(8), nz(8)));
let writer = trunk.segment_writer().unwrap();
let retention = Retention::Tiered {
on_overrun: ArchiveOverrun::Terminate,
cold_window: Duration::from_secs(10),
};
let mut driver = RetentionDriver::new(&trunk, retention, ScriptedSink::new(false))
.expect("Tiered must build a driver");
writer.publish_segment(segment_entry(1));
writer.publish_segment(segment_entry(2)); driver.drive(Timestamp::from_nanos(0));
assert!(driver.terminated, "Terminate must have fired");
assert_eq!(
driver.locate(3, Timestamp::from_nanos(0)),
SegmentLocation::Evicted,
"once terminated, this driver cannot honestly claim Hot for \
anything it is no longer protecting, produced or not"
);
}
}