use std::borrow::Cow;
use std::fmt;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU16, Ordering};
use dashmap::DashMap;
use thiserror::Error;
use crate::Position;
use crate::event::EventRef;
use super::append::{ChunkedVec, PostingSlot, Snap};
use super::{SegmentIndex, TermId, TypeId};
pub struct ActiveTail {
base: Position,
type_column: ChunkedVec<AtomicU16>,
postings: ChunkedVec<PostingSlot>,
tags: Arc<DashMap<Arc<str>, TermId>>,
types: Arc<DashMap<Arc<str>, TypeId>>,
unindexable: AtomicBool,
}
impl fmt::Debug for ActiveTail {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ActiveTail")
.field("base", &self.base)
.field("len", &self.len())
.field("terms", &self.postings.len())
.field("types", &self.types.len())
.finish()
}
}
impl ActiveTail {
pub fn new(base: Position) -> Self {
ActiveTail {
base,
type_column: ChunkedVec::new(),
postings: ChunkedVec::new(),
tags: Arc::new(DashMap::new()),
types: Arc::new(DashMap::new()),
unindexable: AtomicBool::new(false),
}
}
pub fn push(&self, position: Position, event: EventRef<'_>) -> Result<(), TooManyTypes> {
let len = self.type_column.len();
let expected = Position::new(self.base.get() + len as u64);
assert_eq!(
position, expected,
"tail index fed out of order: expected {expected}, got {position}"
);
let type_id = self.intern_type(event.event_type())?;
let local = len;
for tag in event.tags() {
let term = self.intern_tag(tag);
self.postings
.with(term.get(), |slot| slot.push_local(local));
}
self.type_column
.push_with(|slot| slot.store(type_id.get(), Ordering::Relaxed));
Ok(())
}
fn intern_tag(&self, tag: &str) -> TermId {
if let Some(id) = self.tags.get(tag).map(|r| *r) {
return id;
}
let id = TermId(self.postings.push_with(|_| {}));
self.tags.insert(Arc::from(tag), id);
id
}
fn intern_type(&self, name: &str) -> Result<TypeId, TooManyTypes> {
if let Some(id) = self.types.get(name).map(|r| *r) {
return Ok(id);
}
let next = self.types.len();
if next > u16::MAX as usize {
return Err(TooManyTypes {
max: u16::MAX as usize + 1,
});
}
let id = TypeId(next as u16);
self.types.insert(Arc::from(name), id);
Ok(id)
}
pub fn base(&self) -> Position {
self.base
}
pub fn len(&self) -> u32 {
self.type_column.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub(crate) fn mark_unindexable(&self) {
self.unindexable.store(true, Ordering::Release);
}
pub fn is_unindexable(&self) -> bool {
self.unindexable.load(Ordering::Acquire)
}
pub fn view(&self, watermark: Position) -> ActiveView {
self.make_view(self.visible_len(watermark))
}
pub fn view_full(&self) -> ActiveView {
self.make_view(self.len())
}
fn visible_len(&self, watermark: Position) -> u32 {
if watermark.get() >= self.base.get() {
(watermark.get() - self.base.get() + 1).min(u32::MAX as u64) as u32
} else {
0
}
}
fn make_view(&self, upto_len: u32) -> ActiveView {
let type_column = self.type_column.snapshot();
let postings = self.postings.snapshot();
let upto_len = upto_len.min(type_column.covered());
ActiveView {
base: self.base,
upto_len,
type_column,
postings,
tags: Arc::clone(&self.tags),
types: Arc::clone(&self.types),
}
}
pub(crate) fn type_column(&self) -> Vec<u16> {
let snap = self.type_column.snapshot();
(0..self.len())
.map(|i| snap.get(i).map(|s| s.load(Ordering::Relaxed)).unwrap_or(0))
.collect()
}
pub(crate) fn type_names(&self) -> Vec<Arc<str>> {
let mut names: Vec<Arc<str>> = vec![Arc::from(""); self.types.len()];
for entry in self.types.iter() {
names[entry.value().get() as usize] = Arc::clone(entry.key());
}
names
}
pub(crate) fn terms_sorted_with_postings(&self) -> Vec<(Arc<str>, Vec<u32>)> {
let snap = self.postings.snapshot();
let upto = self.len();
let mut terms: Vec<(Arc<str>, Vec<u32>)> = self
.tags
.iter()
.map(|entry| {
let mut postings = Vec::new();
if let Some(slot) = snap.get(entry.value().get()) {
slot.collect_below(upto, &mut postings);
}
(Arc::clone(entry.key()), postings)
})
.collect();
terms.sort_by(|a, b| a.0.cmp(&b.0));
terms
}
}
pub struct ActiveView {
base: Position,
upto_len: u32,
type_column: Snap<AtomicU16>,
postings: Snap<PostingSlot>,
tags: Arc<DashMap<Arc<str>, TermId>>,
types: Arc<DashMap<Arc<str>, TypeId>>,
}
impl SegmentIndex for ActiveView {
fn base(&self) -> Position {
self.base
}
fn len(&self) -> u32 {
self.upto_len
}
fn term_postings(&self, tag: &str) -> Option<Cow<'_, [u32]>> {
let term = self.tags.get(tag).map(|r| *r)?;
let slot = self.postings.get(term.get())?;
let mut out = Vec::new();
slot.collect_below(self.upto_len, &mut out);
Some(Cow::Owned(out))
}
fn term_len(&self, tag: &str) -> Option<u32> {
let term = self.tags.get(tag).map(|r| *r)?;
let slot = self.postings.get(term.get())?;
Some(slot.len().min(self.upto_len))
}
fn type_id(&self, name: &str) -> Option<u16> {
self.types.get(name).map(|r| r.get())
}
fn type_at(&self, local: u32) -> u16 {
self.type_column
.get(local)
.map(|s| s.load(Ordering::Relaxed))
.unwrap_or(0)
}
}
#[derive(Clone, Copy, Debug, Error, PartialEq, Eq)]
#[error("too many distinct event types in one segment (maximum {max})")]
pub struct TooManyTypes {
pub max: usize,
}
const _: fn() = || {
fn is_send<T: Send>() {}
fn is_sync<T: Sync>() {}
is_send::<ActiveTail>();
is_sync::<ActiveTail>();
is_send::<ActiveView>();
is_sync::<ActiveView>();
};
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{Event, EventType, Tag, Tags};
use crate::index::search;
use crate::query::{Query, QueryItem};
use smallvec::SmallVec;
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"").unwrap()
}
fn build(base: u64, events: &[Event]) -> ActiveTail {
let index = ActiveTail::new(Position::new(base));
for (i, ev) in events.iter().enumerate() {
index
.push(Position::new(base + i as u64), ev.as_ref())
.unwrap();
}
index
}
fn postings(index: &ActiveTail, tag: &str) -> Option<Vec<u32>> {
index.view_full().term_postings(tag).map(|c| c.into_owned())
}
#[test]
fn postings_are_ascending_and_correct() {
let events = [
event("E", &["a"]),
event("E", &["b"]),
event("E", &["a", "b"]),
event("E", &["c"]),
];
let index = build(1, &events);
assert_eq!(index.len(), 4);
assert_eq!(postings(&index, "a"), Some(vec![0, 2]));
assert_eq!(postings(&index, "b"), Some(vec![1, 2]));
assert_eq!(postings(&index, "absent"), None);
}
#[test]
fn type_column_tracks_each_event() {
let events = [
event("Registered", &["a"]),
event("Enrolled", &["b"]),
event("Registered", &["c"]),
];
let index = build(1, &events);
let view = index.view_full();
let registered = view.type_id("Registered").unwrap();
let enrolled = view.type_id("Enrolled").unwrap();
assert_ne!(registered, enrolled);
assert_eq!(view.type_at(0), registered);
assert_eq!(view.type_at(1), enrolled);
assert_eq!(view.type_at(2), registered);
assert_eq!(view.type_id("Missing"), None);
}
#[test]
fn base_offsets_local_positions() {
let index = build(100, &[event("E", &["a"])]);
assert_eq!(index.base(), Position::new(100));
assert_eq!(postings(&index, "a"), Some(vec![0]));
}
#[test]
fn view_is_bounded_by_the_watermark() {
let events: Vec<Event> = (0..5).map(|_| event("E", &["a"])).collect();
let index = build(1, &events);
let view = index.view(Position::new(3));
assert_eq!(view.len(), 3);
assert_eq!(view.term_postings("a").unwrap().into_owned(), vec![0, 1, 2]);
let q = Query::item(QueryItem::with_tags(tags(&["a"])));
let got: Vec<u64> = search(&view, &q, Position::ZERO).map(|p| p.get()).collect();
assert_eq!(got, vec![1, 2, 3]);
}
#[test]
fn term_len_is_a_watermark_bounded_upper_bound() {
let events: Vec<Event> = (0..5).map(|_| event("E", &["a"])).collect();
let index = build(1, &events);
assert_eq!(index.view_full().term_len("a"), Some(5));
assert_eq!(index.view(Position::new(3)).term_len("a"), Some(3));
assert_eq!(index.view_full().term_len("absent"), None);
}
#[test]
#[should_panic(expected = "fed out of order")]
fn out_of_order_push_panics() {
let index = ActiveTail::new(Position::new(1));
index
.push(Position::new(1), event("E", &["a"]).as_ref())
.unwrap();
index
.push(Position::new(3), event("E", &["b"]).as_ref())
.unwrap();
}
}