use std::fmt;
use std::fs::{self, File};
use std::io;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::vec::IntoIter;
use thiserror::Error;
use crate::Position;
use crate::event::{DecodeError, EventRef};
use crate::log::set::{LogError, PositionRange, SegmentSet};
use crate::query::Query;
use super::ActiveTail;
use super::recovery::Rebuilder;
use super::search;
use super::segment::{IndexSegment, write_segment_file};
const NAME_DIGITS: usize = 20;
struct SealedIndex {
base: Position,
count: u64,
seg: Option<Arc<IndexSegment>>,
}
impl SealedIndex {
fn max_position(&self) -> Option<u64> {
(self.count > 0).then(|| self.base.get() + self.count - 1)
}
}
pub struct IndexSet {
dir: PathBuf,
sealed: Vec<SealedIndex>,
active: Arc<ActiveTail>,
active_base: Position,
active_span: u64,
active_unindexable: bool,
}
impl IndexSet {
pub fn open(set: &SegmentSet) -> Result<Self, IndexError> {
let dir = set.dir().join("index");
if !dir.exists() {
fs::create_dir_all(&dir).map_err(|source| IndexError::io(&dir, source))?;
if let Some(parent) = dir.parent()
&& !parent.as_os_str().is_empty()
{
sync_dir(parent).map_err(|source| IndexError::io(parent, source))?;
}
}
let mut sealed = Vec::new();
for (base, count) in set.sealed_segments() {
let path = dir.join(index_file_name(base));
let entry = match load_valid(&path, base, count)? {
Some(seg) => SealedIndex {
base,
count,
seg: Some(Arc::new(seg)),
},
None => build_and_seal(set, &dir, base, count)?,
};
sealed.push(entry);
}
let active_base = set.active_base();
let active_count = set.next_position() - active_base;
let rebuilt = rebuild_range(set, active_base, active_count)?;
Ok(IndexSet {
dir,
sealed,
active: Arc::new(rebuilt.index),
active_base,
active_span: rebuilt.count,
active_unindexable: rebuilt.unindexable,
})
}
pub(crate) fn active_tail_arc(&self) -> Arc<ActiveTail> {
Arc::clone(&self.active)
}
pub(crate) fn sealed_index_arcs(&self) -> Vec<Option<Arc<IndexSegment>>> {
self.sealed.iter().map(|s| s.seg.clone()).collect()
}
pub fn push(&mut self, position: Position, event: EventRef<'_>) {
self.active_span += 1;
if self.active_unindexable {
return;
}
if self.active.push(position, event).is_err() {
self.active_unindexable = true;
self.active.mark_unindexable();
#[cfg(feature = "tracing")]
tracing::error!(
"segment at base {} exceeds the per-segment type limit; queries over it will error until it seals and is rebuilt",
self.active_base
);
}
}
pub fn seal_active(&mut self, new_base: Position) {
let base = self.active_base;
let count = self.active_span;
if self.active_unindexable {
self.sealed.push(SealedIndex {
base,
count,
seg: None,
});
} else {
debug_assert_eq!(
count,
self.active.len() as u64,
"active span tracks tail len"
);
let data: Arc<[u8]> = Arc::from(IndexSegment::encode(&self.active));
let seg = IndexSegment::from_bytes(Arc::clone(&data))
.expect("a just-encoded index segment must validate");
self.sealed.push(SealedIndex {
base,
count,
seg: Some(Arc::new(seg)),
});
let path = self.dir.join(index_file_name(base));
if let Err(err) = write_segment_file(&path, &data) {
#[cfg(feature = "tracing")]
tracing::error!("failed to persist index segment {path:?}: {err}");
let _ = err;
}
}
self.active = Arc::new(ActiveTail::new(new_base));
self.active_base = new_base;
self.active_span = 0;
self.active_unindexable = false;
}
fn plan_touched(&self, after: Position) -> Result<Vec<Touched<'_>>, IndexError> {
let mut plan: Vec<Touched<'_>> = Vec::new();
for entry in &self.sealed {
let Some(max) = entry.max_position() else {
continue;
};
if max <= after.get() {
continue; }
match &entry.seg {
Some(seg) => plan.push(Touched::Seg(seg.as_ref())),
None => {
return Err(IndexError::Unindexable {
range: PositionRange {
first: entry.base,
last: Position::new(max),
},
});
}
}
}
if self.active_span > 0 {
let max = self.active_base.get() + self.active_span - 1;
if max > after.get() {
if self.active_unindexable {
return Err(IndexError::Unindexable {
range: PositionRange {
first: self.active_base,
last: Position::new(max),
},
});
}
plan.push(Touched::Active(self.active.as_ref()));
}
}
Ok(plan)
}
pub fn search_all(
&self,
query: &Query,
after: Position,
) -> Result<IntoIter<Position>, IndexError> {
let mut out = Vec::new();
for touched in self.plan_touched(after)? {
match touched {
Touched::Seg(seg) => out.extend(search(seg, query, after)),
Touched::Active(tail) => out.extend(search(&tail.view_full(), query, after)),
}
}
Ok(out.into_iter())
}
pub fn find_match(
&self,
query: &Query,
after: Position,
) -> Result<Option<Position>, IndexError> {
for touched in self.plan_touched(after)? {
let first = match touched {
Touched::Seg(seg) => search(seg, query, after).next(),
Touched::Active(tail) => search(&tail.view_full(), query, after).next(),
};
if first.is_some() {
return Ok(first);
}
}
Ok(None)
}
#[cfg(test)]
pub(crate) fn force_unindexable_for_test(&mut self) {
for entry in &mut self.sealed {
entry.seg = None;
}
self.active_unindexable = true;
self.active.mark_unindexable();
}
}
enum Touched<'a> {
Seg(&'a IndexSegment),
Active(&'a ActiveTail),
}
impl fmt::Debug for IndexSet {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("IndexSet")
.field("dir", &self.dir)
.field("sealed", &self.sealed.len())
.field("active_base", &self.active_base)
.field("active_span", &self.active_span)
.field("active_unindexable", &self.active_unindexable)
.finish()
}
}
fn load_valid(path: &Path, base: Position, count: u64) -> Result<Option<IndexSegment>, IndexError> {
let bytes = match fs::read(path) {
Ok(bytes) => bytes,
Err(err) if err.kind() == io::ErrorKind::NotFound => return Ok(None),
Err(err) => return Err(IndexError::io(path, err)),
};
match IndexSegment::from_bytes(Arc::from(bytes)) {
Ok(seg) if seg.header().base_position == base && seg.header().event_count == count => {
Ok(Some(seg))
}
Ok(_) => {
#[cfg(feature = "tracing")]
tracing::warn!("index segment {path:?} disagrees with the log segment; rebuilding");
Ok(None)
}
Err(_err) => {
#[cfg(feature = "tracing")]
tracing::warn!("index segment {path:?} is corrupt ({_err}); rebuilding from the log");
Ok(None)
}
}
}
fn build_and_seal(
set: &SegmentSet,
dir: &Path,
base: Position,
count: u64,
) -> Result<SealedIndex, IndexError> {
let rebuilt = rebuild_range(set, base, count)?;
if rebuilt.unindexable {
#[cfg(feature = "tracing")]
tracing::error!(
"sealed segment at base {base} exceeds the per-segment type limit; queries over it will error"
);
return Ok(SealedIndex {
base,
count,
seg: None,
});
}
let data: Arc<[u8]> = Arc::from(IndexSegment::encode(&rebuilt.index));
let path = dir.join(index_file_name(base));
if let Err(err) = write_segment_file(&path, &data) {
#[cfg(feature = "tracing")]
tracing::error!("failed to persist rebuilt index segment {path:?}: {err}");
let _ = err;
}
let seg = IndexSegment::from_bytes(Arc::clone(&data))
.expect("a just-encoded index segment must validate");
Ok(SealedIndex {
base,
count,
seg: Some(Arc::new(seg)),
})
}
fn rebuild_range(
set: &SegmentSet,
base: Position,
count: u64,
) -> Result<super::recovery::Rebuilt, IndexError> {
let mut builder = Rebuilder::new(base);
let mut scan = set.scan_from(base);
let mut seen = 0u64;
while seen < count {
match scan.next() {
Some(item) => {
let record = item.map_err(|source| IndexError::Log(Arc::new(source)))?;
let event = EventRef::from_bytes(record.data).map_err(IndexError::Corrupt)?;
builder.feed(record.position, event);
seen += 1;
}
None => break,
}
}
Ok(builder.finish())
}
fn index_file_name(base: Position) -> String {
format!("{:0width$}.idx", base.get(), width = NAME_DIGITS)
}
fn sync_dir(dir: &Path) -> io::Result<()> {
File::open(dir)?.sync_all()
}
#[derive(Debug, Error)]
pub enum IndexError {
#[error("index i/o error at {path}: {source}")]
Io { path: PathBuf, source: io::Error },
#[error("log error while rebuilding the index: {0}")]
Log(Arc<LogError>),
#[error("corrupt event while rebuilding the index: {0}")]
Corrupt(DecodeError),
#[error(
"positions {}..={} cannot be answered from the index (a segment exceeds the per-segment type limit); scan the log for this range",
range.first,
range.last
)]
Unindexable { range: PositionRange },
}
impl IndexError {
fn io(path: &Path, source: io::Error) -> Self {
IndexError::Io {
path: path.to_path_buf(),
source,
}
}
}
const _: fn() = || {
fn is_send<T: Send>() {}
fn is_sync<T: Sync>() {}
is_send::<IndexSet>();
is_sync::<IndexSet>();
};
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{Event, EventType, Tag, Tags};
use crate::log::set::SegmentConfig;
use crate::query::{Matches, 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()
}
fn templates() -> Vec<Event> {
vec![
event("Registered", &["course:c1"]),
event("Enrolled", &["course:c1", "student:s1"]),
event("Renamed", &["student:s1"]),
event("Registered", &["course:c2"]),
event("Enrolled", &["course:c2", "student:s2"]),
event("Renamed", &["student:s2"]),
event("Registered", &["course:c1"]),
event("Enrolled", &["course:c1", "student:s2"]),
]
}
fn build_log(dir: &Path) -> SegmentSet {
let mut set = SegmentSet::open(dir, SegmentConfig::new(512)).unwrap();
let templates = templates();
for i in 0..40 {
set.append_batch(&[templates[i % templates.len()].as_bytes()])
.unwrap();
}
set
}
fn scan_baseline(set: &SegmentSet, query: &Query, after: Position) -> Vec<Position> {
let mut out = Vec::new();
let mut scan = set.scan_after(after);
while let Some(item) = scan.next() {
let record = item.unwrap();
let event = EventRef::from_bytes(record.data).unwrap();
if query.matches(event) {
out.push(record.position);
}
}
out
}
#[test]
fn open_rebuilds_and_answers_across_multiple_segments() {
let dir = TempDir::new().unwrap();
let set = build_log(dir.path());
assert!(
set.sealed_len() >= 1,
"small segments should have sealed some"
);
let index = IndexSet::open(&set).unwrap();
let queries = [
Query::all(),
Query::item(QueryItem::with_tags(tags(&["course:c1"]))),
Query::item(QueryItem::with_tags(tags(&["student:s2"]))),
Query::item(QueryItem::of_types(vec![
EventType::new("Enrolled").unwrap(),
])),
Query::items(vec![
QueryItem::with_tags(tags(&["course:c1"])),
QueryItem::with_tags(tags(&["student:s2"])),
]),
];
let last = set.last_position().get();
for query in &queries {
for after in 0..=last {
let from_index: Vec<Position> = index
.search_all(query, Position::new(after))
.unwrap()
.collect();
let from_scan = scan_baseline(&set, query, Position::new(after));
assert_eq!(from_index, from_scan, "query {query:?} after {after}");
}
}
}
#[test]
fn reopen_uses_persisted_segments() {
let dir = TempDir::new().unwrap();
let set = build_log(dir.path());
drop(IndexSet::open(&set).unwrap());
let idx_dir = dir.path().join("index");
let idx_count = fs::read_dir(&idx_dir).unwrap().count();
assert_eq!(idx_count, set.sealed_len());
let index = IndexSet::open(&set).unwrap();
let q = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
let from_index: Vec<Position> = index.search_all(&q, Position::ZERO).unwrap().collect();
let from_scan = scan_baseline(&set, &q, Position::ZERO);
assert_eq!(from_index, from_scan);
}
#[test]
fn deleted_idx_is_rebuilt() {
let dir = TempDir::new().unwrap();
let set = build_log(dir.path());
drop(IndexSet::open(&set).unwrap());
let idx_dir = dir.path().join("index");
let first = fs::read_dir(&idx_dir)
.unwrap()
.filter_map(Result::ok)
.map(|e| e.path())
.min()
.unwrap();
fs::remove_file(&first).unwrap();
let index = IndexSet::open(&set).unwrap();
let q = Query::all();
let from_index: Vec<Position> = index.search_all(&q, Position::ZERO).unwrap().collect();
let from_scan = scan_baseline(&set, &q, Position::ZERO);
assert_eq!(from_index, from_scan);
}
#[test]
fn corrupt_idx_body_is_rebuilt() {
let dir = TempDir::new().unwrap();
let set = build_log(dir.path());
drop(IndexSet::open(&set).unwrap());
let idx_dir = dir.path().join("index");
let first = fs::read_dir(&idx_dir)
.unwrap()
.filter_map(Result::ok)
.map(|e| e.path())
.min()
.unwrap();
let mut bytes = fs::read(&first).unwrap();
let last = bytes.len() - 1;
bytes[last] ^= 0xFF;
fs::write(&first, &bytes).unwrap();
let index = IndexSet::open(&set).unwrap();
let q = Query::all();
let from_index: Vec<Position> = index.search_all(&q, Position::ZERO).unwrap().collect();
let from_scan = scan_baseline(&set, &q, Position::ZERO);
assert_eq!(from_index, from_scan);
}
#[test]
fn unindexable_segment_errors_when_touched_but_not_when_pruned() {
let dir = TempDir::new().unwrap();
let index = IndexSet {
dir: dir.path().to_path_buf(),
sealed: vec![SealedIndex {
base: Position::new(1),
count: 3,
seg: None,
}],
active: Arc::new(super::super::ActiveTail::new(Position::new(4))),
active_base: Position::new(4),
active_span: 0,
active_unindexable: false,
};
match index.search_all(&Query::all(), Position::ZERO) {
Err(IndexError::Unindexable { range }) => {
assert_eq!(range.first, Position::new(1));
assert_eq!(range.last, Position::new(3));
}
other => panic!("expected Unindexable error, got {other:?}"),
}
let got: Vec<Position> = index
.search_all(&Query::all(), Position::new(3))
.unwrap()
.collect();
assert!(got.is_empty());
}
#[test]
fn live_feed_and_seal_agree_with_scan() {
let dir = TempDir::new().unwrap();
let mut set = SegmentSet::open(dir.path(), SegmentConfig::new(512)).unwrap();
let mut index = IndexSet::open(&set).unwrap();
let templates = templates();
for i in 0..40 {
let ev = &templates[i % templates.len()];
let sealed_before = set.sealed_len();
let range = set.append_batch(&[ev.as_bytes()]).unwrap();
if set.sealed_len() > sealed_before {
index.seal_active(range.first);
}
index.push(range.first, ev.as_ref());
}
assert!(set.sealed_len() >= 1, "should have rolled over");
let last = set.last_position().get();
let q = Query::item(QueryItem::with_tags(tags(&["course:c1"])));
for after in 0..=last {
let from_index: Vec<Position> = index
.search_all(&q, Position::new(after))
.unwrap()
.collect();
let from_scan = scan_baseline(&set, &q, Position::new(after));
assert_eq!(from_index, from_scan, "after {after}");
}
}
#[test]
fn find_match_is_the_first_position_the_scan_would_return() {
let dir = TempDir::new().unwrap();
let set = build_log(dir.path());
let index = IndexSet::open(&set).unwrap();
let queries = [
Query::all(),
Query::item(QueryItem::with_tags(tags(&["course:c1"]))),
Query::item(QueryItem::with_tags(tags(&["student:s2"]))),
Query::item(QueryItem::of_types(vec![
EventType::new("Enrolled").unwrap(),
])),
Query::items(vec![
QueryItem::with_tags(tags(&["course:c1"])),
QueryItem::with_tags(tags(&["student:s2"])),
]),
Query::item(QueryItem::with_tags(tags(&["ghost:x"]))),
];
let last = set.last_position().get();
for query in &queries {
for after in 0..=last {
let found = index.find_match(query, Position::new(after)).unwrap();
let expected = scan_baseline(&set, query, Position::new(after))
.into_iter()
.next();
assert_eq!(found, expected, "query {query:?} after {after}");
}
}
}
#[test]
fn find_match_prunes_and_errors_like_search_all() {
let dir = TempDir::new().unwrap();
let index = IndexSet {
dir: dir.path().to_path_buf(),
sealed: vec![SealedIndex {
base: Position::new(1),
count: 3,
seg: None,
}],
active: Arc::new(super::super::ActiveTail::new(Position::new(4))),
active_base: Position::new(4),
active_span: 0,
active_unindexable: false,
};
match index.find_match(&Query::all(), Position::ZERO) {
Err(IndexError::Unindexable { range }) => {
assert_eq!(range.first, Position::new(1));
assert_eq!(range.last, Position::new(3));
}
other => panic!("expected Unindexable error, got {other:?}"),
}
assert_eq!(
index.find_match(&Query::all(), Position::new(3)).unwrap(),
None
);
}
}