use std::cmp::Ordering;
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as AtomicOrdering};
use std::sync::{Arc, Condvar, Mutex, RwLock};
use std::time::Duration;
use std::vec::IntoIter;
use thiserror::Error;
use seglog::read::Reader;
use crate::Position;
use crate::event::{DecodeError, Event, EventRef};
use crate::index::{
Access, ActiveTail, IndexSegment, SegmentIndex, choose, estimate_matches, search, search_back,
};
use crate::log::set::{
LogError, Record, RecordRef, Scan, ScanBack, Segment, SegmentSet, SegmentSource,
};
use crate::query::{Matches, Query};
use crate::index::IndexSet;
#[cfg(feature = "async")]
pub mod pool;
mod subscribe;
pub use subscribe::{DEFAULT_MAX_BATCH_EVENTS, Subscription};
#[derive(Clone, Copy, Debug)]
pub struct ReadConfig {
pub scan_bias: u32,
}
impl Default for ReadConfig {
fn default() -> Self {
ReadConfig { scan_bias: 4 }
}
}
pub struct Snapshot {
header_size: u64,
sealed_log: Vec<Arc<Segment>>,
sealed_index: Vec<Option<Arc<IndexSegment>>>,
active_log: Arc<Segment>,
active_index: Arc<ActiveTail>,
}
impl Snapshot {
pub(crate) fn capture(set: &SegmentSet, index: &IndexSet) -> Snapshot {
let sealed_log: Vec<Arc<Segment>> = set.sealed_arcs().to_vec();
let index_arcs = index.sealed_index_arcs();
let sealed_index: Vec<Option<Arc<IndexSegment>>> = (0..sealed_log.len())
.map(|i| index_arcs.get(i).cloned().flatten())
.collect();
Snapshot {
header_size: set.header_size(),
sealed_log,
sealed_index,
active_log: set.active_arc(),
active_index: index.active_tail_arc(),
}
}
}
const _: fn() = || {
fn is_send<T: Send>() {}
fn is_sync<T: Sync>() {}
is_send::<Snapshot>();
is_sync::<Snapshot>();
};
impl SegmentSource for Snapshot {
fn header_size(&self) -> u64 {
self.header_size
}
fn segment_count(&self) -> usize {
self.sealed_log.len()
}
fn segment_at(&self, idx: usize) -> Option<&Arc<Segment>> {
match idx.cmp(&self.sealed_log.len()) {
Ordering::Less => self.sealed_log.get(idx),
Ordering::Equal => Some(&self.active_log),
Ordering::Greater => None,
}
}
}
struct Notify {
lock: Mutex<()>,
cv: Condvar,
closed: AtomicBool,
subscribers: AtomicUsize,
#[cfg(feature = "async")]
async_event: event_listener::Event,
}
pub struct ReadCore {
segments: RwLock<Arc<Snapshot>>,
watermark: AtomicU64,
notify: Notify,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum WaitOutcome {
Advanced,
TimedOut,
Closed,
}
impl ReadCore {
pub(crate) fn new(set: &SegmentSet, index: &IndexSet) -> Arc<ReadCore> {
Arc::new(ReadCore {
segments: RwLock::new(Arc::new(Snapshot::capture(set, index))),
watermark: AtomicU64::new(set.last_position().get()),
notify: Notify {
lock: Mutex::new(()),
cv: Condvar::new(),
closed: AtomicBool::new(false),
subscribers: AtomicUsize::new(0),
#[cfg(feature = "async")]
async_event: event_listener::Event::new(),
},
})
}
pub(crate) fn publish_segments(&self, snapshot: Snapshot) {
*self.segments.write().unwrap() = Arc::new(snapshot);
}
pub(crate) fn publish_watermark(&self, tip: Position) {
self.watermark.store(tip.get(), AtomicOrdering::SeqCst);
}
fn head(&self) -> Position {
Position::new(self.watermark.load(AtomicOrdering::Acquire))
}
fn load(&self) -> (Position, Arc<Snapshot>) {
let watermark = Position::new(self.watermark.load(AtomicOrdering::Acquire));
let snapshot = Arc::clone(&self.segments.read().unwrap());
(watermark, snapshot)
}
fn watermark_gate(&self) -> Position {
Position::new(self.watermark.load(AtomicOrdering::SeqCst))
}
pub(crate) fn wake(&self) {
if self.notify.subscribers.load(AtomicOrdering::SeqCst) == 0 {
return;
}
let _guard = self.notify.lock.lock().unwrap();
self.notify.cv.notify_all();
#[cfg(feature = "async")]
self.notify.async_event.notify(usize::MAX);
}
pub(crate) fn close(&self) {
self.notify.closed.store(true, AtomicOrdering::Release);
let _guard = self.notify.lock.lock().unwrap();
self.notify.cv.notify_all();
#[cfg(feature = "async")]
self.notify.async_event.notify(usize::MAX);
}
fn register_subscriber(&self) {
self.notify.subscribers.fetch_add(1, AtomicOrdering::SeqCst);
}
fn deregister_subscriber(&self) {
self.notify
.subscribers
.fetch_sub(1, AtomicOrdering::Relaxed);
}
fn wait_past(&self, cursor: Position, timeout: Option<Duration>) -> WaitOutcome {
let mut guard = self.notify.lock.lock().unwrap();
loop {
if self.notify.closed.load(AtomicOrdering::Acquire) {
return WaitOutcome::Closed;
}
if self.watermark_gate() > cursor {
return WaitOutcome::Advanced;
}
match timeout {
None => guard = self.notify.cv.wait(guard).unwrap(),
Some(dur) => {
let (g, res) = self.notify.cv.wait_timeout(guard, dur).unwrap();
guard = g;
if res.timed_out() {
if self.notify.closed.load(AtomicOrdering::Acquire) {
return WaitOutcome::Closed;
}
if self.watermark_gate() > cursor {
return WaitOutcome::Advanced;
}
return WaitOutcome::TimedOut;
}
}
}
}
}
#[cfg(feature = "async")]
async fn wait_past_async(&self, cursor: Position) -> WaitOutcome {
loop {
if self.notify.closed.load(AtomicOrdering::Acquire) {
return WaitOutcome::Closed;
}
if self.watermark_gate() > cursor {
return WaitOutcome::Advanced;
}
let listener = self.notify.async_event.listen();
if self.notify.closed.load(AtomicOrdering::Acquire) {
return WaitOutcome::Closed;
}
if self.watermark_gate() > cursor {
return WaitOutcome::Advanced;
}
listener.await;
}
}
}
#[derive(Clone)]
pub struct ReadHandle {
core: Arc<ReadCore>,
config: ReadConfig,
}
impl ReadHandle {
pub(crate) fn new(core: Arc<ReadCore>, config: ReadConfig) -> ReadHandle {
ReadHandle { core, config }
}
pub fn read(&self, query: &Query, after: Position, limit: Option<u64>) -> Reads {
let (watermark, snapshot) = self.core.load();
Reads::plan(snapshot, query, after, watermark, &self.config, limit)
}
pub fn read_back(&self, query: &Query, before: Position, limit: Option<u64>) -> Reads {
let (watermark, snapshot) = self.core.load();
let upto = Position::new(before.get().saturating_sub(1).min(watermark.get()));
Reads::plan_back(snapshot, query, upto, watermark, &self.config, limit)
}
pub fn head(&self) -> Position {
self.core.head()
}
pub fn subscribe(&self, query: Query, after: Position) -> Subscription {
Subscription::new(Arc::clone(&self.core), self.config, query, after)
}
}
#[derive(Clone, Copy, Debug)]
pub struct Sequenced<'a> {
pub position: Position,
pub event: EventRef<'a>,
}
#[derive(Debug, Error)]
pub enum ReadError {
#[error("log error during read: {0}")]
Log(Arc<LogError>),
#[error("corrupt event during read: {0}")]
Corrupt(DecodeError),
#[error("read aborted before completion: the worker running it panicked")]
Aborted,
}
pub struct Reads {
watermark: Position,
pending_err: Option<ReadError>,
remaining: Option<u64>,
hit_limit: bool,
mode: Mode,
}
struct CachedReader {
seg_idx: usize,
segment: Arc<Segment>,
reader: Reader<0>,
}
struct IndexedState {
snapshot: Arc<Snapshot>,
positions: IntoIter<Position>,
reader: Option<CachedReader>,
buf: Option<Record>,
}
struct ScanFilteredState {
scan: Scan<Arc<Snapshot>>,
query: Query,
buf: Option<(Position, Event)>,
}
struct ScanBackFilteredState {
scan: ScanBack<Arc<Snapshot>>,
query: Query,
buf: Option<(Position, Event)>,
}
enum Mode {
Scan { scan: Box<Scan<Arc<Snapshot>>> },
ScanFiltered(Box<ScanFilteredState>),
ScanBack { scan: Box<ScanBack<Arc<Snapshot>>> },
ScanBackFiltered(Box<ScanBackFilteredState>),
Indexed(Box<IndexedState>),
Done,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum Direction {
Forward,
Backward,
}
impl Reads {
pub fn watermark(&self) -> Position {
self.watermark
}
pub fn is_exhausted(&self) -> bool {
!self.hit_limit
}
fn plan(
snapshot: Arc<Snapshot>,
query: &Query,
after: Position,
watermark: Position,
config: &ReadConfig,
limit: Option<u64>,
) -> Reads {
Reads::plan_directional(
Direction::Forward,
snapshot,
query,
after,
watermark,
watermark,
config,
limit,
)
}
fn plan_back(
snapshot: Arc<Snapshot>,
query: &Query,
upto: Position,
watermark: Position,
config: &ReadConfig,
limit: Option<u64>,
) -> Reads {
Reads::plan_directional(
Direction::Backward,
snapshot,
query,
Position::ZERO,
upto,
watermark,
config,
limit,
)
}
#[allow(clippy::too_many_arguments)]
fn plan_directional(
direction: Direction,
snapshot: Arc<Snapshot>,
query: &Query,
after: Position,
upto: Position,
watermark: Position,
config: &ReadConfig,
limit: Option<u64>,
) -> Reads {
if after >= upto || limit == Some(0) {
return Reads {
watermark,
pending_err: None,
remaining: limit,
hit_limit: false,
mode: Mode::Done,
};
}
let (estimate, width) = estimate_read(&snapshot, query, after, upto);
let access = choose(estimate, width, config.scan_bias);
#[cfg(feature = "tracing")]
tracing::debug!(
?access,
?direction,
estimate,
width,
scan_bias = config.scan_bias,
"read planner verdict"
);
match access {
Access::Scan => {
let is_all = matches!(*query, Query::All);
let mode = match direction {
Direction::Forward => {
let scan = Scan::start(Arc::clone(&snapshot), after.next(), upto);
if is_all {
Mode::Scan {
scan: Box::new(scan),
}
} else {
Mode::ScanFiltered(Box::new(ScanFilteredState {
scan,
query: query.clone(),
buf: None,
}))
}
}
Direction::Backward => {
let scan = ScanBack::start(Arc::clone(&snapshot), after.next(), upto);
if is_all {
Mode::ScanBack {
scan: Box::new(scan),
}
} else {
Mode::ScanBackFiltered(Box::new(ScanBackFilteredState {
scan,
query: query.clone(),
buf: None,
}))
}
}
};
Reads {
watermark,
pending_err: None,
remaining: limit,
hit_limit: false,
mode,
}
}
Access::Index => {
let planned = match direction {
Direction::Forward => plan_positions(&snapshot, query, after, upto, limit),
Direction::Backward => plan_positions_back(&snapshot, query, upto, limit),
};
match planned {
Ok(positions) => Reads {
watermark,
pending_err: None,
remaining: limit,
hit_limit: false,
mode: Mode::Indexed(Box::new(IndexedState {
snapshot,
positions: positions.into_iter(),
reader: None,
buf: None,
})),
},
Err(err) => Reads {
watermark,
pending_err: Some(err),
remaining: limit,
hit_limit: false,
mode: Mode::Done,
},
}
}
}
}
#[allow(clippy::should_implement_trait)]
pub fn next(&mut self) -> Option<Result<Sequenced<'_>, ReadError>> {
if let Some(err) = self.pending_err.take() {
self.mode = Mode::Done;
return Some(Err(err));
}
match &mut self.remaining {
Some(0) => {
self.hit_limit = true;
self.mode = Mode::Done;
return None;
}
Some(n) => *n -= 1,
None => {}
}
match &mut self.mode {
Mode::Done => None,
Mode::Scan { scan } => scan.next().map(decode_record),
Mode::ScanBack { scan } => scan.next().map(decode_record),
Mode::ScanFiltered(state) => {
let ScanFilteredState { scan, query, buf } = state.as_mut();
next_filtered(scan, query, buf)
}
Mode::ScanBackFiltered(state) => {
let ScanBackFilteredState { scan, query, buf } = state.as_mut();
next_filtered(scan, query, buf)
}
Mode::Indexed(state) => {
let IndexedState {
snapshot,
positions,
reader,
buf,
} = state.as_mut();
let position = positions.next()?;
let Some((seg_idx, _)) = snapshot.locate(position) else {
return Some(Err(ReadError::Log(Arc::new(LogError::NotFound {
position,
}))));
};
if reader.as_ref().map(|c| c.seg_idx) != Some(seg_idx) {
let segment = Arc::clone(snapshot.segment_at(seg_idx).unwrap());
match segment.open_reader() {
Ok(r) => {
*reader = Some(CachedReader {
seg_idx,
segment,
reader: r,
})
}
Err(err) => return Some(Err(ReadError::Log(Arc::new(err)))),
}
}
let cached = reader.as_mut().unwrap();
let local = (position.get() - cached.segment.base_position().get()) as usize;
match cached.segment.read_at_local(&mut cached.reader, local) {
Ok(Some(record)) => {
*buf = Some(record);
let bytes = &buf.as_ref().unwrap().data;
match EventRef::from_bytes(bytes) {
Ok(event) => Some(Ok(Sequenced { position, event })),
Err(err) => Some(Err(ReadError::Corrupt(err))),
}
}
Ok(None) => Some(Err(ReadError::Log(Arc::new(LogError::NotFound {
position,
})))),
Err(err) => Some(Err(ReadError::Log(Arc::new(err)))),
}
}
}
}
pub fn collect_owned(mut self) -> Result<Vec<(Position, crate::event::Event)>, ReadError> {
let mut out = Vec::new();
while let Some(item) = self.next() {
let seq = item?;
out.push((seq.position, seq.event.to_owned()));
}
Ok(out)
}
}
fn estimate_read(
snapshot: &Snapshot,
query: &Query,
after: Position,
watermark: Position,
) -> (u64, u64) {
let wm = watermark.get();
let mut estimate: u64 = 0;
let mut width: u64 = 0;
for (seg, index) in snapshot.sealed_log.iter().zip(snapshot.sealed_index.iter()) {
let base = seg.base_position();
let count = seg.event_count();
if count == 0 {
continue;
}
let effective_max = (base.get() + count - 1).min(wm);
if effective_max <= after.get() {
continue;
}
let seg_width = segment_width(after, base, effective_max);
width += seg_width;
estimate += match index {
Some(index_seg) => estimate_segment(index_seg.as_ref(), query, seg_width),
None => seg_width,
};
}
let active_base = snapshot.active_log.base_position();
if watermark >= active_base {
let seg_width = segment_width(after, active_base, wm);
width += seg_width;
estimate += if snapshot.active_index.is_unindexable() {
seg_width
} else {
estimate_segment(&snapshot.active_index.view(watermark), query, seg_width)
};
}
(estimate.min(width), width)
}
fn segment_width(after: Position, base: Position, effective_max: u64) -> u64 {
let first = first_after(after, base).get();
(effective_max + 1).saturating_sub(first)
}
fn estimate_segment<I: SegmentIndex>(index: &I, query: &Query, seg_width: u64) -> u64 {
let estimate = estimate_matches(index, query, seg_width);
#[cfg(feature = "tracing")]
if let Query::Items(items) = query {
for (item, spec) in items.iter().enumerate() {
tracing::trace!(
segment_base = index.base().get(),
item,
estimate = crate::index::estimate_item(index, spec, seg_width),
seg_width,
"planner per-item estimate"
);
}
}
estimate
}
fn plan_positions(
snapshot: &Arc<Snapshot>,
query: &Query,
after: Position,
watermark: Position,
limit: Option<u64>,
) -> Result<Vec<Position>, ReadError> {
let mut out = Vec::new();
let wm = watermark.get();
let cap = limit.map(|k| k as usize);
for (seg, index) in snapshot.sealed_log.iter().zip(snapshot.sealed_index.iter()) {
let base = seg.base_position();
let count = seg.event_count();
if count == 0 {
continue;
}
let effective_max = (base.get() + count - 1).min(wm);
if effective_max <= after.get() {
continue; }
match index {
Some(index_seg) => {
let iter = search(index_seg.as_ref(), query, after).take_while(|p| p.get() <= wm);
extend_capped(&mut out, iter, cap);
}
None => scan_positions_into(
snapshot,
query,
first_after(after, base),
Position::new(effective_max),
&mut out,
cap,
)?,
}
if let Some(k) = cap
&& out.len() >= k
{
out.truncate(k);
return Ok(out);
}
}
let active_base = snapshot.active_log.base_position();
if watermark >= active_base {
if snapshot.active_index.is_unindexable() {
scan_positions_into(
snapshot,
query,
first_after(after, active_base),
watermark,
&mut out,
cap,
)?;
} else {
let view = snapshot.active_index.view(watermark);
extend_capped(&mut out, search(&view, query, after), cap);
}
}
if let Some(k) = cap {
out.truncate(k);
}
Ok(out)
}
fn first_after(after: Position, base: Position) -> Position {
Position::new(after.get().max(base.get().saturating_sub(1)) + 1)
}
trait RecordScan {
fn next_record(&mut self) -> Option<Result<RecordRef<'_>, LogError>>;
}
impl RecordScan for Scan<Arc<Snapshot>> {
fn next_record(&mut self) -> Option<Result<RecordRef<'_>, LogError>> {
self.next()
}
}
impl RecordScan for ScanBack<Arc<Snapshot>> {
fn next_record(&mut self) -> Option<Result<RecordRef<'_>, LogError>> {
self.next()
}
}
fn extend_capped(
out: &mut Vec<Position>,
iter: impl Iterator<Item = Position>,
cap: Option<usize>,
) {
match cap {
Some(k) => out.extend(iter.take(k.saturating_sub(out.len()))),
None => out.extend(iter),
}
}
fn scan_positions_matching<S: RecordScan>(
mut scan: S,
query: &Query,
out: &mut Vec<Position>,
cap: Option<usize>,
) -> Result<(), ReadError> {
while let Some(item) = scan.next_record() {
let record = item.map_err(|err| ReadError::Log(Arc::new(err)))?;
let event = EventRef::from_bytes(record.data).map_err(ReadError::Corrupt)?;
if query.matches(event) {
out.push(record.position);
if let Some(k) = cap
&& out.len() >= k
{
break;
}
}
}
Ok(())
}
fn decode_record(item: Result<RecordRef<'_>, LogError>) -> Result<Sequenced<'_>, ReadError> {
match item {
Ok(record) => match EventRef::from_bytes(record.data) {
Ok(event) => Ok(Sequenced {
position: record.position,
event,
}),
Err(err) => Err(ReadError::Corrupt(err)),
},
Err(err) => Err(ReadError::Log(Arc::new(err))),
}
}
fn next_filtered<'a, S: RecordScan>(
scan: &mut S,
query: &Query,
buf: &'a mut Option<(Position, Event)>,
) -> Option<Result<Sequenced<'a>, ReadError>> {
loop {
match scan.next_record()? {
Ok(record) => {
let event = match EventRef::from_bytes(record.data) {
Ok(event) => event,
Err(err) => return Some(Err(ReadError::Corrupt(err))),
};
if query.matches(event) {
*buf = Some((record.position, event.to_owned()));
break;
}
}
Err(err) => return Some(Err(ReadError::Log(Arc::new(err)))),
}
}
let (position, event) = buf.as_ref().unwrap();
Some(Ok(Sequenced {
position: *position,
event: event.as_ref(),
}))
}
fn scan_positions_into(
snapshot: &Arc<Snapshot>,
query: &Query,
first: Position,
upto: Position,
out: &mut Vec<Position>,
cap: Option<usize>,
) -> Result<(), ReadError> {
if let Some(k) = cap
&& out.len() >= k
{
return Ok(());
}
scan_positions_matching(
Scan::start(Arc::clone(snapshot), first, upto),
query,
out,
cap,
)
}
fn plan_positions_back(
snapshot: &Arc<Snapshot>,
query: &Query,
upto: Position,
limit: Option<u64>,
) -> Result<Vec<Position>, ReadError> {
let mut out = Vec::new();
let wm = upto.get(); let cap = limit.map(|k| k as usize);
let active_base = snapshot.active_log.base_position();
if upto >= active_base {
if snapshot.active_index.is_unindexable() {
scan_positions_back_into(snapshot, query, active_base, upto, &mut out, cap)?;
} else {
let view = snapshot.active_index.view(upto);
let before = Position::new(upto.get().saturating_add(1));
extend_capped(&mut out, search_back(&view, query, before), cap);
}
if let Some(k) = cap
&& out.len() >= k
{
out.truncate(k);
return Ok(out);
}
}
for (seg, index) in snapshot
.sealed_log
.iter()
.zip(snapshot.sealed_index.iter())
.rev()
{
let base = seg.base_position();
let count = seg.event_count();
if count == 0 {
continue;
}
if base.get() > wm {
continue; }
let effective_max = (base.get() + count - 1).min(wm);
match index {
Some(index_seg) => {
let before = Position::new(effective_max.saturating_add(1));
extend_capped(
&mut out,
search_back(index_seg.as_ref(), query, before),
cap,
);
}
None => scan_positions_back_into(
snapshot,
query,
base,
Position::new(effective_max),
&mut out,
cap,
)?,
}
if let Some(k) = cap
&& out.len() >= k
{
out.truncate(k);
return Ok(out);
}
}
if let Some(k) = cap {
out.truncate(k);
}
Ok(out)
}
fn scan_positions_back_into(
snapshot: &Arc<Snapshot>,
query: &Query,
first: Position,
upto: Position,
out: &mut Vec<Position>,
cap: Option<usize>,
) -> Result<(), ReadError> {
if let Some(k) = cap
&& out.len() >= k
{
return Ok(());
}
scan_positions_matching(
ScanBack::start(Arc::clone(snapshot), first, upto),
query,
out,
cap,
)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{Event, EventType, Tag, Tags};
use crate::index::IndexSet;
use crate::log::set::SegmentConfig;
use crate::query::QueryItem;
use smallvec::SmallVec;
use tempfile::TempDir;
fn tags(items: &[&str]) -> Tags {
Tags::new(
items
.iter()
.map(|s| Tag::new(*s).unwrap())
.collect::<SmallVec<[Tag; 4]>>(),
)
.unwrap()
}
fn event(ty: &str, tag_strs: &[&str]) -> Event {
Event::new(&EventType::new(ty).unwrap(), &tags(tag_strs), b"x").unwrap()
}
#[test]
fn plan_positions_is_clamped_to_the_pinned_watermark() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
for _ in 0..60 {
set.append_batch(&[event("Enrolled", &["course:c1"]).as_bytes()])
.unwrap();
}
assert!(set.sealed_len() >= 2, "need several sealed segments");
let index = IndexSet::open(&set).unwrap();
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let query = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let last = set.last_position().get();
for wm in 1..=last {
let planned =
plan_positions(&snapshot, &query, Position::ZERO, Position::new(wm), None).unwrap();
let got: Vec<u64> = planned.iter().map(|p| p.get()).collect();
let expected: Vec<u64> = (1..=wm).collect();
assert_eq!(got, expected, "watermark {wm}");
}
}
#[test]
fn active_unindexable_reader_scans_the_log_for_complete_results() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(1 << 16)).unwrap();
for _ in 0..6 {
set.append_batch(&[event("Enrolled", &["course:c1"]).as_bytes()])
.unwrap();
}
let last = set.last_position();
let active_base = set.active_base();
let tail = ActiveTail::new(active_base);
let mut scan = set.scan_from(active_base);
for _ in 0..4 {
let record = scan.next().unwrap().unwrap();
let ev = EventRef::from_bytes(record.data).unwrap();
tail.push(record.position, ev).unwrap();
}
drop(scan);
tail.mark_unindexable();
assert_eq!(
tail.len(),
4,
"tail is truncated below the durable tip of 6"
);
let snapshot = Arc::new(Snapshot {
header_size: set.header_size(),
sealed_log: Vec::new(),
sealed_index: Vec::new(),
active_log: set.active_arc(),
active_index: Arc::new(tail),
});
let query = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let planned = plan_positions(&snapshot, &query, Position::ZERO, last, None).unwrap();
let got: Vec<u64> = planned.iter().map(|p| p.get()).collect();
assert_eq!(got, (1..=6).collect::<Vec<_>>());
}
#[test]
fn plan_positions_honors_the_limit() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
for _ in 0..60 {
set.append_batch(&[event("Enrolled", &["course:c1"]).as_bytes()])
.unwrap();
}
assert!(set.sealed_len() >= 2, "need several sealed segments");
let index = IndexSet::open(&set).unwrap();
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let query = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let watermark = set.last_position();
let full = plan_positions(&snapshot, &query, Position::ZERO, watermark, None).unwrap();
assert_eq!(full.len(), 60);
for limit in [0u64, 1, 5, 30, 59, 60, 61, 1000] {
let planned =
plan_positions(&snapshot, &query, Position::ZERO, watermark, Some(limit)).unwrap();
let want = &full[..(limit as usize).min(full.len())];
assert_eq!(planned, want, "limit {limit}");
}
}
#[test]
fn reads_next_respects_the_limit_and_reports_exhaustion() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
for _ in 0..40 {
set.append_batch(&[event("Enrolled", &["course:c1"]).as_bytes()])
.unwrap();
}
let index = IndexSet::open(&set).unwrap();
let collect = |limit: Option<u64>, scan_bias: u32| -> (Vec<u64>, bool) {
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let watermark = set.last_position();
let mut reads = Reads::plan(
snapshot,
&Query::item(QueryItem::with_tags(tags(&["course:c1"]))),
Position::ZERO,
watermark,
&ReadConfig { scan_bias },
limit,
);
let mut out = Vec::new();
while let Some(item) = reads.next() {
out.push(item.unwrap().position.get());
}
(out, reads.is_exhausted())
};
for scan_bias in [1u32, u32::MAX] {
let (full, full_exhausted) = collect(None, scan_bias);
assert_eq!(full, (1..=40).collect::<Vec<_>>(), "bias {scan_bias}");
assert!(full_exhausted, "an unlimited read is exhausted at its end");
let (capped, capped_exhausted) = collect(Some(10), scan_bias);
assert_eq!(capped, (1..=10).collect::<Vec<_>>(), "bias {scan_bias}");
assert!(!capped_exhausted, "a capped read is not exhausted");
let (zero, zero_exhausted) = collect(Some(0), scan_bias);
assert!(zero.is_empty(), "a zero cap yields nothing");
assert!(
!zero_exhausted,
"a zero cap is a limit stop, not exhaustion"
);
let (over, over_exhausted) = collect(Some(1000), scan_bias);
assert_eq!(
over,
(1..=40).collect::<Vec<_>>(),
"a cap above the result returns all"
);
assert!(over_exhausted, "an under-cap read drains and is exhausted");
}
}
#[test]
fn plan_positions_back_mirrors_plan_positions() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
for _ in 0..60 {
set.append_batch(&[event("Enrolled", &["course:c1"]).as_bytes()])
.unwrap();
}
assert!(set.sealed_len() >= 2, "need several sealed segments");
let index = IndexSet::open(&set).unwrap();
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let query = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let last = set.last_position();
let forward = plan_positions(&snapshot, &query, Position::ZERO, last, None).unwrap();
let want: Vec<Position> = forward.iter().rev().copied().collect();
assert_eq!(
plan_positions_back(&snapshot, &query, last, None).unwrap(),
want
);
for upto in [1u64, 5, 30, 59, 60] {
let back = plan_positions_back(&snapshot, &query, Position::new(upto), None).unwrap();
let want: Vec<Position> = forward
.iter()
.filter(|p| p.get() <= upto)
.rev()
.copied()
.collect();
assert_eq!(back, want, "upto {upto}");
}
let full: Vec<Position> = forward.iter().rev().copied().collect();
for limit in [0u64, 1, 5, 30, 60, 61] {
let back = plan_positions_back(&snapshot, &query, last, Some(limit)).unwrap();
assert_eq!(
back,
full[..(limit as usize).min(full.len())],
"limit {limit}"
);
}
}
#[test]
fn reads_back_is_the_reverse_of_reads_forward() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
for i in 0..40u64 {
let carries: &[&str] = if i % 2 == 0 {
&["course:c1"]
} else {
&["course:c1", "student:s1"]
};
set.append_batch(&[event("Enrolled", carries).as_bytes()])
.unwrap();
}
let index = IndexSet::open(&set).unwrap();
let collect_forward = |query: &Query, scan_bias: u32| -> Vec<u64> {
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let wm = set.last_position();
let mut reads = Reads::plan(
snapshot,
query,
Position::ZERO,
wm,
&ReadConfig { scan_bias },
None,
);
let mut out = Vec::new();
while let Some(item) = reads.next() {
out.push(item.unwrap().position.get());
}
out
};
let collect_back =
|query: &Query, scan_bias: u32, limit: Option<u64>| -> (Vec<u64>, bool) {
let snapshot = Arc::new(Snapshot::capture(&set, &index));
let wm = set.last_position();
let mut reads =
Reads::plan_back(snapshot, query, wm, wm, &ReadConfig { scan_bias }, limit);
let mut out = Vec::new();
while let Some(item) = reads.next() {
out.push(item.unwrap().position.get());
}
(out, reads.is_exhausted())
};
let queries = [
Query::all(),
Query::item(QueryItem::with_tags(tags(&["student:s1"]))),
];
for query in &queries {
for scan_bias in [1u32, u32::MAX] {
let want: Vec<u64> = collect_forward(query, scan_bias)
.into_iter()
.rev()
.collect();
let (back, exhausted) = collect_back(query, scan_bias, None);
assert_eq!(back, want, "query {query:?} bias {scan_bias}");
assert!(
exhausted,
"an unlimited reverse read is exhausted at its end"
);
let (capped, capped_exhausted) = collect_back(query, scan_bias, Some(5));
let want_capped: Vec<u64> = want.iter().take(5).copied().collect();
assert_eq!(
capped, want_capped,
"capped query {query:?} bias {scan_bias}"
);
assert!(!capped_exhausted, "a capped reverse read is not exhausted");
}
}
}
}