1use crate::common::PageID;
2use crate::common::Position;
3use crate::events_tree_nodes::{
4 EventInternalNode, EventLeafNode, EventOverflowNode, EventRecord, EventValue,
5 deserialize_metadata, metadata_serialized_size, serialize_metadata_into,
6 validate_event_record_for_append,
7};
8use crate::mvcc::{Mvcc, Writer};
9use crate::node::Node;
10use crate::page::{PAGE_HEADER_SIZE, Page};
11use std::collections::HashMap;
12use std::sync::Arc;
13use umadb_dcb::{DcbError, DcbResult};
14
15fn write_overflow_chain(mvcc: &Mvcc, writer: &mut Writer, data: &[u8]) -> DcbResult<PageID> {
17 let payload_cap = mvcc.page_size.saturating_sub(PAGE_HEADER_SIZE + 8);
19 if payload_cap == 0 {
20 return Err(DcbError::DatabaseCorrupted(
21 "Page size too small to store overflow data".to_string(),
22 ));
23 }
24 let chunk_count = data.len().div_ceil(payload_cap);
26 let mut chunks: Vec<&[u8]> = Vec::with_capacity(chunk_count);
27 let mut i = 0;
28 while i < data.len() {
29 let end = (i + payload_cap).min(data.len());
30 chunks.push(&data[i..end]);
31 i = end;
32 }
33 if chunks.is_empty() {
34 let page_id = writer.alloc_page_id();
36 let node = EventOverflowNode {
37 next: PageID(0),
38 data: Vec::new(),
39 };
40 let page = Page::new(page_id, Node::EventOverflow(node));
41 writer.insert_dirty(page)?;
42 return Ok(page_id);
43 }
44 let mut next_id = PageID(0);
45 for chunk in chunks.iter().rev() {
46 let page_id = writer.alloc_page_id();
47 let node = EventOverflowNode {
48 next: next_id,
49 data: (*chunk).to_vec(),
50 };
51 let page = Page::new(page_id, Node::EventOverflow(node));
52 writer.insert_dirty(page)?;
53 next_id = page_id;
54 }
55 Ok(next_id)
56}
57
58fn read_overflow_chain(
59 mvcc: &Mvcc,
60 dirty: &HashMap<PageID, Page>,
61 mut page_id: PageID,
62) -> DcbResult<Vec<u8>> {
63 let mut out: Vec<u8> = Vec::new();
64 while page_id.0 != 0 {
65 page_id = with_page(mvcc, dirty, page_id, |page| match &page.node {
66 Node::EventOverflow(node) => {
67 out.extend_from_slice(&node.data);
68 Ok(node.next)
69 }
70 _ => Err(DcbError::DatabaseCorrupted(
71 "Expected EventOverflow node".to_string(),
72 )),
73 })?;
74 }
75 Ok(out)
76}
77
78fn with_page<T, F>(
79 mvcc: &Mvcc,
80 dirty: &HashMap<PageID, Page>,
81 page_id: PageID,
82 f: F,
83) -> DcbResult<T>
84where
85 F: FnOnce(&Page) -> DcbResult<T>,
86{
87 if let Some(page) = dirty.get(&page_id) {
88 f(page)
89 } else {
90 let page = mvcc.read_page(page_id)?;
91 f(&page)
92 }
93}
94
95fn data_and_metadata_to_overflow_chain(
100 mvcc: &Mvcc,
101 writer: &mut Writer,
102 record_data: &Vec<u8>,
103 metadata: &Vec<(String, String)>,
104) -> DcbResult<(PageID, usize)> {
105 let metadata_len = if metadata.is_empty() {
106 0
107 } else {
108 metadata_serialized_size(metadata)
109 };
110
111 let mut combined = vec![0u8; metadata_len + record_data.len()];
113
114 if metadata_len > 0 {
115 serialize_metadata_into(metadata, &mut combined[..metadata_len])?;
117 }
118
119 combined[metadata_len..].copy_from_slice(record_data);
121
122 let root_id = write_overflow_chain(mvcc, writer, &combined)?;
123 Ok((root_id, metadata_len))
124}
125
126fn materialize_event_value(
127 mvcc: &Mvcc,
128 dirty: &HashMap<PageID, Page>,
129 value: &EventValue,
130) -> DcbResult<EventRecord> {
131 match value {
132 EventValue::Inline(rec) => Ok(rec.clone()),
133 EventValue::Overflow {
134 event_type,
135 data_len,
136 tags,
137 root_id,
138 uuid,
139 metadata_len,
140 } => {
141 let all_data = read_overflow_chain(mvcc, dirty, *root_id)?;
142 let all_data_len = all_data.len();
143 if (all_data_len as u64) != *data_len + *metadata_len {
144 return Err(DcbError::DatabaseCorrupted(format!(
145 "Overflow data length mismatch, data: {data_len}, metadata: {metadata_len}"
146 )));
147 }
148 let split = *metadata_len as usize;
150 let metadata = if *metadata_len > 0 {
151 deserialize_metadata(&all_data[..split])?
152 } else {
153 Vec::new()
154 };
155 let data = all_data[split..].to_vec();
156 Ok(EventRecord {
157 event_type: event_type.clone(),
158 data,
159 tags: tags.clone(),
160 uuid: *uuid,
161 metadata,
162 })
163 }
164 }
165}
166
167pub fn event_tree_append(
169 mvcc: &Mvcc,
170 writer: &mut Writer,
171 event_record: EventRecord,
172 position: Position,
173) -> DcbResult<()> {
174 let event_values_and_size_diffs =
175 validate_event_record_for_append(mvcc.page_size, event_record)?;
176 event_tree_append_event_value(mvcc, writer, event_values_and_size_diffs, position)
177}
178
179pub fn event_tree_append_event_value(
185 mvcc: &Mvcc,
186 writer: &mut Writer,
187 event_values_and_size_diffs: ((EventValue, usize), Option<(EventValue, usize)>),
188 position: Position,
189) -> DcbResult<()> {
190 let verbose = mvcc.verbose;
191 let mut current_page_id: PageID = writer.events_tree_root_id;
197
198 let mut stack: Vec<PageID> = Vec::new();
200 loop {
201 let current_page_ref = writer.get_page_ref(mvcc, current_page_id)?;
202 if matches!(current_page_ref.node, Node::EventLeaf(_)) {
203 break;
204 }
205 if let Node::EventInternal(internal_node) = ¤t_page_ref.node {
206 if verbose {
207 println!("{:?} is internal node", current_page_ref.page_id);
208 }
209 stack.push(current_page_id);
210 current_page_id = *internal_node.child_ids.last().ok_or_else(|| {
211 DcbError::DatabaseCorrupted("Internal node has no children".to_string())
212 })?;
213 } else {
214 return Err(DcbError::DatabaseCorrupted(
215 "Expected EventInternal node".to_string(),
216 ));
217 }
218 }
219 if verbose {
220 println!("{current_page_id:?} is leaf node");
221 }
222
223 let dirty_page_id = { writer.get_dirty_page_id(current_page_id)? };
225 let replacement_info: Option<(PageID, PageID)> = {
226 if dirty_page_id != current_page_id {
227 Some((current_page_id, dirty_page_id))
228 } else {
229 None
230 }
231 };
232
233 let serialized_size = writer
234 .get_page_ref(mvcc, dirty_page_id)?
235 .calc_serialized_size();
236
237 let (event_value, size_diff) = match event_values_and_size_diffs {
256 ((EventValue::Inline(record), _), Some((mut overflow_value, overflow_size_diff))) => {
257 let (overflow_root_id, overflow_metadata_len) =
258 data_and_metadata_to_overflow_chain(mvcc, writer, &record.data, &record.metadata)?;
259 if let EventValue::Overflow {
260 root_id,
261 metadata_len,
262 ..
263 } = &mut overflow_value
264 {
265 *root_id = overflow_root_id;
266 *metadata_len = overflow_metadata_len as u64;
267 } else {
268 return Err(DcbError::InternalError(
269 "Shouldn't get here when setting root_id on event overflow value".to_string(),
270 ));
271 }
272
273 (overflow_value, overflow_size_diff)
274 }
275 ((inline_value, inline_size_diff), None) => (inline_value, inline_size_diff),
276 _ => {
277 return Err(DcbError::InternalError(
278 "Shouldn't get here when matching event_values_and_size_diffs".to_string(),
279 ));
280 }
281 };
282
283 let mut popped: Option<EventValue> = None;
285
286 if serialized_size + size_diff <= mvcc.page_size {
287 let dirty_leaf_page = writer.get_mut_dirty(dirty_page_id)?;
289 match &mut dirty_leaf_page.node {
290 Node::EventLeaf(node) => {
291 if verbose {
292 println!(
293 "Pushing {:?} at {:?} onto {:?}",
294 event_value, position, dirty_page_id,
295 );
296 }
297 node.keys.push(position);
298 node.values.push(event_value);
299 }
308 _ => {
309 return Err(DcbError::DatabaseCorrupted(
310 "Expected EventLeaf node".to_string(),
311 ));
312 }
313 }
314 } else {
315 popped = Some(event_value);
317 }
318
319 let mut split_info: Option<(Position, PageID)> = None;
321
322 if let Some(event_value) = popped {
323 let new_leaf_page_id = writer.alloc_page_id();
325 let new_leaf_node = EventLeafNode {
326 keys: vec![position],
327 values: vec![event_value],
328 };
329 let new_leaf_page = Page::new(new_leaf_page_id, Node::EventLeaf(new_leaf_node.clone()));
330 if verbose {
339 println!(
340 "Created new leaf {:?}: {:?}",
341 new_leaf_page_id, new_leaf_page.node
342 );
343 }
344 writer.insert_dirty(new_leaf_page)?;
345 if verbose {
346 println!("Promoting {position:?} and {new_leaf_page_id:?}");
347 }
348 split_info = Some((position, new_leaf_page_id));
349 }
350
351 let mut current_replacement_info = replacement_info;
353 while let Some(parent_page_id) = stack.pop() {
354 let dirty_page_id = writer.get_dirty_page_id(parent_page_id)?;
356 let parent_replacement_info: Option<(PageID, PageID)> = {
357 if dirty_page_id != parent_page_id {
358 Some((parent_page_id, dirty_page_id))
359 } else {
360 None
361 }
362 };
363 let dirty_internal_page = writer.get_mut_dirty(dirty_page_id)?;
365
366 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
367 if let Some((old_id, new_id)) = current_replacement_info {
368 dirty_internal_node.replace_last_child_id(old_id, new_id)?;
369 if verbose {
370 println!(
371 "Replaced {old_id:?} with {new_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
372 );
373 }
374 } else if verbose {
375 println!("Nothing to replace in {dirty_page_id:?}")
376 }
377 } else {
378 return Err(DcbError::DatabaseCorrupted(
379 "Expected EventInternal node".to_string(),
380 ));
381 }
382
383 if let Some((promoted_key, promoted_page_id)) = split_info {
384 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
385 dirty_internal_node
387 .append_promoted_key_and_page_id(promoted_key, promoted_page_id)?;
388
389 if verbose {
390 println!(
391 "Appended promoted key {promoted_key:?} and child {promoted_page_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
392 );
393 }
394 } else {
395 return Err(DcbError::DatabaseCorrupted(
396 "Expected EventInternal node".to_string(),
397 ));
398 }
399 }
400
401 if dirty_internal_page.calc_serialized_size() > mvcc.page_size {
404 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
405 if verbose {
406 println!("Splitting internal {dirty_page_id:?}...");
407 }
408 if dirty_internal_node.keys.len() < 3 || dirty_internal_node.child_ids.len() < 4 {
411 return Err(DcbError::DatabaseCorrupted(
412 "Cannot split internal node with too few keys/children".to_string(),
413 ));
414 }
415
416 let (promoted_key, new_keys, new_child_ids) = dirty_internal_node.split_off()?;
418
419 assert_eq!(
421 dirty_internal_node.keys.len() + 1,
422 dirty_internal_node.child_ids.len()
423 );
424
425 let new_internal_node = EventInternalNode {
426 keys: new_keys,
427 child_ids: new_child_ids,
428 };
429
430 assert_eq!(
432 new_internal_node.keys.len() + 1,
433 new_internal_node.child_ids.len()
434 );
435
436 let new_internal_page_id = writer.alloc_page_id();
438 let new_internal_page =
439 Page::new(new_internal_page_id, Node::EventInternal(new_internal_node));
440 if verbose {
441 println!(
442 "Created internal {:?}: {:?}",
443 new_internal_page_id, new_internal_page.node
444 );
445 }
446 writer.insert_dirty(new_internal_page)?;
447
448 split_info = Some((promoted_key, new_internal_page_id));
449 } else {
450 return Err(DcbError::DatabaseCorrupted(
451 "Expected EventInternal node".to_string(),
452 ));
453 }
454 } else {
455 split_info = None;
456 }
457 current_replacement_info = parent_replacement_info;
458 }
459
460 if let Some((old_id, new_id)) = current_replacement_info {
461 if writer.events_tree_root_id == old_id {
462 writer.events_tree_root_id = new_id;
463 if verbose {
464 println!("Replaced root {old_id:?} with {new_id:?}");
465 }
466 } else {
467 return Err(DcbError::RootIDMismatch(old_id.0, new_id.0));
468 }
469 }
470
471 if let Some((promoted_key, promoted_page_id)) = split_info {
472 let new_internal_node = EventInternalNode {
474 keys: vec![promoted_key],
475 child_ids: vec![writer.events_tree_root_id, promoted_page_id],
476 };
477
478 let new_root_page_id = writer.alloc_page_id();
479 let new_root_page = Page::new(new_root_page_id, Node::EventInternal(new_internal_node));
480 if verbose {
481 println!(
482 "Created new internal root {:?}: {:?}",
483 new_root_page_id, new_root_page.node
484 );
485 }
486 writer.insert_dirty(new_root_page)?;
487
488 writer.events_tree_root_id = new_root_page_id;
489 }
490
491 Ok(())
492}
493
494pub fn event_tree_lookup(
495 mvcc: &Mvcc,
496 dirty: &HashMap<PageID, Page>,
497 events_tree_root_id: PageID,
498 position: Position,
499) -> DcbResult<EventRecord> {
500 enum LookupStep {
501 Next(PageID),
502 Found(EventRecord),
503 }
504
505 let mut current_page_id: PageID = events_tree_root_id;
506 loop {
507 let step = with_page(mvcc, dirty, current_page_id, |page| match &page.node {
508 Node::EventInternal(internal) => {
509 let idx = match internal.keys.binary_search(&position) {
511 Ok(i) => i + 1,
512 Err(i) => i,
513 };
514 if idx >= internal.child_ids.len() {
515 return Err(DcbError::DatabaseCorrupted(
516 "Child index out of bounds in event tree".to_string(),
517 ));
518 }
519 Ok(LookupStep::Next(internal.child_ids[idx]))
520 }
521 Node::EventLeaf(leaf) => match leaf.keys.binary_search(&position) {
522 Ok(i) => {
523 let rec = materialize_event_value(mvcc, dirty, &leaf.values[i])?;
524 Ok(LookupStep::Found(rec))
525 }
526 Err(_) => Err(DcbError::DatabaseCorrupted(format!(
527 "Event at position {position:?} not found",
528 ))),
529 },
530 _ => Err(DcbError::DatabaseCorrupted(format!(
531 "Expected EventInternal or EventLeaf node in event tree, got {}",
532 page.node.type_name()
533 ))),
534 })?;
535
536 match step {
537 LookupStep::Next(next_page_id) => current_page_id = next_page_id,
538 LookupStep::Found(rec) => return Ok(rec),
539 }
540 }
541}
542
543pub struct EventIterator<'a> {
544 pub mvcc: &'a Mvcc,
545 pub dirty: &'a HashMap<PageID, Page>,
546 pub stack: Vec<(PageID, Option<usize>)>,
547 pub page_cache: HashMap<PageID, Arc<Page>>,
548 pub start: Option<Position>, pub backwards: bool,
550}
551
552impl<'a> EventIterator<'a> {
553 pub fn new(
554 mvcc: &'a Mvcc,
555 dirty: &'a HashMap<PageID, Page>,
556 events_tree_root_id: PageID,
557 start: Option<Position>,
558 backwards: bool,
559 ) -> Self {
560 let next_position = (events_tree_root_id, None);
561 Self {
562 mvcc,
563 dirty,
564 stack: vec![next_position],
565 page_cache: HashMap::new(),
566 start,
567 backwards,
568 }
569 }
570
571 pub fn next_batch(
572 &mut self,
573 batch_size: u32,
574 cancel: Option<&std::sync::Arc<std::sync::atomic::AtomicBool>>,
575 ) -> DcbResult<Vec<(Position, EventRecord)>> {
576 let mut result: Vec<(Position, EventRecord)> = Vec::with_capacity(batch_size as usize);
577 if batch_size == 0 {
578 return Ok(result);
579 }
580 if self.backwards && self.start == Some(Position(0)) {
581 return Ok(result);
582 }
583 while result.len() < batch_size as usize {
584 if let Some(c) = cancel {
585 if c.load(std::sync::atomic::Ordering::Relaxed) {
586 return Err(DcbError::CancelledByUser());
587 }
588 }
589 let Some((page_id, mut stacked_idx)) = self.stack.pop() else {
590 break; };
592
593 let mut remove_page = false;
595 let mut push_revisit: Option<(PageID, Option<usize>)> = None;
596 let mut push_child: Option<(PageID, Option<usize>)> = None; let mut emit_event: Option<(Position, EventRecord)> = None;
598
599 {
600 let page_ref: &Page = if let Some(p) = self.dirty.get(&page_id) {
602 p
603 } else if let Some(p) = self.page_cache.get(&page_id) {
604 p.as_ref()
605 } else {
606 let page_arc = self.mvcc.read_page(page_id)?;
607 self.page_cache.insert(page_id, page_arc);
608 self.page_cache
609 .get(&page_id)
610 .expect("page should be in cache")
611 .as_ref()
612 };
613
614 match &page_ref.node {
615 Node::EventInternal(internal) => {
616 if stacked_idx.is_none() && !internal.keys.is_empty() {
618 stacked_idx = match &self.start {
624 Some(from) => match internal.keys.binary_search(from) {
625 Ok(i) => Some(i + 1),
626 Err(i) => Some(i),
627 },
628 None => {
629 if !self.backwards {
630 Some(0)
631 } else {
632 Some(internal.child_ids.len() - 1)
633 }
634 }
635 };
636 }
637
638 if let Some(child_ids_idx) = stacked_idx {
639 push_child = Some((internal.child_ids[child_ids_idx], None));
643 if !self.backwards {
645 if child_ids_idx + 1 < internal.child_ids.len() {
646 push_revisit = Some((page_id, Some(child_ids_idx + 1)));
649 } else {
650 remove_page = true;
652 }
654 } else if child_ids_idx > 0 {
655 push_revisit = Some((page_id, Some(child_ids_idx - 1)));
658 } else {
659 remove_page = true;
661 }
663 } else {
664 remove_page = true
667 };
668 }
669 Node::EventLeaf(leaf) => {
670 if stacked_idx.is_none() {
672 let values_len = leaf.values.len();
675
676 stacked_idx = if values_len > 0 {
677 match &self.start {
678 Some(from) => match leaf.keys.binary_search(from) {
679 Ok(i) => Some(i),
680 Err(i) => {
681 if !self.backwards {
682 Some(i)
683 } else {
684 Some(i - 1)
685 }
686 }
687 },
688 None => {
689 if !self.backwards {
690 Some(0)
691 } else {
692 Some(values_len - 1)
693 }
694 }
695 }
696 } else {
697 None
698 }
699 }
700
701 if let Some(values_idx) = stacked_idx {
702 if values_idx < leaf.values.len() {
704 let event_position = leaf.keys[values_idx];
705 let event_record = materialize_event_value(
706 self.mvcc,
707 self.dirty,
708 &leaf.values[values_idx],
709 )?;
710 emit_event = Some((event_position, event_record));
712
713 if !self.backwards {
714 if values_idx + 1 < leaf.values.len() {
715 push_revisit = Some((page_id, Some(values_idx + 1)));
717 } else {
719 remove_page = true;
721 }
723 } else if values_idx > 0 {
724 push_revisit = Some((page_id, Some(values_idx - 1)));
726 } else {
728 remove_page = true;
730 }
732 } else {
733 remove_page = true;
736 }
737 } else {
738 remove_page = true;
740 }
741 }
742 _ => {
743 return Err(DcbError::DatabaseCorrupted(format!(
744 "Expected EventInternal or EventLeaf node in event tree, got {}",
745 page_ref.node.type_name()
746 )));
747 }
748 }
749 }
750
751 if let Some(revisit) = push_revisit {
753 self.stack.push(revisit);
755 }
756 if let Some((child_id, child_start_idx)) = push_child {
757 self.stack.push((child_id, child_start_idx));
758 }
759 if let Some((event_position, event_record)) = emit_event {
760 result.push((event_position, event_record));
761 }
762 if remove_page {
763 self.page_cache.remove(&page_id);
764 }
765 }
766 Ok(result)
767 }
768}
769
770#[cfg(test)]
771mod tests {
772 use super::*;
773 use crate::mvcc::{Mvcc, StorageOptions};
774 use crate::node::Node;
775 use rand::random;
776 use serial_test::serial;
777 use tempfile::tempdir;
778
779 static VERBOSE: bool = false;
780
781 fn construct_db(page_size: usize) -> (tempfile::TempDir, Mvcc) {
783 let temp_dir = tempdir().unwrap();
784 let db_path = temp_dir.path().join("mvcc-test.db");
785 let db = Mvcc::new(
786 VERBOSE,
787 StorageOptions::default()
788 .db_path(db_path)
789 .page_size(page_size),
790 )
791 .unwrap();
792 (temp_dir, db)
793 }
794
795 #[test]
796 #[serial]
797 fn test_append_event_to_empty_leaf_root() {
798 let (_temp_dir, db) = construct_db(64);
800
801 let mut writer = db.writer().unwrap();
803
804 let position = writer.issue_position();
806
807 let record = EventRecord {
809 event_type: "UserCreated".to_string(),
810 data: vec![1, 2, 3, 4],
811 tags: vec!["users".to_string(), "creation".to_string()],
812 uuid: None,
813 metadata: Vec::new(),
814 };
815
816 event_tree_append(&db, &mut writer, record.clone(), position).unwrap();
818
819 let new_root_id = writer.events_tree_root_id;
821 assert!(writer.dirty.contains_key(&new_root_id));
822 let page = writer.dirty.get(&new_root_id).unwrap();
823 match &page.node {
824 Node::EventLeaf(node) => {
825 assert_eq!(vec![position], node.keys);
826 assert_eq!(
827 vec![crate::events_tree_nodes::EventValue::Inline(record.clone())],
828 node.values
829 );
830 }
831 _ => panic!("Expected EventLeaf node"),
832 }
833
834 db.commit(&mut writer).unwrap();
836
837 let header_page = db.get_latest_header_page().unwrap();
839 let header = header_page.as_header_node().unwrap();
840 let persisted_page = db.read_page(header.events_tree_root_id).unwrap();
841 match &persisted_page.node {
842 Node::EventLeaf(node) => {
843 assert_eq!(vec![position], node.keys);
844 assert_eq!(
845 vec![crate::events_tree_nodes::EventValue::Inline(record)],
846 node.values
847 );
848 }
849 _ => panic!("Expected EventLeaf node after commit"),
850 }
851 }
852
853 #[test]
854 #[serial]
855 fn test_insert_events_until_split_leaf_one_writer() {
856 let (_temp_dir, db) = construct_db(256);
858
859 let mut writer = db.writer().unwrap();
861
862 let mut has_split_leaf = false;
863 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
864
865 while !has_split_leaf {
867 let position = writer.issue_position();
869 let record = EventRecord {
870 event_type: "UserCreated".to_string(),
871 data: (0..8).map(|_| random::<u8>()).collect(),
872 tags: vec!["users".to_string(), "creation".to_string()],
873 uuid: None,
874 metadata: Vec::new(),
875 };
876 appended.push((position, record.clone()));
877
878 event_tree_append(&db, &mut writer, record, position).unwrap();
880
881 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
883 match &root_page.node {
884 Node::EventInternal(_) => {
885 has_split_leaf = true;
886 }
887 _ => {}
888 }
889 }
890
891 let mut copy_inserted = appended.clone();
893
894 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
896 let root_node = match &root_page.node {
897 Node::EventInternal(node) => node,
898 _ => panic!("Expected EventInternal node"),
899 };
900
901 for (i, &child_id) in root_node.child_ids.iter().enumerate() {
903 let child_page = writer.dirty.get(&child_id).unwrap();
904 assert_eq!(child_id, child_page.page_id);
905
906 let child_node = match &child_page.node {
907 Node::EventLeaf(node) => node,
908 _ => panic!("Expected EventLeaf node"),
909 };
910
911 if i > 0 {
913 assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
914 }
915
916 for (k, &key) in child_node.keys.iter().enumerate() {
918 let record = &child_node.values[k];
919 let (appended_position, appended_record) = copy_inserted.remove(0);
920 assert_eq!(appended_position, key);
921 assert_eq!(appended_record, record.clone());
922 }
923 }
924 }
925
926 #[test]
927 #[serial]
928 fn test_insert_events_until_split_leaf_many_writers() {
929 let (_temp_dir, db) = construct_db(256);
931
932 let mut has_split_leaf = false;
933 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
934
935 while !has_split_leaf {
937 let mut writer = db.writer().unwrap();
939
940 let position = writer.issue_position();
942 let record = EventRecord {
943 event_type: "UserCreated".to_string(),
944 data: (0..8).map(|_| random::<u8>()).collect(),
945 tags: vec!["users".to_string(), "creation".to_string()],
946 uuid: None,
947 metadata: Vec::new(),
948 };
949 appended.push((position, record.clone()));
950
951 event_tree_append(&db, &mut writer, record, position).unwrap();
953
954 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
956 match &root_page.node {
957 Node::EventInternal(_) => {
958 has_split_leaf = true;
959 }
960 _ => {}
961 }
962
963 db.commit(&mut writer).unwrap();
964 }
965
966 let mut copy_inserted = appended.clone();
968
969 let writer = db.writer().unwrap();
971
972 let root_page = db.read_page(writer.events_tree_root_id).unwrap();
974 let root_node = match &root_page.node {
975 Node::EventInternal(node) => node,
976 _ => panic!("Expected EventInternal node"),
977 };
978
979 for (i, &child_id) in root_node.child_ids.iter().enumerate() {
981 let child_page = db.read_page(child_id).unwrap();
982 assert_eq!(child_id, child_page.page_id);
983
984 let child_node = match &child_page.node {
985 Node::EventLeaf(node) => node,
986 _ => panic!("Expected EventLeaf node"),
987 };
988
989 if i > 0 {
991 assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
992 }
993
994 for (k, &key) in child_node.keys.iter().enumerate() {
996 let record = &child_node.values[k];
997 let (appended_position, appended_record) = copy_inserted.remove(0);
998 assert_eq!(appended_position, key);
999 assert_eq!(appended_record, record.clone());
1000 }
1001 }
1002 }
1003
1004 #[test]
1005 #[serial]
1006 fn test_insert_events_until_split_internal_one_writer() {
1007 let (_temp_dir, db) = construct_db(256);
1009
1010 let mut writer = db.writer().unwrap();
1012
1013 let mut has_split_internal = false;
1014 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1015
1016 while !has_split_internal {
1018 let position = writer.issue_position();
1020 let record = EventRecord {
1021 event_type: "UserCreated".to_string(),
1022 data: (0..8).map(|_| random::<u8>()).collect(),
1023 tags: vec!["users".to_string(), "creation".to_string()],
1024 uuid: None,
1025 metadata: Vec::new(),
1026 };
1027 appended.push((position, record.clone()));
1028
1029 event_tree_append(&db, &mut writer, record, position).unwrap();
1031
1032 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1034 match &root_page.node {
1035 Node::EventInternal(root_node) => {
1036 if !root_node.child_ids.is_empty() {
1038 let child_id = root_node.child_ids[0];
1039 if let Some(child_page) = writer.dirty.get(&child_id) {
1040 match &child_page.node {
1041 Node::EventInternal(_) => {
1042 has_split_internal = true;
1043 }
1044 _ => {}
1045 }
1046 }
1047 }
1048 }
1049 _ => {}
1050 }
1051 }
1052
1053 let mut copy_inserted = appended.clone();
1055
1056 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1058 let root_node = match &root_page.node {
1059 Node::EventInternal(node) => node,
1060 _ => panic!("Expected EventInternal node"),
1061 };
1062
1063 for &child_id in root_node.child_ids.iter() {
1065 let child_page = writer.dirty.get(&child_id).unwrap();
1066 assert_eq!(child_id, child_page.page_id);
1067
1068 let child_node = match &child_page.node {
1069 Node::EventInternal(node) => node,
1070 _ => panic!("Expected EventInternal node"),
1071 };
1072
1073 for &grand_child_id in child_node.child_ids.iter() {
1074 let grand_child_page = writer.dirty.get(&grand_child_id).unwrap();
1075 assert_eq!(grand_child_id, grand_child_page.page_id);
1076
1077 let grand_child_node = match &grand_child_page.node {
1078 Node::EventLeaf(node) => node,
1079 _ => panic!("Expected EventLeaf node"),
1080 };
1081
1082 for (k, &key) in grand_child_node.keys.iter().enumerate() {
1084 let record = &grand_child_node.values[k];
1085 let (appended_position, appended_record) = copy_inserted.remove(0);
1086 assert_eq!(appended_position, key);
1087 assert_eq!(appended_record, record.clone());
1088 }
1089 }
1090 }
1091 }
1092
1093 #[test]
1094 #[serial]
1095 fn test_insert_events_until_split_internal_many_writers() {
1096 let (_temp_dir, db) = construct_db(512);
1098
1099 let mut has_split_internal = false;
1100 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1101
1102 while !has_split_internal {
1104 let mut writer = db.writer().unwrap();
1106
1107 let position = writer.issue_position();
1109 let record = EventRecord {
1110 event_type: "UserCreated".to_string(),
1111 data: (0..8).map(|_| random::<u8>()).collect(),
1112 tags: vec!["users".to_string(), "creation".to_string()],
1113 uuid: None,
1114 metadata: Vec::new(),
1115 };
1116 appended.push((position, record.clone()));
1117
1118 event_tree_append(&db, &mut writer, record, position).unwrap();
1120
1121 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1123 match &root_page.node {
1124 Node::EventInternal(root_node) => {
1125 if !root_node.child_ids.is_empty() {
1127 let child_id = root_node.child_ids[0];
1128 if let Some(child_page) = writer.dirty.get(&child_id) {
1129 match &child_page.node {
1130 Node::EventInternal(_) => {
1131 has_split_internal = true;
1132 }
1133 _ => {}
1134 }
1135 }
1136 }
1137 }
1138 _ => {}
1139 }
1140 db.commit(&mut writer).unwrap();
1141 }
1142
1143 let mut copy_inserted = appended.clone();
1145
1146 let writer = db.writer().unwrap();
1148
1149 let root_page = db.read_page(writer.events_tree_root_id).unwrap();
1151 let root_node = match &root_page.node {
1152 Node::EventInternal(node) => node,
1153 _ => panic!("Expected EventInternal node"),
1154 };
1155
1156 for &child_id in root_node.child_ids.iter() {
1158 let child_page = db.read_page(child_id).unwrap();
1159 assert_eq!(child_id, child_page.page_id);
1160
1161 let child_node = match &child_page.node {
1162 Node::EventInternal(node) => node,
1163 _ => panic!("Expected EventInternal node"),
1164 };
1165
1166 for &grand_child_id in child_node.child_ids.iter() {
1167 let grand_child_page = db.read_page(grand_child_id).unwrap();
1168 assert_eq!(grand_child_id, grand_child_page.page_id);
1169
1170 let grand_child_node = match &grand_child_page.node {
1171 Node::EventLeaf(node) => node,
1172 _ => panic!("Expected EventLeaf node"),
1173 };
1174
1175 for (k, &key) in grand_child_node.keys.iter().enumerate() {
1177 let record = &grand_child_node.values[k];
1178 let (appended_position, appended_record) = copy_inserted.remove(0);
1179 assert_eq!(appended_position, key);
1181 assert_eq!(appended_record, record.clone());
1182 }
1183 }
1184 }
1185 assert_eq!(0, copy_inserted.len());
1186 }
1187
1188 #[test]
1189 #[serial]
1190 fn test_read_events_all() {
1191 let (_temp_dir, db) = construct_db(512);
1193
1194 let mut has_split_internal = false;
1195 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1196
1197 while !has_split_internal {
1199 let mut writer = db.writer().unwrap();
1201
1202 let position = writer.issue_position();
1204 let record = EventRecord {
1205 event_type: "UserCreated".to_string(),
1206 data: (0..8).map(|_| random::<u8>()).collect(),
1207 tags: vec!["users".to_string(), "creation".to_string()],
1208 uuid: None,
1209 metadata: Vec::new(),
1210 };
1211 appended.push((position, record.clone()));
1212
1213 event_tree_append(&db, &mut writer, record, position).unwrap();
1215
1216 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1218 match &root_page.node {
1219 Node::EventInternal(root_node) => {
1220 if !root_node.child_ids.is_empty() {
1222 let child_id = root_node.child_ids[0];
1223 if let Some(child_page) = writer.dirty.get(&child_id) {
1224 match &child_page.node {
1225 Node::EventInternal(_) => {
1226 has_split_internal = true;
1227 }
1228 _ => {}
1229 }
1230 }
1231 }
1232 }
1233 _ => {}
1234 }
1235 db.commit(&mut writer).unwrap();
1236 }
1237
1238 let copy_inserted = appended.clone();
1240
1241 let reader = db.reader().unwrap();
1243 let events_tree_root_id = reader.events_tree_root_id;
1244 let reader_tsn = reader.tsn;
1245
1246 let dirty = HashMap::new();
1247 let mut events_iterator = EventIterator::new(&db, &dirty, events_tree_root_id, None, false);
1248
1249 {
1251 assert!(
1253 db.reader_tsns.contains_tsn(reader_tsn),
1254 "TSN should remain registered until reader is dropped"
1255 );
1256 }
1257
1258 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1260 loop {
1261 let batch = events_iterator.next_batch(3, None).unwrap();
1262 if batch.is_empty() {
1263 break;
1264 }
1265 scanned.extend(batch);
1266
1267 assert!(
1270 db.reader_tsns.contains_tsn(reader_tsn),
1271 "TSN should remain registered until reader is dropped"
1272 );
1273 }
1274
1275 assert_eq!(copy_inserted.len(), scanned.len());
1276 for (i, expected) in copy_inserted.iter().enumerate() {
1277 assert_eq!(expected.0, scanned[i].0);
1278 assert_eq!(expected.1, scanned[i].1);
1279 }
1280
1281 let dirty = HashMap::new();
1283 for (pos, expected_rec) in copy_inserted.iter() {
1284 let found = event_tree_lookup(&db, &dirty, events_tree_root_id, *pos).unwrap();
1285 assert_eq!(expected_rec, &found);
1286 }
1287
1288 assert!(
1290 events_iterator.page_cache.is_empty(),
1291 "EventIterator page_cache should be empty after full scan"
1292 );
1293
1294 {
1296 assert!(
1298 db.reader_tsns.contains_tsn(reader_tsn),
1299 "TSN should remain registered until reader is dropped"
1300 );
1301 }
1302
1303 drop(reader);
1305 {
1306 assert!(
1308 !db.reader_tsns.contains_tsn(reader_tsn),
1309 "TSN should be removed after reader is dropped"
1310 );
1311 assert_eq!(0, db.reader_tsns.live_reader_count());
1312 }
1313 }
1314
1315 #[test]
1316 #[serial]
1317 fn test_read_events_from_forwards() {
1318 let (_temp_dir, db) = construct_db(128);
1320
1321 let mut has_split_internal = false;
1322 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1323
1324 while !has_split_internal {
1326 let mut writer = db.writer().unwrap();
1328
1329 let position = writer.issue_position();
1331 let record = EventRecord {
1332 event_type: "UserCreated".to_string(),
1333 data: (0..8).map(|_| random::<u8>()).collect(),
1334 tags: vec!["users".to_string(), "creation".to_string()],
1335 uuid: None,
1336 metadata: Vec::new(),
1337 };
1338 appended.push((position, record.clone()));
1339
1340 event_tree_append(&db, &mut writer, record, position).unwrap();
1342
1343 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1345 match &root_page.node {
1346 Node::EventInternal(root_node) => {
1347 if !root_node.child_ids.is_empty() {
1349 let child_id = root_node.child_ids[0];
1350 if let Some(child_page) = writer.dirty.get(&child_id) {
1351 match &child_page.node {
1352 Node::EventInternal(_) => {
1353 has_split_internal = true;
1354 }
1355 _ => {}
1356 }
1357 }
1358 }
1359 }
1360 _ => {}
1361 }
1362 db.commit(&mut writer).unwrap();
1363 }
1364
1365 let appended_len = appended.len();
1366 for i in 0..appended_len + 2 {
1370 let from = Position(i as u64);
1371
1372 let reader = db.reader().unwrap();
1374 let events_tree_root_id = reader.events_tree_root_id;
1375 let reader_tsn = reader.tsn;
1376 let dirty = HashMap::new();
1377 let mut events_iterator =
1378 EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), false);
1379
1380 {
1382 assert!(
1384 db.reader_tsns.contains_tsn(reader_tsn),
1385 "TSN should remain registered until reader is dropped"
1386 );
1387 }
1388
1389 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1391 loop {
1392 let batch = events_iterator.next_batch(3, None).unwrap();
1393 if batch.is_empty() {
1394 break;
1395 }
1396 scanned.extend(batch);
1397
1398 assert!(
1401 db.reader_tsns.contains_tsn(reader_tsn),
1402 "TSN should remain registered until reader is dropped"
1403 );
1404 }
1405
1406 let num_to_skip = i.max(1) - 1;
1408 let expected: Vec<(Position, EventRecord)> =
1409 appended.clone().into_iter().skip(num_to_skip).collect();
1410 assert_eq!(expected.len(), scanned.len());
1411 for (i, exp) in expected.iter().enumerate() {
1413 assert_eq!(exp.0, scanned[i].0);
1414 assert_eq!(exp.1, scanned[i].1);
1415 }
1416
1417 assert!(
1419 events_iterator.page_cache.is_empty(),
1420 "EventIterator page_cache should be empty after filtered scan"
1421 );
1422
1423 drop(reader);
1425 {
1426 assert!(
1428 !db.reader_tsns.contains_tsn(reader_tsn),
1429 "TSN should be removed after reader is dropped"
1430 );
1431 }
1432 }
1433 }
1434
1435 #[test]
1436 #[serial]
1437 fn test_read_events_from_backwards() {
1438 let (_temp_dir, db) = construct_db(128);
1440
1441 let mut has_split_internal = false;
1442 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1443
1444 while !has_split_internal {
1446 let mut writer = db.writer().unwrap();
1448
1449 let position = writer.issue_position();
1451 let record = EventRecord {
1452 event_type: "UserCreated".to_string(),
1453 data: (0..8).map(|_| random::<u8>()).collect(),
1454 tags: vec!["users".to_string(), "creation".to_string()],
1455 uuid: None,
1456 metadata: Vec::new(),
1457 };
1458 appended.push((position, record.clone()));
1459
1460 event_tree_append(&db, &mut writer, record, position).unwrap();
1462
1463 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1465 match &root_page.node {
1466 Node::EventInternal(root_node) => {
1467 if !root_node.child_ids.is_empty() {
1469 let child_id = root_node.child_ids[0];
1470 if let Some(child_page) = writer.dirty.get(&child_id) {
1471 match &child_page.node {
1472 Node::EventInternal(_) => {
1473 has_split_internal = true;
1474 }
1475 _ => {}
1476 }
1477 }
1478 }
1479 }
1480 _ => {}
1481 }
1482 db.commit(&mut writer).unwrap();
1483 }
1484
1485 let appended_len = appended.len();
1486 for i in 0..appended_len + 2 {
1490 let from = Position(i as u64);
1491
1492 let reader = db.reader().unwrap();
1494 let events_tree_root_id = reader.events_tree_root_id;
1495 let reader_tsn = reader.tsn;
1496 let dirty = HashMap::new();
1497 let mut events_iterator =
1498 EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), true);
1499
1500 {
1502 assert!(
1504 db.reader_tsns.contains_tsn(reader_tsn),
1505 "TSN should remain registered until reader is dropped"
1506 );
1507 }
1508
1509 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1511 loop {
1512 let batch = events_iterator.next_batch(3, None).unwrap();
1513 if batch.is_empty() {
1514 break;
1515 }
1516 scanned.extend(batch);
1517
1518 assert!(
1521 db.reader_tsns.contains_tsn(reader_tsn),
1522 "TSN should remain registered until reader is dropped"
1523 );
1524 }
1525
1526 let mut expected: Vec<(Position, EventRecord)> = appended.clone();
1528 expected.truncate(i);
1529 expected.reverse();
1530 assert_eq!(expected.len(), scanned.len());
1531 for (i, exp) in expected.iter().enumerate() {
1533 assert_eq!(exp.0, scanned[i].0);
1534 assert_eq!(exp.1, scanned[i].1);
1535 }
1536
1537 assert!(
1539 events_iterator.page_cache.is_empty(),
1540 "EventIterator page_cache should be empty after filtered scan"
1541 );
1542
1543 drop(reader);
1545 {
1546 assert!(
1548 !db.reader_tsns.contains_tsn(reader_tsn),
1549 "TSN should be removed after reader is dropped"
1550 );
1551 }
1552 }
1553 }
1554
1555 #[test]
1556 #[serial]
1557 fn test_large_event_data_exact_page_size() {
1558 let (_tmp, db) = construct_db(512);
1559 let mut writer = db.writer().unwrap();
1561 let pos = writer.issue_position();
1562 let data = vec![0xAB; 512];
1563 let event = EventRecord {
1564 event_type: "Big".into(),
1565 data: data.clone(),
1566 tags: vec![],
1567 uuid: None,
1568 metadata: Vec::new(),
1569 };
1570 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1571 db.commit(&mut writer).unwrap();
1572
1573 let reader = db.reader().unwrap();
1575 let dirty = HashMap::new();
1576 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1577 assert_eq!(event, got);
1578
1579 let header_page = db.get_latest_header_page().unwrap();
1581 let header = header_page.as_header_node().unwrap();
1582 let root = db.read_page(header.events_tree_root_id).unwrap();
1583 match &root.node {
1584 Node::EventInternal(internal) => {
1585 let leaf_id = *internal.child_ids.last().unwrap();
1587 let leaf_page = db.read_page(leaf_id).unwrap();
1588 match &leaf_page.node {
1589 Node::EventLeaf(leaf) => match &leaf.values[0] {
1590 EventValue::Overflow { data_len, .. } => {
1591 assert_eq!(*data_len as usize, data.len())
1592 }
1593 _ => panic!("Expected Overflow for large event"),
1594 },
1595 _ => panic!("Expected EventLeaf child"),
1596 }
1597 }
1598 Node::EventLeaf(leaf) => match &leaf.values[0] {
1599 EventValue::Overflow { data_len, .. } => {
1600 assert_eq!(*data_len as usize, data.len())
1601 }
1602 _ => panic!("Expected Overflow for large event"),
1603 },
1604 _ => panic!("Unexpected root node type"),
1605 }
1606 }
1607
1608 #[test]
1609 #[serial]
1610 fn test_large_event_data_four_times_page_size() {
1611 let (_tmp, db) = construct_db(512);
1612 let mut writer = db.writer().unwrap();
1613 let pos = writer.issue_position();
1614 let data = vec![0xCD; 512 * 4];
1615 let event = EventRecord {
1616 event_type: "Bigger".into(),
1617 data: data.clone(),
1618 tags: vec![],
1619 uuid: None,
1620 metadata: Vec::new(),
1621 };
1622 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1623 db.commit(&mut writer).unwrap();
1624
1625 let reader = db.reader().unwrap();
1627 let dirty = HashMap::new();
1628 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1629 assert_eq!(event, got);
1630
1631 let header_page = db.get_latest_header_page().unwrap();
1633 let header = header_page.as_header_node().unwrap();
1634 let root = db.read_page(header.events_tree_root_id).unwrap();
1635 let check_leaf = |leaf: &EventLeafNode| match &leaf.values[0] {
1636 EventValue::Overflow { data_len, .. } => {
1637 assert_eq!(*data_len as usize, data.len())
1638 }
1639 _ => panic!("Expected Overflow for very large event"),
1640 };
1641 match &root.node {
1642 Node::EventInternal(internal) => {
1643 let leaf_id = *internal.child_ids.last().unwrap();
1644 let leaf_page = db.read_page(leaf_id).unwrap();
1645 match &leaf_page.node {
1646 Node::EventLeaf(leaf) => check_leaf(&leaf),
1647 _ => panic!("Expected leaf"),
1648 }
1649 }
1650 Node::EventLeaf(leaf) => check_leaf(&leaf),
1651 _ => panic!("Unexpected root node type"),
1652 }
1653 }
1654
1655 #[test]
1656 #[serial]
1657 fn test_inline_event_metadata_roundtrip() {
1658 let (_tmp, db) = construct_db(4096);
1659 let mut writer = db.writer().unwrap();
1660 let pos = writer.issue_position();
1661
1662 let mut metadata = Vec::new();
1663 metadata.push(("source".to_string(), "web".to_string()));
1664 metadata.push(("correlation_id".to_string(), "abc-123".to_string()));
1665
1666 let event = EventRecord {
1667 event_type: "SmallWithMetadata".into(),
1668 data: vec![1, 2, 3, 4],
1669 tags: vec!["t".into()],
1670 uuid: None,
1671 metadata: metadata.clone(),
1672 };
1673 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1674 db.commit(&mut writer).unwrap();
1675
1676 let reader = db.reader().unwrap();
1678 let dirty = HashMap::new();
1679 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1680 assert_eq!(event, got);
1681 assert_eq!(metadata, got.metadata);
1682
1683 let header_page = db.get_latest_header_page().unwrap();
1685 let header = header_page.as_header_node().unwrap();
1686 let root = db.read_page(header.events_tree_root_id).unwrap();
1687 match &root.node {
1688 Node::EventLeaf(leaf) => match &leaf.values[0] {
1689 EventValue::Inline(rec) => assert_eq!(metadata, rec.metadata),
1690 _ => panic!("Expected Inline value"),
1691 },
1692 _ => panic!("Expected EventLeaf root for a single small event"),
1693 }
1694 }
1695
1696 #[test]
1697 #[serial]
1698 fn test_overflow_event_metadata_roundtrip() {
1699 let (_tmp, db) = construct_db(512);
1700 let mut writer = db.writer().unwrap();
1701 let pos = writer.issue_position();
1702
1703 let mut metadata = Vec::new();
1704 metadata.push(("source".to_string(), "bulk-import".to_string()));
1705 metadata.push(("schema".to_string(), "v3".to_string()));
1706
1707 let data = vec![0xAB; 512 * 4];
1709 let event = EventRecord {
1710 event_type: "BigWithMetadata".into(),
1711 data: data.clone(),
1712 tags: vec!["t".into()],
1713 uuid: Some(uuid::Uuid::new_v4()),
1714 metadata: metadata.clone(),
1715 };
1716 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1717 db.commit(&mut writer).unwrap();
1718
1719 let reader = db.reader().unwrap();
1722 let dirty = HashMap::new();
1723 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1724 assert_eq!(event, got);
1725 assert_eq!(data, got.data);
1726 assert_eq!(metadata, got.metadata);
1727
1728 let header_page = db.get_latest_header_page().unwrap();
1731 let header = header_page.as_header_node().unwrap();
1732 let root = db.read_page(header.events_tree_root_id).unwrap();
1733 let check_leaf = |leaf: &EventLeafNode| match &leaf.values[0] {
1734 EventValue::Overflow {
1735 data_len,
1736 metadata_len,
1737 ..
1738 } => {
1739 assert_eq!(*data_len as usize, data.len());
1740 assert!(*metadata_len > 0, "expected non-zero metadata_len");
1741 }
1742 _ => panic!("Expected Overflow value for large event"),
1743 };
1744 match &root.node {
1745 Node::EventInternal(internal) => {
1746 let leaf_id = *internal.child_ids.last().unwrap();
1747 let leaf_page = db.read_page(leaf_id).unwrap();
1748 match &leaf_page.node {
1749 Node::EventLeaf(leaf) => check_leaf(leaf),
1750 _ => panic!("Expected EventLeaf child"),
1751 }
1752 }
1753 Node::EventLeaf(leaf) => check_leaf(leaf),
1754 _ => panic!("Unexpected root node type"),
1755 }
1756 }
1757
1758 #[test]
1759 #[serial]
1760 fn test_direct_overflow_event_metadata_roundtrip() {
1761 let (_tmp, db) = construct_db(4096);
1764 let mut writer = db.writer().unwrap();
1765 let pos = writer.issue_position();
1766
1767 let mut metadata = Vec::new();
1768 metadata.push(("origin".to_string(), "direct-overflow".to_string()));
1769
1770 let data = vec![0xCD; (u16::MAX as usize) + 1024];
1771 let event = EventRecord {
1772 event_type: "HugeWithMetadata".into(),
1773 data: data.clone(),
1774 tags: vec!["t".into()],
1775 uuid: None,
1776 metadata: metadata.clone(),
1777 };
1778 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1779 db.commit(&mut writer).unwrap();
1780
1781 let reader = db.reader().unwrap();
1782 let dirty = HashMap::new();
1783 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1784 assert_eq!(event, got);
1785 assert_eq!(data, got.data);
1786 assert_eq!(metadata, got.metadata);
1787 }
1788
1789 #[test]
1790 #[serial]
1791 fn test_append_rejects_oversized_metadata_key() {
1792 use crate::events_tree_nodes::MAX_METADATA_ENTRY_LEN;
1793
1794 let (_tmp, db) = construct_db(4096);
1795 let mut writer = db.writer().unwrap();
1796 let pos = writer.issue_position();
1797
1798 let mut metadata = Vec::new();
1799 metadata.push(("x".repeat(MAX_METADATA_ENTRY_LEN + 1), "source".to_string()));
1800 let event = EventRecord {
1801 event_type: "TooBigMetadata".into(),
1802 data: vec![1, 2, 3],
1803 tags: vec!["t".into()],
1804 uuid: None,
1805 metadata,
1806 };
1807
1808 match event_tree_append(&db, &mut writer, event, pos) {
1811 Err(DcbError::InvalidArgument(_)) => {}
1812 other => panic!("Expected InvalidArgument, got {other:?}"),
1813 }
1814
1815 db.commit(&mut writer).unwrap();
1818 let reader = db.reader().unwrap();
1819 let dirty = HashMap::new();
1820 assert!(
1821 event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).is_err(),
1822 "rejected event must not be retrievable"
1823 );
1824 }
1825
1826 #[test]
1827 #[serial]
1828 fn test_append_rejects_oversized_metadata_value() {
1829 use crate::events_tree_nodes::MAX_METADATA_ENTRY_LEN;
1830
1831 let (_tmp, db) = construct_db(4096);
1832 let mut writer = db.writer().unwrap();
1833 let pos = writer.issue_position();
1834
1835 let mut metadata = Vec::new();
1836 metadata.push(("source".to_string(), "x".repeat(MAX_METADATA_ENTRY_LEN + 1)));
1837 let event = EventRecord {
1838 event_type: "TooBigMetadata".into(),
1839 data: vec![1, 2, 3],
1840 tags: vec!["t".into()],
1841 uuid: None,
1842 metadata,
1843 };
1844
1845 match event_tree_append(&db, &mut writer, event, pos) {
1848 Err(DcbError::InvalidArgument(_)) => {}
1849 other => panic!("Expected InvalidArgument, got {other:?}"),
1850 }
1851
1852 db.commit(&mut writer).unwrap();
1855 let reader = db.reader().unwrap();
1856 let dirty = HashMap::new();
1857 assert!(
1858 event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).is_err(),
1859 "rejected event must not be retrievable"
1860 );
1861 }
1862
1863 #[test]
1864 #[serial]
1865 fn test_append_rejects_oversized_event_type() {
1866 use crate::events_tree_nodes::MAX_EVENT_TYPE_LEN;
1867
1868 let (_tmp, db) = construct_db(4096);
1869 let mut writer = db.writer().unwrap();
1870 let pos = writer.issue_position();
1871
1872 let event = EventRecord {
1873 event_type: "x".repeat(MAX_EVENT_TYPE_LEN + 1),
1874 data: vec![1, 2, 3],
1875 tags: vec!["t".into()],
1876 uuid: None,
1877 metadata: Vec::new(),
1878 };
1879
1880 match event_tree_append(&db, &mut writer, event, pos) {
1881 Err(DcbError::InvalidArgument(_)) => {}
1882 other => panic!("Expected InvalidArgument, got {other:?}"),
1883 }
1884 }
1885
1886 #[test]
1887 #[serial]
1888 fn test_append_rejects_oversized_tag() {
1889 use crate::events_tree_nodes::MAX_TAG_LEN;
1890
1891 let (_tmp, db) = construct_db(4096);
1892 let mut writer = db.writer().unwrap();
1893 let pos = writer.issue_position();
1894
1895 let event = EventRecord {
1896 event_type: "Type".into(),
1897 data: vec![1, 2, 3],
1898 tags: vec!["x".repeat(MAX_TAG_LEN + 1)],
1899 uuid: None,
1900 metadata: Vec::new(),
1901 };
1902
1903 match event_tree_append(&db, &mut writer, event, pos) {
1904 Err(DcbError::InvalidArgument(_)) => {}
1905 other => panic!("Expected InvalidArgument, got {other:?}"),
1906 }
1907 }
1908 #[test]
1959 #[serial]
1960 fn test_read_events_all_backwards_with_internal_nodes() {
1961 let (_temp_dir, db) = construct_db(128);
1963
1964 let mut has_split_internal = false;
1965 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1966
1967 while !has_split_internal {
1969 let mut writer = db.writer().unwrap();
1970 let position = writer.issue_position();
1971 let record = EventRecord {
1972 event_type: "TestEvent".to_string(),
1973 data: vec![1, 2, 3, 4],
1974 tags: vec!["test".to_string()],
1975 uuid: None,
1976 metadata: Vec::new(),
1977 };
1978 appended.push((position, record.clone()));
1979 event_tree_append(&db, &mut writer, record, position).unwrap();
1980
1981 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1982 if let Node::EventInternal(root_node) = &root_page.node {
1983 if !root_node.child_ids.is_empty() {
1984 let child_id = root_node.child_ids[0];
1985 if let Some(child_page) = writer.dirty.get(&child_id) {
1986 if let Node::EventInternal(_) = &child_page.node {
1987 has_split_internal = true;
1988 }
1989 }
1990 }
1991 }
1992 db.commit(&mut writer).unwrap();
1993 }
1994
1995 let reader = db.reader().unwrap();
1997 let dirty = HashMap::new();
1998 let mut events_iterator = EventIterator::new(
1999 &db,
2000 &dirty,
2001 reader.events_tree_root_id,
2002 None, true, );
2005
2006 let mut scanned = Vec::new();
2007 loop {
2008 let batch = events_iterator.next_batch(10, None).unwrap();
2009 if batch.is_empty() {
2010 break;
2011 }
2012 scanned.extend(batch);
2013 }
2014
2015 let mut expected = appended.clone();
2017 expected.reverse();
2018
2019 assert_eq!(scanned.len(), expected.len());
2020 for i in 0..expected.len() {
2021 assert_eq!(scanned[i].0, expected[i].0);
2022 }
2023 }
2024
2025 #[test]
2026 #[serial]
2027 fn test_append_event_with_tag_larger_than_page_size() {
2028 use crate::events_tree_nodes::MAX_TAG_LEN;
2033
2034 let page_size = 512;
2036 let (_tmp, db) = construct_db(page_size);
2037 let mut writer = db.writer().unwrap();
2038 let pos = writer.issue_position();
2039
2040 let large_tag = "t".repeat(600);
2043 assert!(large_tag.len() > page_size);
2044 assert!(large_tag.len() < MAX_TAG_LEN);
2045
2046 let event = EventRecord {
2047 event_type: "Type".into(),
2048 data: vec![1, 2, 3],
2049 tags: vec![large_tag],
2050 uuid: None,
2051 metadata: Vec::new(),
2052 };
2053
2054 match event_tree_append(&db, &mut writer, event, pos) {
2055 Err(DcbError::InvalidArgument(message)) => {
2056 assert!(
2057 message.contains("event too large for page size"),
2058 "unexpected InvalidArgument message: {message}"
2059 );
2060 }
2061 other => panic!("Expected DatabaseCorrupted, got {other:?}"),
2062 }
2063 }
2064}