use std::path::Path;
use crate::common::{PageID, Position};
use crate::events_tree::{EventIterator, event_tree_append, event_tree_lookup};
use crate::events_tree_nodes::EventRecord;
use crate::mvcc::{Mvcc, Writer};
use crate::node::Node;
use crate::page::Page;
use crate::tags_tree::{TagsTreeIterator, tags_tree_insert};
use crate::tags_tree_nodes::{TagHash, get_tag_key_width};
use crate::tracking_tree_nodes::{TrackingInternalNode, TrackingLeafNode};
use itertools::Itertools;
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
use umadb_dcb::{
DcbAppendCondition, DcbEvent, DcbEventStoreSync, DcbQuery, DcbReadResponseSync, DcbResult,
DcbSequencedEvent, DcbError, TrackingInfo,
};
use uuid::Uuid;
pub static DEFAULT_PAGE_SIZE: usize = 4096;
pub const DEFAULT_DB_FILENAME: &str = "uma.db";
pub const DB_SCHEMA_VERSION: u32 = 1;
pub struct UmaDb {
pub mvcc: Arc<Mvcc>,
}
impl UmaDb {
pub fn new<P: AsRef<Path>>(path: P) -> DcbResult<Self> {
let p = path.as_ref();
let file_path = if p.is_dir() {
p.join(DEFAULT_DB_FILENAME)
} else {
p.to_path_buf()
};
let mvcc = Mvcc::new(&file_path, DEFAULT_PAGE_SIZE, false)?;
Ok(Self {
mvcc: Arc::new(mvcc),
})
}
pub fn from_arc(mvcc: Arc<Mvcc>) -> Self {
Self { mvcc }
}
pub fn get_tracking_info(&self, source: &str) -> DcbResult<Option<u64>> {
let reader = self.mvcc.reader()?;
let mut pid = reader.tracking_tree_root_id;
if pid == PageID(0) {
return Ok(None);
}
loop {
let page = self.mvcc.read_page(pid)?;
match &page.node {
Node::TrackingLeaf(node) => return Ok(node.get(source).map(|p| p.0)),
Node::TrackingInternal(internal) => {
let idx = match internal.keys.binary_search_by(|k| k.as_str().cmp(source)) {
Ok(i) => i + 1,
Err(i) => i,
};
if idx >= internal.child_ids.len() {
return Err(DcbError::DatabaseCorrupted(
"tracking internal child index out of bounds".to_string(),
));
}
pid = internal.child_ids[idx];
}
other => {
return Err(DcbError::DatabaseCorrupted(format!(
"Invalid tracking node type: {}",
other.type_name()
)));
}
}
}
}
pub fn append_batch(
&self,
mut items: Vec<(
Vec<DcbEvent>,
Option<DcbAppendCondition>,
Option<TrackingInfo>,
)>,
) -> DcbResult<Vec<DcbResult<u64>>> {
let total = items.len();
let mvcc = &self.mvcc;
let mut writer = mvcc.writer()?;
let mut results: Vec<DcbResult<u64>> = Vec::with_capacity(total);
let mut abort_idx: Option<usize> = None;
let mut abort_err: Option<DcbError> = None;
for (idx, (events, condition, tracking)) in items.drain(..).enumerate() {
if abort_idx.is_some() {
break;
}
let res = Self::process_append_request(
events,
condition,
tracking,
mvcc,
&mut writer,
None,
);
match &res {
Ok(_) => results.push(res),
Err(e) if is_integrity_error(e) => results.push(Err(clone_dcb_error(e))),
Err(e) => {
abort_idx = Some(idx);
abort_err = Some(clone_dcb_error(e));
results.push(Err(clone_dcb_error(e))); }
}
}
if let (Some(failed_at), Some(orig_err)) = (abort_idx, abort_err) {
let shadow = shadow_for_batch_abort(&orig_err);
for i in 0..results.len() {
if i != failed_at {
results[i] = Err(clone_dcb_error(&shadow));
}
}
while results.len() < total {
results.push(Err(clone_dcb_error(&shadow)));
}
return Ok(results);
}
mvcc.commit(&mut writer)?;
Ok(results)
}
pub fn process_append_request(
events: Vec<DcbEvent>,
condition: Option<DcbAppendCondition>,
tracking_info: Option<TrackingInfo>,
mvcc: &Arc<Mvcc>,
writer: &mut Writer,
cancel: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> DcbResult<u64> {
if let Some(cond) = condition {
let from = cond.after.map(|after| Position(after + 1));
let read_result1 = read_conditional(
mvcc,
&writer.dirty,
writer.events_tree_root_id,
writer.tags_tree_root_id,
cond.fail_if_events_match.clone(),
from,
false,
Some(1),
false,
cancel.clone(),
);
match read_result1 {
Ok(found_vec) => {
if let Some(matched) = found_vec.first() {
return match is_request_idempotent(
mvcc,
&writer.dirty,
writer.events_tree_root_id,
writer.tags_tree_root_id,
&events,
cond.fail_if_events_match.clone(),
from,
cancel.clone(),
) {
Ok(Some(last_recorded_position)) => Ok(last_recorded_position),
Ok(None) => {
let msg = format!(
"condition: {:?} matched: {:?}, ",
cond.clone(),
matched,
);
Err(DcbError::IntegrityError(msg))
}
Err(err) => {
Err(err)
}
};
}
}
Err(e) => {
return Err(e);
}
}
}
if let Some(tracking_info) = tracking_info
&& let Err(e) = tracking_upsert(
mvcc,
writer,
&tracking_info.source,
Position(tracking_info.position),
)
{
return Err(e);
}
if events.is_empty() {
return Ok(0);
}
match unconditional_append(mvcc, writer, events) {
Ok(last) => Ok(last),
Err(e) => Err(e),
}
}
}
impl DcbEventStoreSync for UmaDb {
fn read(
&self,
query: Option<DcbQuery>,
start: Option<u64>,
backwards: bool,
limit: Option<u32>,
_subscribe: bool, ) -> DcbResult<Box<dyn DcbReadResponseSync + Send + 'static>> {
let mvcc = &self.mvcc;
let reader = mvcc.reader()?;
let last_committed_position = reader.next_position.0.saturating_sub(1);
let q = query.unwrap_or(DcbQuery { items: vec![] });
let from = start.map(Position);
let events = read_conditional(
mvcc,
&HashMap::new(),
reader.events_tree_root_id,
reader.tags_tree_root_id,
q,
from,
backwards,
limit,
false,
None,
)?;
let head = if limit.is_none() {
if last_committed_position == 0 {
None
} else {
Some(last_committed_position)
}
} else {
events.last().map(|e| e.position)
};
Ok(Box::new(ReadResponse {
events: VecDeque::from(events),
head,
}))
}
fn head(&self) -> DcbResult<Option<u64>> {
let db = &self.mvcc;
let header_page = db.get_latest_header_page()?;
let header = header_page.as_header_node()?;
let last = header.next_position.0.saturating_sub(1);
if last == 0 { Ok(None) } else { Ok(Some(last)) }
}
fn get_tracking_info(&self, source: &str) -> DcbResult<Option<u64>> {
UmaDb::get_tracking_info(self, source)
}
fn append(
&self,
events: Vec<DcbEvent>,
condition: Option<DcbAppendCondition>,
tracking_info: Option<TrackingInfo>,
) -> DcbResult<u64> {
if events.is_empty() {
return Ok(0);
}
let mvcc = &self.mvcc;
let mut writer = mvcc.writer()?;
let result = Self::process_append_request(
events,
condition,
tracking_info,
mvcc,
&mut writer,
None,
);
mvcc.commit(&mut writer)?;
result
}
}
struct ReadResponse {
events: VecDeque<DcbSequencedEvent>,
head: Option<u64>,
}
fn tracking_upsert(mvcc: &Mvcc, writer: &mut Writer, source: &str, pos: Position) -> DcbResult<()> {
let key_len = source.len();
if key_len > u8::MAX as usize {
return Err(DcbError::InvalidArgument(format!(
"tracking source too long ({} > 255)",
key_len
)));
}
let root = writer.tracking_tree_root_id;
if root == PageID(0) {
let mut node = TrackingLeafNode::new();
node.keys.push(source.to_string());
node.values.push(pos);
let new_root_id = writer.alloc_page_id();
let page = Page::new(new_root_id, Node::TrackingLeaf(node));
writer.insert_dirty(page)?;
writer.tracking_tree_root_id = new_root_id;
return Ok(());
}
let mut stack: Vec<(PageID, usize)> = Vec::new();
let mut current_id = root;
loop {
let page = writer.get_page_ref(mvcc, current_id)?;
match &page.node {
Node::TrackingLeaf(_) => break,
Node::TrackingInternal(internal) => {
let child_idx = internal.child_index_for_key(source);
if child_idx >= internal.child_ids.len() {
return Err(DcbError::DatabaseCorrupted(
"tracking internal child index out of bounds".to_string(),
));
}
let next = internal.child_ids[child_idx];
stack.push((current_id, child_idx));
current_id = next;
}
other => {
return Err(DcbError::DatabaseCorrupted(format!(
"Invalid tracking node type: {}",
other.type_name()
)));
}
}
}
{
let page = writer.get_page_ref(mvcc, current_id)?;
let Node::TrackingLeaf(leaf) = &page.node else {
return Err(DcbError::DatabaseCorrupted(
"Expected TrackingLeaf".to_string(),
));
};
if let Some(existing) = leaf.get(source)
&& pos.0 <= existing.0
{
return Err(DcbError::IntegrityError(format!(
"non-increasing tracking position for source '{source}': {} <= {}",
pos.0, existing.0
)));
}
}
let dirty_leaf_id = writer.get_dirty_page_id(current_id)?;
let mut replacement_info: Option<(PageID, PageID)> =
(dirty_leaf_id != current_id).then_some((current_id, dirty_leaf_id));
let mut split_info: Option<(String, PageID)> = None;
{
let leaf_page = writer.get_mut_dirty(dirty_leaf_id)?;
{
let Node::TrackingLeaf(ref mut node) = leaf_page.node else {
return Err(DcbError::DatabaseCorrupted(
"Dirty tracking page not a leaf".to_string(),
));
};
match node.keys.binary_search_by(|k| k.as_str().cmp(source)) {
Ok(i) => node.values[i] = pos,
Err(ins) => {
node.keys.insert(ins, source.to_string());
node.values.insert(ins, pos);
}
}
}
if leaf_page.calc_serialized_size() > mvcc.page_size {
let promoted_key: String;
let right_id: PageID;
{
let Node::TrackingLeaf(ref mut node) = leaf_page.node else {
return Err(DcbError::DatabaseCorrupted(
"Dirty tracking page not a leaf".to_string(),
));
};
if node.keys.len() < 2 {
return Err(DcbError::DatabaseCorrupted(
"Cannot split tracking leaf with too few keys".to_string(),
));
}
let mid = node.keys.len() / 2;
let right_keys = node.keys.split_off(mid);
let right_vals = node.values.split_off(mid);
promoted_key = right_keys
.first()
.ok_or_else(|| DcbError::DatabaseCorrupted("empty right split".to_string()))?
.clone();
let right_leaf = TrackingLeafNode {
keys: right_keys,
values: right_vals,
};
right_id = writer.alloc_page_id();
let right_page = Page::new(right_id, Node::TrackingLeaf(right_leaf));
writer.insert_dirty(right_page)?;
}
split_info = Some((promoted_key, right_id));
}
}
while let Some((parent_id, child_idx)) = stack.pop() {
let dirty_parent_id = writer.get_dirty_page_id(parent_id)?;
let parent_replacement_info =
(dirty_parent_id != parent_id).then_some((parent_id, dirty_parent_id));
if let Some((old_id, new_id)) = replacement_info.take() {
let parent_page = writer.get_mut_dirty(dirty_parent_id)?;
let Node::TrackingInternal(ref mut internal) = parent_page.node else {
return Err(DcbError::DatabaseCorrupted(
"Expected TrackingInternal".to_string(),
));
};
internal.replace_child_id_at(child_idx, old_id, new_id)?;
}
if let Some((prom_key, new_child_id)) = split_info.take() {
let need_split: bool;
{
let parent_page = writer.get_mut_dirty(dirty_parent_id)?;
let Node::TrackingInternal(ref mut internal) = parent_page.node else {
return Err(DcbError::DatabaseCorrupted(
"Expected TrackingInternal".to_string(),
));
};
internal.insert_promoted_at(child_idx, prom_key, new_child_id);
need_split = parent_page.calc_serialized_size() > mvcc.page_size;
}
if need_split {
let parent_page = writer.get_mut_dirty(dirty_parent_id)?;
let Node::TrackingInternal(ref mut internal) = parent_page.node else {
return Err(DcbError::DatabaseCorrupted(
"Expected TrackingInternal".to_string(),
));
};
let (promote_up, right_keys, right_child_ids) = internal.split_off()?;
let right_internal = TrackingInternalNode {
keys: right_keys,
child_ids: right_child_ids,
};
let right_internal_id = writer.alloc_page_id();
let right_internal_page =
Page::new(right_internal_id, Node::TrackingInternal(right_internal));
writer.insert_dirty(right_internal_page)?;
split_info = Some((promote_up, right_internal_id));
}
}
replacement_info = parent_replacement_info;
}
if let Some((old_id, new_id)) = replacement_info.take() {
if writer.tracking_tree_root_id == old_id {
writer.tracking_tree_root_id = new_id;
} else {
return Err(DcbError::RootIDMismatch(old_id.0, new_id.0));
}
}
if let Some((prom_key, right_id)) = split_info.take() {
let new_root_id = writer.alloc_page_id();
let left_id = writer.tracking_tree_root_id;
let new_root = TrackingInternalNode {
keys: vec![prom_key],
child_ids: vec![left_id, right_id],
};
let new_root_page = Page::new(new_root_id, Node::TrackingInternal(new_root));
writer.insert_dirty(new_root_page)?;
writer.tracking_tree_root_id = new_root_id;
}
Ok(())
}
impl Iterator for ReadResponse {
type Item = DcbResult<DcbSequencedEvent>;
fn next(&mut self) -> Option<Self::Item> {
self.events.pop_front().map(Ok)
}
}
impl DcbReadResponseSync for ReadResponse {
fn head(&mut self) -> DcbResult<Option<u64>> {
Ok(self.head)
}
fn collect_with_head(&mut self) -> DcbResult<(Vec<DcbSequencedEvent>, Option<u64>)> {
let events = self.events.drain(..).collect();
Ok((events, self.head))
}
fn next_batch(&mut self) -> DcbResult<Vec<DcbSequencedEvent>> {
let batch = self.events.drain(..).collect();
Ok(batch)
}
}
pub fn unconditional_append(
mvcc: &Mvcc,
writer: &mut Writer,
events: Vec<DcbEvent>,
) -> DcbResult<u64> {
let mut last_pos_u64: u64 = 0;
for ev in events.into_iter() {
let position = writer.issue_position();
last_pos_u64 = position.0;
for tag in ev.tags.iter() {
let tag_hash: TagHash = tag_to_hash(tag);
tags_tree_insert(mvcc, writer, tag_hash, position)?;
}
let record = EventRecord {
event_type: ev.event_type,
data: ev.data,
tags: ev.tags,
uuid: ev.uuid,
};
event_tree_append(mvcc, writer, record, position)?;
}
Ok(last_pos_u64)
}
pub fn read_conditional(
mvcc: &Mvcc,
dirty: &HashMap<PageID, Page>,
events_tree_root_id: PageID,
tags_tree_root_id: PageID,
query: DcbQuery,
start: Option<Position>,
backwards: bool,
limit: Option<u32>,
force_sequential_read: bool,
cancel: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> DcbResult<Vec<DcbSequencedEvent>> {
const SCAN_BATCH_SIZE: u32 = 256;
if let Some(0) = limit {
return Ok(Vec::new());
}
if query.items.is_empty() {
let mut iter = EventIterator::new(mvcc, dirty, events_tree_root_id, start, backwards);
let mut out: Vec<DcbSequencedEvent> = Vec::new();
'outer_all: loop {
if let Some(ref c) = cancel {
if c.load(std::sync::atomic::Ordering::Relaxed) {
return Err(DcbError::CancelledByUser());
}
}
let batch = iter.next_batch(limit.unwrap_or(SCAN_BATCH_SIZE), cancel.as_ref())?;
if batch.is_empty() {
break;
}
for (pos, rec) in batch.into_iter() {
out.push(DcbSequencedEvent {
position: pos.0,
event: DcbEvent {
event_type: rec.event_type,
data: rec.data,
tags: rec.tags,
uuid: rec.uuid,
},
});
if let Some(lim) = limit
&& out.len() >= lim as usize
{
break 'outer_all;
}
}
}
return Ok(out);
}
let all_items_have_tags = query.items.iter().all(|it| !it.tags.is_empty());
if !all_items_have_tags || force_sequential_read {
let mut iter = EventIterator::new(mvcc, dirty, events_tree_root_id, start, backwards);
let mut out: Vec<DcbSequencedEvent> = Vec::new();
let matches_item = |rec: &EventRecord| -> bool {
for item in &query.items {
let type_ok =
item.types.is_empty() || item.types.iter().any(|t| t == &rec.event_type);
if !type_ok {
continue;
}
let tags_ok = item.tags.iter().all(|t| rec.tags.iter().any(|et| et == t));
if type_ok && tags_ok {
return true;
}
}
false
};
'outer_fallback: loop {
if let Some(ref c) = cancel {
if c.load(std::sync::atomic::Ordering::Relaxed) {
return Err(DcbError::CancelledByUser());
}
}
let batch = iter.next_batch(SCAN_BATCH_SIZE, cancel.as_ref())?;
if batch.is_empty() {
break;
}
for (pos, rec) in batch.into_iter() {
if matches_item(&rec) {
out.push(DcbSequencedEvent {
position: pos.0,
event: DcbEvent {
event_type: rec.event_type,
data: rec.data,
tags: rec.tags,
uuid: rec.uuid,
},
});
if let Some(lim) = limit
&& out.len() >= lim as usize
{
break 'outer_fallback;
}
}
}
}
return Ok(out);
}
let mut tag_qiis: HashMap<String, Vec<usize>> = HashMap::with_capacity(query.items.len() * 2);
let mut qi_tags: Vec<HashSet<String>> = Vec::with_capacity(query.items.len());
for (qiid, item) in query.items.iter().enumerate() {
qi_tags.push(item.tags.iter().cloned().collect());
for tag in &item.tags {
tag_qiis.entry(tag.clone()).or_default().push(qiid);
}
}
struct PositionTagQiidIterator<I>
where
I: Iterator<Item = Position>,
{
inner: I,
tag: String,
qiids: Vec<usize>,
}
impl<I> PositionTagQiidIterator<I>
where
I: Iterator<Item = Position>,
{
fn new(inner: I, tag: String, qiids: Vec<usize>) -> Self {
Self { inner, tag, qiids }
}
}
impl<I> Iterator for PositionTagQiidIterator<I>
where
I: Iterator<Item = Position>,
{
type Item = (Position, String, Vec<usize>);
fn next(&mut self) -> Option<Self::Item> {
self.inner
.next()
.map(|p| (p, self.tag.clone(), self.qiids.clone()))
}
}
let mut tag_iters: Vec<PositionTagQiidIterator<_>> = Vec::new();
for (tag, qiids) in tag_qiis.iter() {
let tag_hash: TagHash = tag_to_hash(tag);
let positions_iter =
TagsTreeIterator::new(mvcc, dirty, tags_tree_root_id, tag_hash, start, backwards); tag_iters.push(PositionTagQiidIterator::new(
positions_iter,
tag.clone(),
qiids.clone(),
));
}
let merged = tag_iters
.into_iter()
.kmerge_by(|a, b| if !backwards { a.0 < b.0 } else { a.0 > b.0 });
struct GroupByPositionIterator<I>
where
I: Iterator<Item = (Position, String, Vec<usize>)>,
{
inner: I,
current_pos: Option<Position>,
tags: HashSet<String>,
qiis: HashSet<usize>,
finished: bool,
}
impl<I> GroupByPositionIterator<I>
where
I: Iterator<Item = (Position, String, Vec<usize>)>,
{
fn new(inner: I) -> Self {
Self {
inner,
current_pos: None,
tags: HashSet::new(),
qiis: HashSet::new(),
finished: false,
}
}
}
impl<I> Iterator for GroupByPositionIterator<I>
where
I: Iterator<Item = (Position, String, Vec<usize>)>,
{
type Item = (Position, HashSet<String>, HashSet<usize>);
fn next(&mut self) -> Option<Self::Item> {
if self.finished {
return None;
}
for (pos, tag, qiids) in self.inner.by_ref() {
if self.current_pos.is_none() {
self.current_pos = Some(pos);
} else if self.current_pos.unwrap() != pos {
let out_pos = self.current_pos.unwrap();
let out_tags = std::mem::take(&mut self.tags);
let out_qiis = std::mem::take(&mut self.qiis);
self.current_pos = Some(pos);
self.tags.insert(tag);
for q in qiids {
self.qiis.insert(q);
}
return Some((out_pos, out_tags, out_qiis));
}
self.tags.insert(tag);
for q in qiids {
self.qiis.insert(q);
}
}
if let Some(p) = self.current_pos.take() {
self.finished = true;
let out_tags = std::mem::take(&mut self.tags);
let out_qiis = std::mem::take(&mut self.qiis);
return Some((p, out_tags, out_qiis));
}
None
}
}
let mut out: Vec<DcbSequencedEvent> = Vec::new();
for (pos, tags_present, qiis_present) in GroupByPositionIterator::new(merged) {
if let Some(ref c) = cancel {
if c.load(std::sync::atomic::Ordering::Relaxed) {
return Err(DcbError::CancelledByUser());
}
}
let matching_qiis: Vec<usize> = qiis_present
.iter()
.copied()
.filter(|&qii| qi_tags[qii].is_subset(&tags_present))
.collect();
if matching_qiis.is_empty() {
continue;
}
let rec = event_tree_lookup(mvcc, dirty, events_tree_root_id, pos)?;
let mut match_ok = false;
'matchcheck: for qii in matching_qiis.iter().copied() {
let item = &query.items[qii];
let type_ok = item.types.is_empty() || item.types.iter().any(|t| t == &rec.event_type);
if !type_ok {
continue;
}
let tags_ok = item.tags.iter().all(|t| rec.tags.iter().any(|et| et == t));
if tags_ok {
match_ok = true;
break 'matchcheck;
}
}
if !match_ok {
continue;
}
out.push(DcbSequencedEvent {
position: pos.0,
event: DcbEvent {
event_type: rec.event_type,
data: rec.data,
tags: rec.tags,
uuid: rec.uuid,
},
});
if let Some(lim) = limit
&& out.len() >= lim as usize
{
break;
}
}
Ok(out)
}
#[inline(always)]
pub fn tag_to_hash_v5uuid(tag: &str) -> TagHash {
let u = Uuid::new_v5(&Uuid::NAMESPACE_URL, tag.as_bytes());
u.into_bytes()
}
#[inline(always)]
pub fn tag_to_hash_crc64(tag: &str) -> TagHash {
const SALT: [u8; 4] = [0x9E, 0x37, 0x79, 0xB9];
let mut hasher1 = crc32fast::Hasher::new();
hasher1.update(tag.as_bytes());
let a = hasher1.finalize();
let mut hasher2 = crc32fast::Hasher::new();
hasher2.update(tag.as_bytes());
hasher2.update(&SALT);
let b = hasher2.finalize();
let value = ((a as u64) << 32) | (b as u64);
let mut out: TagHash = [0u8; crate::tags_tree_nodes::TAG_HASH_LEN];
out[..8].copy_from_slice(&value.to_le_bytes());
out
}
#[inline]
pub fn tag_to_hash(tag: &str) -> TagHash {
if get_tag_key_width() == 16 {
tag_to_hash_v5uuid(tag)
} else {
tag_to_hash_crc64(tag)
}
}
pub fn is_request_idempotent(
mvcc: &Arc<Mvcc>,
dirty: &HashMap<PageID, Page>,
events_tree_root_id: PageID,
tags_tree_root_id: PageID,
events: &Vec<DcbEvent>,
fail_if_events_match: DcbQuery,
start: Option<Position>,
cancel: Option<Arc<std::sync::atomic::AtomicBool>>,
) -> DcbResult<Option<u64>> {
let submitted_events_len = events.len();
let mut submitted_event_ids: Vec<Option<Uuid>> = vec![];
for submitted_event in events {
if submitted_event.uuid.is_some() {
submitted_event_ids.push(submitted_event.uuid);
}
}
if submitted_events_len == submitted_event_ids.len()
&& submitted_events_len as u64 <= u32::MAX as u64
{
let read_result = read_conditional(
mvcc,
dirty,
events_tree_root_id,
tags_tree_root_id,
fail_if_events_match,
start,
false,
Some(submitted_events_len as u32),
false,
cancel,
);
match read_result {
Ok(found_events) => {
let mut found_event_ids: Vec<Option<Uuid>> = vec![];
let found_events_len = found_events.len();
if found_events_len == submitted_events_len {
let last_found_event = &found_events[found_events_len - 1];
let last_found_event_position = last_found_event.position;
for found_event in found_events {
found_event_ids.push(found_event.event.uuid);
}
if found_event_ids == submitted_event_ids {
return Ok(Some(last_found_event_position));
}
}
}
Err(e) => {
return Err(e);
}
}
}
Ok(None)
}
pub fn is_integrity_error(e: &DcbError) -> bool {
matches!(e, DcbError::IntegrityError(_))
}
pub fn clone_dcb_error(src: &DcbError) -> DcbError {
match src {
DcbError::AuthenticationError(err) => DcbError::AuthenticationError(err.to_string()),
DcbError::InitializationError(err) => DcbError::InitializationError(err.to_string()),
DcbError::Io(err) => DcbError::Io(std::io::Error::other(err.to_string())),
DcbError::IntegrityError(s) => DcbError::IntegrityError(s.clone()),
DcbError::Corruption(s) => DcbError::Corruption(s.clone()),
DcbError::InvalidArgument(s) => DcbError::InvalidArgument(s.clone()),
DcbError::PageNotFound(id) => DcbError::PageNotFound(*id),
DcbError::DirtyPageNotFound(id) => DcbError::DirtyPageNotFound(*id),
DcbError::RootIDMismatch(old_id, new_id) => DcbError::RootIDMismatch(*old_id, *new_id),
DcbError::DatabaseCorrupted(s) => DcbError::DatabaseCorrupted(s.clone()),
DcbError::InternalError(s) => DcbError::InternalError(s.clone()),
DcbError::SerializationError(s) => DcbError::SerializationError(s.clone()),
DcbError::DeserializationError(s) => DcbError::DeserializationError(s.clone()),
DcbError::PageAlreadyFreed(id) => DcbError::PageAlreadyFreed(*id),
DcbError::PageAlreadyDirty(id) => DcbError::PageAlreadyDirty(*id),
DcbError::TransportError(err) => DcbError::TransportError(err.clone()),
DcbError::CancelledByUser() => DcbError::CancelledByUser(),
}
}
pub fn shadow_for_batch_abort(src: &DcbError) -> DcbError {
let msg = "batch aborted due to internal error".to_string();
match src {
DcbError::AuthenticationError(_) => DcbError::AuthenticationError(msg),
DcbError::InitializationError(_) => DcbError::InitializationError(msg),
DcbError::Io(_) => DcbError::Io(std::io::Error::other(msg)),
DcbError::IntegrityError(_) => DcbError::IntegrityError(msg),
DcbError::Corruption(_) => DcbError::Corruption(msg),
DcbError::InvalidArgument(_) => DcbError::InvalidArgument(msg),
DcbError::DatabaseCorrupted(_) => DcbError::DatabaseCorrupted(msg),
DcbError::InternalError(_) => DcbError::InternalError(msg),
DcbError::SerializationError(_) => DcbError::SerializationError(msg),
DcbError::DeserializationError(_) => DcbError::DeserializationError(msg),
DcbError::TransportError(_) => DcbError::TransportError(msg),
DcbError::PageNotFound(id) => DcbError::PageNotFound(*id),
DcbError::DirtyPageNotFound(id) => DcbError::DirtyPageNotFound(*id),
DcbError::RootIDMismatch(a, b) => DcbError::RootIDMismatch(*a, *b),
DcbError::PageAlreadyFreed(id) => DcbError::PageAlreadyFreed(*id),
DcbError::PageAlreadyDirty(id) => DcbError::PageAlreadyDirty(*id),
DcbError::CancelledByUser() => DcbError::CancelledByUser(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::page::Page;
use serial_test::serial;
use std::collections::HashMap;
use tempfile::tempdir;
use umadb_dcb::{
DcbAppendCondition, DcbEvent, DcbEventStoreSync, DcbQuery, DcbQueryItem, DcbError,
};
use uuid::Uuid;
#[test]
#[serial]
fn tracking_get_none_on_new_db() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("tracking-none.db");
let uma = UmaDb::new(db_path).unwrap();
let pos = uma.get_tracking_info("source-A").unwrap();
assert!(pos.is_none());
}
#[test]
#[serial]
fn tracking_source_length_too_long_errors() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("tracking-longkey.db");
let uma = UmaDb::new(db_path).unwrap();
let long_key = "a".repeat(256);
let ev = DcbEvent::new().event_type("T").data([1u8]);
let err = uma
.append(
vec![ev],
None,
Some(TrackingInfo {
source: long_key,
position: 1,
}),
)
.unwrap_err();
match err {
DcbError::InvalidArgument(msg) => assert!(msg.contains("too long")),
other => panic!("unexpected error: {:?}", other),
}
}
#[test]
#[serial]
fn tracking_leaf_split_creates_internal_root_and_lookups_work() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("tracking-split.db");
let mvcc = Mvcc::new(&db_path, 256, false).unwrap();
let uma = UmaDb::from_arc(Arc::new(mvcc));
let base_event = DcbEvent::new().event_type("T").data([0u8]);
for i in 0..50u32 {
let key = format!("k{:03}", i);
let ev = base_event.clone();
uma.append(
vec![ev],
None,
Some(TrackingInfo {
source: key.clone(),
position: (i + 1) as u64,
}),
)
.unwrap();
}
assert_eq!(Some(1), uma.get_tracking_info("k000").unwrap());
assert_eq!(Some(25), uma.get_tracking_info("k024").unwrap());
assert_eq!(Some(50), uma.get_tracking_info("k049").unwrap());
assert_eq!(None, uma.get_tracking_info("k999").unwrap());
}
#[test]
#[serial]
fn tracking_internal_node_splits_under_load() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("tracking-internal-split.db");
let mvcc = Mvcc::new(&db_path, 128, false).unwrap();
let uma = UmaDb::from_arc(Arc::new(mvcc));
let base_event = DcbEvent::new().event_type("T").data([0u8]);
let mut observed: HashMap<String, u64> = HashMap::new();
let mut detected_internal_split = false;
for i in 0..200u32 {
let key = format!("s{:03}xxxx", i); let ev = base_event.clone();
let pos = (i + 1) as u64;
uma.append(
vec![ev],
None,
Some(TrackingInfo {
source: key.clone(),
position: pos,
}),
)
.unwrap();
observed.insert(key, pos);
if i % 5 == 4 {
let reader = uma.mvcc.reader().unwrap();
let root_id = reader.tracking_tree_root_id;
if root_id != PageID(0) {
let root = uma.mvcc.read_page(root_id).unwrap();
if let Node::TrackingInternal(root_internal) = &root.node {
let first_child_id = root_internal.child_ids[0];
let first_child = uma.mvcc.read_page(first_child_id).unwrap();
if matches!(first_child.node, Node::TrackingInternal(_)) {
detected_internal_split = true;
break;
}
}
}
}
}
assert!(
detected_internal_split,
"Exceeded safety limit without causing tracking internal split"
);
for (k, expected_pos) in &observed {
let got = uma.get_tracking_info(&k).unwrap();
assert_eq!(
Some(*expected_pos),
got,
"tracking info mismatch for key {k}"
);
}
let mut updated: HashMap<String, u64> = HashMap::new();
for (k, prev_pos) in &observed {
let new_pos = *prev_pos + 1000;
uma.append(
vec![base_event.clone()],
None,
Some(TrackingInfo {
source: k.clone(),
position: new_pos,
}),
)
.unwrap();
updated.insert(k.clone(), new_pos);
}
for (k, expected_pos) in updated {
let got = uma.get_tracking_info(&k).unwrap();
assert_eq!(
Some(expected_pos),
got,
"after update: tracking info mismatch for key {k}"
);
}
}
#[test]
#[serial]
fn append_with_tracking_create_and_update_and_monotonic_enforced() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("tracking-append.db");
let uma = UmaDb::new(db_path).unwrap();
let ev = DcbEvent::new()
.event_type("T1")
.data(vec![1, 2, 3])
.tags(["x", "y"]);
let last = uma
.append(
vec![ev.clone()],
None,
Some(TrackingInfo {
source: "src1".into(),
position: 5,
}),
)
.unwrap();
assert_eq!(1, last);
assert_eq!(Some(5), uma.get_tracking_info("src1").unwrap());
let err = uma
.append(
vec![ev.clone()],
None,
Some(TrackingInfo {
source: "src1".into(),
position: 5,
}),
)
.err()
.expect("expected error");
match err {
DcbError::IntegrityError(msg) => {
assert!(msg.contains("non-increasing tracking position"))
}
other => panic!("unexpected error: {:?}", other),
}
let last2 = uma
.append(
vec![ev],
None,
Some(TrackingInfo {
source: "src1".into(),
position: 6,
}),
)
.unwrap();
assert_eq!(2, last2);
assert_eq!(Some(6), uma.get_tracking_info("src1").unwrap());
}
fn read_conditional(
mvcc: &Mvcc,
events_tree_root_id: PageID,
tags_tree_root_id: PageID,
query: DcbQuery,
start: Option<Position>,
backwards: bool,
limit: Option<u32>,
) -> DcbResult<Vec<DcbSequencedEvent>> {
super::read_conditional(
mvcc,
&HashMap::<PageID, Page>::new(),
events_tree_root_id,
tags_tree_root_id,
query,
start,
backwards,
limit,
false,
None,
)
}
static VERBOSE: bool = false;
fn standard_events() -> Vec<DcbEvent> {
let shared_tags = vec![
"alpha".to_string(),
"beta".to_string(),
"gamma".to_string(),
"delta".to_string(),
"epsilon".to_string(),
];
let mut input: Vec<DcbEvent> = Vec::new();
for i in 0..10u8 {
let t1 = shared_tags[(i % 5) as usize].clone();
let t2 = shared_tags[((i + 2) % 5) as usize].clone();
input.push(DcbEvent {
event_type: format!("Type{}", i),
data: vec![i, i + 1, i + 2],
tags: vec![t1, t2],
uuid: None,
});
}
input
}
fn setup_db_with_standard_events() -> (tempfile::TempDir, Mvcc, Vec<DcbEvent>) {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("mvcc-api-test.db");
let db = Mvcc::new(db_path.as_ref(), DEFAULT_PAGE_SIZE, VERBOSE).unwrap();
let input = standard_events();
let mut writer = db.writer().unwrap();
let last = unconditional_append(&db, &mut writer, input.clone()).unwrap();
db.commit(&mut writer).unwrap();
let header_page = db.get_latest_header_page().unwrap();
let header = header_page.as_header_node().unwrap();
let head = header.next_position.0.saturating_sub(1);
assert_eq!(last, head);
(temp_dir, db, input)
}
#[test]
#[serial]
fn empty_query_after_and_limit() {
let (_tmp, mut mvcc, input) = setup_db_with_standard_events();
let reader = mvcc.reader().unwrap();
let all = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(all.len(), input.len());
assert!(all.windows(2).all(|w| w[0].position < w[1].position));
let first = all[0].position;
let tail = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(first + 1)),
false,
None,
)
.unwrap();
assert_eq!(tail.len(), input.len() - 1);
let last = all.last().unwrap().position;
let none = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(last + 1)),
false,
None,
)
.unwrap();
assert!(none.is_empty());
let lim0 = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(1)),
false,
Some(0),
)
.unwrap();
assert!(lim0.is_empty());
let lim3 = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(1)),
false,
Some(3),
)
.unwrap();
assert_eq!(lim3.len(), 3);
let lim20 = read_conditional(
&mut mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(1)),
false,
Some(20),
)
.unwrap();
assert_eq!(lim20.len(), input.len());
}
#[test]
#[serial]
fn tags_only_single_tag_after_and_limit() {
let (_tmp, mut db, _input) = setup_db_with_standard_events();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string()],
}],
};
let reader = db.reader().unwrap();
let res = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(res.len(), 4);
assert!(
res.iter()
.all(|e| e.event.tags.iter().any(|t| t == "alpha"))
);
assert!(res.windows(2).all(|w| w[0].position < w[1].position));
let positions: Vec<u64> = res.iter().map(|e| e.position).collect();
let after_first = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(positions[0] + 1)),
false,
None,
)
.unwrap();
assert_eq!(after_first.len(), positions.len() - 1);
let after_last = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(*positions.last().unwrap() + 1)),
false,
None,
)
.unwrap();
assert!(after_last.is_empty());
let lim0 = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
Some(0),
)
.unwrap();
assert!(lim0.is_empty());
let lim1 = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
Some(1),
)
.unwrap();
assert_eq!(lim1.len(), 1);
let lim10 = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(1)),
false,
Some(10),
)
.unwrap();
assert_eq!(lim10.len(), 4);
}
#[test]
#[serial]
fn tags_only_multi_tag_and() {
let (_tmp, mut db, _input) = setup_db_with_standard_events();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string(), "gamma".to_string()],
}],
};
let reader = db.reader().unwrap();
let res = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(res.len(), 2);
assert!(
res.iter()
.all(|e| e.event.tags.iter().any(|t| t == "alpha"))
);
assert!(
res.iter()
.all(|e| e.event.tags.iter().any(|t| t == "gamma"))
);
}
#[test]
#[serial]
fn types_plus_tags_index_path() {
let (_tmp, mut db, _input) = setup_db_with_standard_events();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["Type0".to_string()],
tags: vec!["alpha".to_string()],
}],
};
let reader = db.reader().unwrap();
let res = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(res.len(), 1);
assert_eq!(res[0].event.event_type, "Type0");
assert!(res[0].event.tags.iter().any(|t| t == "alpha"));
}
#[test]
#[serial]
fn or_semantics_and_deduplication() {
let (_tmp, mut db, _input) = setup_db_with_standard_events();
let alpha_only = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string()],
}],
};
let reader = db.reader().unwrap();
let alpha_positions: Vec<u64> = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
alpha_only.clone(),
Some(Position(1)),
false,
None,
)
.unwrap()
.into_iter()
.map(|e| e.position)
.collect();
let query = DcbQuery {
items: vec![
DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string()],
},
DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string(), "gamma".to_string()],
},
],
};
let res = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
query,
Some(Position(1)),
false,
None,
)
.unwrap();
let res_positions: Vec<u64> = res.into_iter().map(|e| e.position).collect();
assert_eq!(res_positions, alpha_positions);
}
#[test]
#[serial]
fn fallback_types_only_after_and_limit() {
let temp_dir = tempdir().unwrap();
let db_path = temp_dir.path().join("mvcc-fallback-types-only.db");
let mut db = Mvcc::new(db_path.as_ref(), DEFAULT_PAGE_SIZE, VERBOSE).unwrap();
let events = vec![
DcbEvent {
event_type: "TypeA".to_string(),
data: vec![1],
tags: vec!["x".to_string()],
uuid: None,
},
DcbEvent {
event_type: "TypeB".to_string(),
data: vec![2],
tags: vec!["y".to_string()],
uuid: None,
},
DcbEvent {
event_type: "TypeA".to_string(),
data: vec![3],
tags: vec!["z".to_string()],
uuid: None,
},
];
let mut writer = db.writer().unwrap();
let last = unconditional_append(&db, &mut writer, events).unwrap();
db.commit(&mut writer).unwrap();
let header_page = db.get_latest_header_page().unwrap();
let header = header_page.as_header_node().unwrap();
let head = header.next_position.0.saturating_sub(1);
assert_eq!(last, head);
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["TypeA".to_string()],
tags: vec![],
}],
};
let reader = db.reader().unwrap();
let res = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(res.len(), 2);
assert!(res.iter().all(|e| e.event.event_type == "TypeA"));
let first_pos = res[0].position;
let res_after = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(first_pos + 1)),
false,
None,
)
.unwrap();
assert_eq!(res_after.len(), 1);
let res_lim1 = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(1)),
false,
Some(1),
)
.unwrap();
assert_eq!(res_lim1.len(), 1);
}
#[test]
#[serial]
fn fallback_empty_item_matches_all() {
let (_tmp, mut db, input) = setup_db_with_standard_events();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec![],
}],
};
let reader = db.reader().unwrap();
let all = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
assert_eq!(all.len(), input.len());
let first = all[1].position;
let tail = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(first)),
false,
None,
)
.unwrap();
assert_eq!(tail.len(), input.len() - 1);
let lim5 = read_conditional(
&mut db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(1)),
false,
Some(5),
)
.unwrap();
assert_eq!(lim5.len(), 5);
}
#[test]
#[serial]
fn test_event_store() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
assert_eq!(None, store.head().unwrap());
let events = vec![
DcbEvent {
event_type: "TypeA".to_string(),
data: vec![1],
tags: vec!["foo".to_string()],
uuid: None,
},
DcbEvent {
event_type: "TypeB".to_string(),
data: vec![2],
tags: vec!["bar".to_string(), "foo".to_string()],
uuid: None,
},
];
let last = store.append(events.clone(), None, None).unwrap();
assert!(last > 0);
assert_eq!(store.head().unwrap(), Some(last));
let mut resp = store.read(None, None, false, None, false).unwrap();
let (all, head) = resp.collect_with_head().unwrap();
assert_eq!(head, Some(last));
assert_eq!(all.len(), 2);
assert_eq!(all[0].event.event_type, "TypeA");
assert_eq!(all[1].event.event_type, "TypeB");
let mut resp_lim1 = store.read(None, None, false, Some(1), false).unwrap();
let (only_one, head_lim1) = resp_lim1.collect_with_head().unwrap();
assert_eq!(only_one.len(), 1);
assert_eq!(only_one[0].event.event_type, "TypeA");
assert_eq!(head_lim1, Some(only_one[0].position));
let query = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["foo".to_string()],
}],
};
let mut resp2 = store.read(Some(query), None, false, None, false).unwrap();
let out2 = resp2.next_batch().unwrap();
assert_eq!(out2.len(), 2);
assert!(out2.iter().all(|e| e.event.tags.iter().any(|t| t == "foo")));
let first_pos = all[0].position + 1;
let mut resp3 = store
.read(None, Some(first_pos), false, None, false)
.unwrap();
let out3 = resp3.next_batch().unwrap();
assert_eq!(out3.len(), 1);
assert_eq!(out3[0].event.event_type, "TypeB");
let cond_pass = DcbAppendCondition {
fail_if_events_match: DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["foo".to_string()],
}],
},
after: Some(last),
};
let ok_last = store
.append(
vec![DcbEvent {
event_type: "TypeC".to_string(),
data: vec![3],
tags: vec!["baz".to_string()],
uuid: None,
}],
Some(cond_pass),
None,
)
.expect("append with passing condition should succeed");
assert!(ok_last > last);
assert_eq!(store.head().unwrap(), Some(ok_last));
let cond_fail = DcbAppendCondition {
fail_if_events_match: DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["foo".to_string()],
}],
},
after: Some(0),
};
let before_head = store.head().unwrap();
let res = store.append(
vec![DcbEvent {
event_type: "TypeD".to_string(),
data: vec![4],
tags: vec!["qux".to_string()],
uuid: None,
}],
Some(cond_fail),
None,
);
match res {
Err(DcbError::IntegrityError(_)) => {}
other => panic!("Expected IntegrityError, got {:?}", other),
}
assert_eq!(store.head().unwrap(), before_head);
}
#[test]
fn test_append_batch_mixed_conditions() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
let e1 = DcbEvent {
event_type: "A".into(),
data: b"1".to_vec(),
tags: vec!["t1".into()],
uuid: None,
};
let e2 = DcbEvent {
event_type: "B".into(),
data: b"2".to_vec(),
tags: vec!["t2".into()],
uuid: None,
};
let e3 = DcbEvent {
event_type: "C".into(),
data: b"3".to_vec(),
tags: vec!["t3".into()],
uuid: None,
};
let items = vec![
(vec![e1.clone()], None, None),
(
vec![e2.clone()],
Some(DcbAppendCondition {
fail_if_events_match: DcbQuery::default(),
after: None,
}),
None,
),
(
vec![e3.clone()],
Some(DcbAppendCondition {
fail_if_events_match: DcbQuery::default(),
after: Some(10),
}),
None,
),
];
let results = store.append_batch(items).unwrap();
assert_eq!(results.len(), 3);
match &results[0] {
Ok(pos) => assert_eq!(*pos, 1),
Err(e) => panic!("unexpected error for first item: {:?}", e),
}
match &results[1] {
Ok(pos) => panic!("expected integrity error, got Ok({})", pos),
Err(e) => assert!(matches!(e, DcbError::IntegrityError(_))),
}
match &results[2] {
Ok(pos) => assert_eq!(*pos, 2),
Err(e) => panic!("unexpected error for third item: {:?}", e),
}
let (events, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].event.data, e1.data);
assert_eq!(events[1].event.data, e3.data);
assert_eq!(head, Some(2));
}
#[test]
fn test_append_batch_dirty_visibility_with_tags() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
let e1 = DcbEvent {
event_type: "T".into(),
data: b"one".to_vec(),
tags: vec!["x".into()],
uuid: None,
};
let e2 = DcbEvent {
event_type: "T".into(),
data: b"two".to_vec(),
tags: vec!["y".into()],
uuid: None,
};
let e3 = DcbEvent {
event_type: "T".into(),
data: b"three".to_vec(),
tags: vec!["z".into()],
uuid: None,
};
let query_tag_x = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["x".into()],
}],
};
let items = vec![
(vec![e1.clone()], None, None),
(
vec![e2.clone()],
Some(DcbAppendCondition {
fail_if_events_match: query_tag_x.clone(),
after: None,
}),
None,
),
(
vec![e3.clone()],
Some(DcbAppendCondition {
fail_if_events_match: query_tag_x.clone(),
after: Some(1),
}),
None,
),
];
let results = store.append_batch(items).unwrap();
assert_eq!(results.len(), 3);
match &results[0] {
Ok(pos) => assert_eq!(*pos, 1),
Err(e) => panic!("unexpected error for first item: {:?}", e),
}
match &results[1] {
Ok(pos) => panic!("expected integrity error, got Ok({})", pos),
Err(e) => assert!(matches!(e, DcbError::IntegrityError(_))),
}
match &results[2] {
Ok(pos) => assert_eq!(*pos, 2),
Err(e) => panic!("unexpected error for third item: {:?}", e),
}
let (events, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(events.len(), 2);
assert_eq!(events[0].event.data, e1.data);
assert_eq!(events[1].event.data, e3.data);
assert_eq!(head, Some(2));
let (tagx_events, _) = store
.read_with_head(Some(query_tag_x.clone()), None, false, None)
.unwrap();
assert_eq!(tagx_events.len(), 1);
assert_eq!(tagx_events[0].event.data, e1.data);
}
#[test]
fn test_append_batch_dirty_visibility_with_types_small_and_big_overflow() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
let small = DcbEvent {
event_type: "S".into(),
data: b"sm".to_vec(),
tags: vec!["tS".into()],
uuid: None,
};
let big_data_len = DEFAULT_PAGE_SIZE * 3; let big = DcbEvent {
event_type: "B".into(),
data: vec![0xAB; big_data_len],
tags: vec!["tB".into()],
uuid: None,
};
let filler1 = DcbEvent {
event_type: "X".into(),
data: b"x".to_vec(),
tags: vec![],
uuid: None,
};
let filler2 = DcbEvent {
event_type: "Y".into(),
data: b"y".to_vec(),
tags: vec![],
uuid: None,
};
let final_ok = DcbEvent {
event_type: "C".into(),
data: b"c".to_vec(),
tags: vec![],
uuid: None,
};
let q_type_s = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["S".into()],
tags: vec![],
}],
};
let q_type_b = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["B".into()],
tags: vec![],
}],
};
let items = vec![
(vec![small.clone()], None, None),
(
vec![filler1.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_type_s.clone(),
after: None,
}),
None,
),
(vec![big.clone()], None, None),
(
vec![filler2.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_type_b.clone(),
after: None,
}),
None,
),
(
vec![final_ok.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_type_b.clone(),
after: Some(2),
}),
None,
),
];
let results = store.append_batch(items).unwrap();
assert_eq!(results.len(), 5);
match &results[0] {
Ok(pos) => assert_eq!(*pos, 1),
other => panic!("unexpected for item0: {:?}", other),
}
match &results[1] {
Err(DcbError::IntegrityError(_)) => {}
other => {
panic!("expected IntegrityError for item1, got {:?}", other)
}
}
match &results[2] {
Ok(pos) => assert_eq!(*pos, 2),
other => panic!("unexpected for item2: {:?}", other),
}
match &results[3] {
Err(DcbError::IntegrityError(_)) => {}
other => {
panic!("expected IntegrityError for item3, got {:?}", other)
}
}
match &results[4] {
Ok(pos) => assert_eq!(*pos, 3),
other => panic!("unexpected for item4: {:?}", other),
}
let (events, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(events[0].event.event_type, small.event_type);
assert_eq!(events[1].event.event_type, big.event_type);
assert_eq!(events[2].event.event_type, final_ok.event_type);
assert_eq!(head, Some(3));
let (small_by_type, _) = store
.read_with_head(Some(q_type_s.clone()), None, false, None)
.unwrap();
assert_eq!(small_by_type.len(), 1);
assert_eq!(small_by_type[0].event.event_type, "S");
let (big_by_type, _) = store
.read_with_head(Some(q_type_b.clone()), None, false, None)
.unwrap();
assert_eq!(big_by_type.len(), 1);
assert_eq!(big_by_type[0].event.event_type, "B");
assert_eq!(big_by_type[0].event.data.len(), big_data_len);
assert!(big_by_type[0].event.data.iter().all(|&b| b == 0xAB));
}
#[test]
fn test_append_batch_dirty_visibility_with_tags_and_types_small_and_big_overflow() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
let small = DcbEvent {
event_type: "S".into(),
data: b"sm".to_vec(),
tags: vec!["x".into()],
uuid: None,
};
let big_data_len = DEFAULT_PAGE_SIZE * 3; let big = DcbEvent {
event_type: "B".into(),
data: vec![0xCD; big_data_len],
tags: vec!["y".into()],
uuid: None,
};
let filler1 = DcbEvent {
event_type: "X".into(),
data: b"x".to_vec(),
tags: vec![],
uuid: None,
};
let filler2 = DcbEvent {
event_type: "Y".into(),
data: b"y".to_vec(),
tags: vec![],
uuid: None,
};
let final_ok = DcbEvent {
event_type: "C".into(),
data: b"c".to_vec(),
tags: vec![],
uuid: None,
};
let q_s_and_x = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["S".into()],
tags: vec!["x".into()],
}],
};
let q_b_and_y = DcbQuery {
items: vec![DcbQueryItem {
types: vec!["B".into()],
tags: vec!["y".into()],
}],
};
let items = vec![
(vec![small.clone()], None, None),
(
vec![filler1.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_s_and_x.clone(),
after: None,
}),
None,
),
(vec![big.clone()], None, None),
(
vec![filler2.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_b_and_y.clone(),
after: None,
}),
None,
),
(
vec![final_ok.clone()],
Some(DcbAppendCondition {
fail_if_events_match: q_b_and_y.clone(),
after: Some(2),
}),
None,
),
];
let results = store.append_batch(items).unwrap();
assert_eq!(results.len(), 5);
match &results[0] {
Ok(pos) => assert_eq!(*pos, 1),
other => panic!("unexpected for item0: {:?}", other),
}
match &results[1] {
Err(DcbError::IntegrityError(_)) => {}
other => {
panic!("expected IntegrityError for item1, got {:?}", other)
}
}
match &results[2] {
Ok(pos) => assert_eq!(*pos, 2),
other => panic!("unexpected for item2: {:?}", other),
}
match &results[3] {
Err(DcbError::IntegrityError(_)) => {}
other => {
panic!("expected IntegrityError for item3, got {:?}", other)
}
}
match &results[4] {
Ok(pos) => assert_eq!(*pos, 3),
other => panic!("unexpected for item4: {:?}", other),
}
let (events, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(events.len(), 3);
assert_eq!(events[0].event.event_type, small.event_type);
assert_eq!(events[1].event.event_type, big.event_type);
assert_eq!(events[2].event.event_type, final_ok.event_type);
assert_eq!(head, Some(3));
let (small_combined, _) = store
.read_with_head(Some(q_s_and_x.clone()), None, false, None)
.unwrap();
assert_eq!(small_combined.len(), 1);
assert_eq!(small_combined[0].event.event_type, "S");
assert!(small_combined[0].event.tags.iter().any(|t| t == "x"));
let (big_combined, _) = store
.read_with_head(Some(q_b_and_y.clone()), None, false, None)
.unwrap();
assert_eq!(big_combined.len(), 1);
assert_eq!(big_combined[0].event.event_type, "B");
assert!(big_combined[0].event.tags.iter().any(|t| t == "y"));
assert_eq!(big_combined[0].event.data.len(), big_data_len);
assert!(big_combined[0].event.data.iter().all(|&b| b == 0xCD));
}
#[test]
fn test_append_event_with_uuid_is_maintained_and_activated_append_idempotency() {
let temp_dir = tempdir().unwrap();
let store = UmaDb::new(temp_dir.path()).unwrap();
let condition1 = Some(DcbAppendCondition {
fail_if_events_match: DcbQuery { items: vec![] },
after: None,
});
let event1 = DcbEvent {
event_type: "type1".to_string(),
data: b"data1".to_vec(),
tags: vec!["tag1".to_string()],
uuid: Some(Uuid::new_v4()),
};
let mut commit_position1 = store
.append(vec![event1.clone()], condition1.clone(), None)
.unwrap();
assert_eq!(1, commit_position1);
let (result, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(1, result.len());
assert_eq!(Some(1), head);
assert_eq!(event1.uuid, result[0].event.uuid);
commit_position1 = store
.append(vec![event1.clone()], condition1.clone(), None)
.unwrap();
assert_eq!(1, commit_position1);
let (result, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(1, result.len());
assert_eq!(Some(1), head);
assert_eq!(event1.uuid, result[0].event.uuid);
let event2 = DcbEvent {
event_type: "type2".to_string(),
data: b"data2".to_vec(),
tags: vec!["tag2".to_string()],
uuid: Some(Uuid::new_v4()),
};
let mut commit_position2 = store.append(vec![event2.clone()], None, None).unwrap();
assert_eq!(2, commit_position2);
let (result, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(2, result.len());
assert_eq!(Some(2), head);
assert_eq!(event1.uuid, result[0].event.uuid);
assert_eq!(event2.uuid, result[1].event.uuid);
commit_position1 = store
.append(vec![event1.clone()], condition1.clone(), None)
.unwrap();
assert_eq!(1, commit_position1);
commit_position2 = store
.append(
vec![event1.clone(), event2.clone()],
condition1.clone(),
None,
)
.unwrap();
assert_eq!(2, commit_position2);
let (result, head) = store.read_with_head(None, None, false, None).unwrap();
assert_eq!(2, result.len());
assert_eq!(Some(2), head);
assert_eq!(event1.uuid, result[0].event.uuid);
assert_eq!(event2.uuid, result[1].event.uuid);
let result = store.append(vec![event2.clone()], condition1.clone(), None);
assert!(matches!(result, Err(DcbError::IntegrityError(_))));
let result = store.append(
vec![event2.clone(), event1.clone()],
condition1.clone(),
None,
);
assert!(matches!(result, Err(DcbError::IntegrityError(_))));
}
#[test]
#[serial]
fn empty_query_backwards_from_and_limit() {
let (_tmp, mvcc, _input) = setup_db_with_standard_events();
let reader = mvcc.reader().unwrap();
let fwd = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(1)),
false,
None,
)
.unwrap();
let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
assert!(!fwd_pos.is_empty());
assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
let back_all = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
None,
true,
None,
)
.unwrap();
let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
let mut fwd_rev = fwd_pos.clone();
fwd_rev.reverse();
assert_eq!(fwd_rev, back_all_pos);
let last = *fwd_pos.last().unwrap();
let back_from_last = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(last)),
true,
None,
)
.unwrap();
let back_from_last_pos: Vec<u64> = back_from_last.iter().map(|e| e.position).collect();
assert_eq!(back_from_last_pos, fwd_rev);
let back_from_before_last = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
Some(Position(last - 1)),
true,
None,
)
.unwrap();
let back_from_before_last_pos: Vec<u64> =
back_from_before_last.iter().map(|e| e.position).collect();
assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
let back_lim3 = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
DcbQuery { items: vec![] },
None,
true,
Some(3),
)
.unwrap();
let back_lim3_pos: Vec<u64> = back_lim3.iter().map(|e| e.position).collect();
assert_eq!(back_lim3_pos, fwd_rev[..3.min(fwd_rev.len())].to_vec());
}
#[test]
#[serial]
fn tags_only_single_tag_backwards() {
let (_tmp, mvcc, _input) = setup_db_with_standard_events();
let reader = mvcc.reader().unwrap();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string()],
}],
};
let fwd = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
assert!(!fwd_pos.is_empty());
assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
let back_all = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
None,
true,
None,
)
.unwrap();
let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
let mut fwd_rev = fwd_pos.clone();
fwd_rev.reverse();
assert_eq!(back_all_pos, fwd_rev);
let last = *fwd_pos.last().unwrap();
let back_from_before_last = read_conditional(
&mvcc,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
Some(Position(last - 1)),
true,
None,
)
.unwrap();
let back_from_before_last_pos: Vec<u64> =
back_from_before_last.iter().map(|e| e.position).collect();
assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
}
#[test]
#[serial]
fn tags_only_multi_tag_and_backwards() {
let (_tmp, db, _input) = setup_db_with_standard_events();
let reader = db.reader().unwrap();
let qi = DcbQuery {
items: vec![DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string(), "gamma".to_string()],
}],
};
let fwd = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
assert!(!fwd_pos.is_empty());
assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
assert!(
fwd.iter()
.all(|e| e.event.tags.iter().any(|t| t == "alpha"))
);
assert!(
fwd.iter()
.all(|e| e.event.tags.iter().any(|t| t == "gamma"))
);
let back_all = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
None,
true,
None,
)
.unwrap();
let mut fwd_rev = fwd_pos.clone();
fwd_rev.reverse();
let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
assert_eq!(back_all_pos, fwd_rev);
let last = *fwd_pos.last().unwrap();
let back_from_last = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(last)),
true,
None,
)
.unwrap();
let back_from_last_pos: Vec<u64> = back_from_last.iter().map(|e| e.position).collect();
assert_eq!(back_from_last_pos, fwd_rev);
let back_from_before_last = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(last - 1)),
true,
None,
)
.unwrap();
let back_from_before_last_pos: Vec<u64> =
back_from_before_last.iter().map(|e| e.position).collect();
assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
let back_lim2 = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
None,
true,
Some(2),
)
.unwrap();
let back_lim2_pos: Vec<u64> = back_lim2.iter().map(|e| e.position).collect();
assert_eq!(back_lim2_pos, fwd_rev[..2.min(fwd_rev.len())].to_vec());
}
#[test]
#[serial]
fn tags_multi_item_two_tags_each_backwards() {
let (_tmp, db, _input) = setup_db_with_standard_events();
let reader = db.reader().unwrap();
let qi = DcbQuery {
items: vec![
DcbQueryItem {
types: vec![],
tags: vec!["alpha".to_string(), "gamma".to_string()],
},
DcbQueryItem {
types: vec![],
tags: vec!["beta".to_string(), "delta".to_string()],
},
],
};
let fwd = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
Some(Position(1)),
false,
None,
)
.unwrap();
let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
assert!(!fwd_pos.is_empty());
assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
assert!(fwd.iter().all(|e| {
let tags = &e.event.tags;
let has = |a: &str| tags.iter().any(|t| t == a);
(has("alpha") && has("gamma")) || (has("beta") && has("delta"))
}));
let back_all = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi.clone(),
None,
true,
None,
)
.unwrap();
let mut fwd_rev = fwd_pos.clone();
fwd_rev.reverse();
let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
assert_eq!(back_all_pos, fwd_rev);
let back_lim1 = read_conditional(
&db,
reader.events_tree_root_id,
reader.tags_tree_root_id,
qi,
None,
true,
Some(1),
)
.unwrap();
assert_eq!(back_lim1.len(), 1);
assert_eq!(back_lim1[0].position, *fwd_rev.first().unwrap());
}
}