1use crate::common::PageID;
2use crate::common::Position;
3use crate::events_tree_nodes::{
4 EventInternalNode, EventLeafNode, EventOverflowNode, EventRecord, EventValue,
5};
6use crate::mvcc::{Mvcc, Writer};
7use crate::node::Node;
8use crate::page::{PAGE_HEADER_SIZE, Page};
9use std::collections::HashMap;
10use umadb_dcb::{DCBError, DCBResult};
11
12fn write_overflow_chain(mvcc: &Mvcc, writer: &mut Writer, data: &[u8]) -> DCBResult<PageID> {
14 let payload_cap = mvcc.page_size.saturating_sub(PAGE_HEADER_SIZE + 8);
16 if payload_cap == 0 {
17 return Err(DCBError::DatabaseCorrupted(
18 "Page size too small to store overflow data".to_string(),
19 ));
20 }
21 let mut chunks: Vec<&[u8]> = Vec::new();
23 let mut i = 0;
24 while i < data.len() {
25 let end = (i + payload_cap).min(data.len());
26 chunks.push(&data[i..end]);
27 i = end;
28 }
29 if chunks.is_empty() {
30 let page_id = writer.alloc_page_id();
32 let node = EventOverflowNode {
33 next: PageID(0),
34 data: Vec::new(),
35 };
36 let page = Page::new(page_id, Node::EventOverflow(node));
37 writer.insert_dirty(page)?;
38 return Ok(page_id);
39 }
40 let mut next_id = PageID(0);
41 for chunk in chunks.iter().rev() {
42 let page_id = writer.alloc_page_id();
43 let node = EventOverflowNode {
44 next: next_id,
45 data: (*chunk).to_vec(),
46 };
47 let page = Page::new(page_id, Node::EventOverflow(node));
48 writer.insert_dirty(page)?;
49 next_id = page_id;
50 }
51 Ok(next_id)
52}
53
54fn read_overflow_chain(
55 mvcc: &Mvcc,
56 dirty: &HashMap<PageID, Page>,
57 mut page_id: PageID,
58) -> DCBResult<Vec<u8>> {
59 let mut out: Vec<u8> = Vec::new();
60 while page_id.0 != 0 {
61 let page = if let Some(p) = dirty.get(&page_id) {
63 p.clone()
64 } else {
65 mvcc.read_page(page_id)?
66 };
67 match page.node {
68 Node::EventOverflow(node) => {
69 out.extend_from_slice(&node.data);
70 page_id = node.next;
71 }
72 _ => {
73 return Err(DCBError::DatabaseCorrupted(
74 "Expected EventOverflow node".to_string(),
75 ));
76 }
77 }
78 }
79 Ok(out)
80}
81
82fn materialize_event_value(
83 mvcc: &Mvcc,
84 dirty: &HashMap<PageID, Page>,
85 value: &EventValue,
86) -> DCBResult<EventRecord> {
87 match value {
88 EventValue::Inline(rec) => Ok(rec.clone()),
89 EventValue::Overflow {
90 event_type,
91 data_len,
92 tags,
93 root_id,
94 uuid,
95 } => {
96 let data = read_overflow_chain(mvcc, dirty, *root_id)?;
97 if (data.len() as u64) != *data_len {
98 return Err(DCBError::DatabaseCorrupted(
99 "Overflow data length mismatch".to_string(),
100 ));
101 }
102 Ok(EventRecord {
103 event_type: event_type.clone(),
104 data,
105 tags: tags.clone(),
106 uuid: *uuid,
107 })
108 }
109 }
110}
111
112pub fn event_tree_append(
118 mvcc: &Mvcc,
119 writer: &mut Writer,
120 event: EventRecord,
121 position: Position,
122) -> DCBResult<()> {
123 let verbose = mvcc.verbose;
124 if verbose {
125 println!("Appending event: {position:?} {event:?}");
126 println!("Root is {:?}", writer.events_tree_root_id);
127 }
128 let mut current_page_id: PageID = writer.events_tree_root_id;
130
131 let mut stack: Vec<PageID> = Vec::new();
133 loop {
134 let current_page_ref = writer.get_page_ref(mvcc, current_page_id)?;
135 if matches!(current_page_ref.node, Node::EventLeaf(_)) {
136 break;
137 }
138 if let Node::EventInternal(internal_node) = ¤t_page_ref.node {
139 if verbose {
140 println!("{:?} is internal node", current_page_ref.page_id);
141 }
142 stack.push(current_page_id);
143 current_page_id = *internal_node
144 .child_ids
145 .last()
146 .expect("Internal node should have some children");
147 } else {
148 return Err(DCBError::DatabaseCorrupted(
149 "Expected EventInternal node".to_string(),
150 ));
151 }
152 }
153 if verbose {
154 println!("{current_page_id:?} is leaf node");
155 }
156
157 let pending_value = if event.data.len() > u16::MAX as usize {
159 let root_id = write_overflow_chain(mvcc, writer, &event.data)?;
160 EventValue::Overflow {
161 event_type: event.event_type.clone(),
162 data_len: event.data.len() as u64,
163 tags: event.tags.clone(),
164 root_id,
165 uuid: event.uuid,
166 }
167 } else {
168 EventValue::Inline(event)
169 };
170
171 let dirty_page_id = { writer.get_dirty_page_id(current_page_id)? };
173 let replacement_info: Option<(PageID, PageID)> = {
174 if dirty_page_id != current_page_id {
175 Some((current_page_id, dirty_page_id))
176 } else {
177 None
178 }
179 };
180
181 let mut popped: Option<(Position, EventValue)> = None;
183
184 {
186 let dirty_leaf_page = writer.get_mut_dirty(dirty_page_id)?;
187 match &mut dirty_leaf_page.node {
188 Node::EventLeaf(node) => {
189 node.keys.push(position);
190 node.values.push(pending_value);
191
192 let serialized_size = dirty_leaf_page.calc_serialized_size();
194 if serialized_size > mvcc.page_size {
195 if let Node::EventLeaf(dirty_leaf_node) = &mut dirty_leaf_page.node {
196 let (last_key, last_value) = dirty_leaf_node.pop_last_key_and_value()?;
197 if verbose {
198 println!(
199 "Split leaf {:?}: {:?}",
200 dirty_page_id,
201 dirty_leaf_node.clone()
202 );
203 }
204 popped = Some((last_key, last_value));
205 } else {
206 return Err(DCBError::DatabaseCorrupted(
207 "Expected EventLeaf node".to_string(),
208 ));
209 }
210 }
211 }
212 _ => {
213 return Err(DCBError::DatabaseCorrupted(
214 "Expected EventLeaf node at event tree root".to_string(),
215 ));
216 }
217 }
218 }
219
220 let mut split_info: Option<(Position, PageID)> = None;
222
223 if let Some((last_key, mut last_value)) = popped {
224 let new_leaf_page_id = writer.alloc_page_id();
226 let mut new_leaf_node = EventLeafNode {
227 keys: vec![last_key],
228 values: vec![last_value.clone()],
229 };
230 let mut new_leaf_page = Page::new(new_leaf_page_id, Node::EventLeaf(new_leaf_node.clone()));
231 let serialized_size = new_leaf_page.calc_serialized_size();
232 if serialized_size > mvcc.page_size
233 && let EventValue::Inline(rec) = last_value
234 {
235 let root_id = write_overflow_chain(mvcc, writer, &rec.data)?;
236 last_value = EventValue::Overflow {
237 event_type: rec.event_type,
238 data_len: rec.data.len() as u64,
239 tags: rec.tags,
240 root_id,
241 uuid: rec.uuid,
242 };
243 new_leaf_node = EventLeafNode {
244 keys: vec![last_key],
245 values: vec![last_value.clone()],
246 };
247 new_leaf_page = Page::new(new_leaf_page_id, Node::EventLeaf(new_leaf_node.clone()));
248 }
250 if verbose {
257 println!(
258 "Created new leaf {:?}: {:?}",
259 new_leaf_page_id, new_leaf_page.node
260 );
261 }
262 writer.insert_dirty(new_leaf_page)?;
263 if verbose {
264 println!("Promoting {last_key:?} and {new_leaf_page_id:?}");
265 }
266 split_info = Some((last_key, new_leaf_page_id));
267 }
268
269 let mut current_replacement_info = replacement_info;
271 while let Some(parent_page_id) = stack.pop() {
272 let dirty_page_id = { writer.get_dirty_page_id(parent_page_id)? };
274 let parent_replacement_info: Option<(PageID, PageID)> = {
275 if dirty_page_id != parent_page_id {
276 Some((parent_page_id, dirty_page_id))
277 } else {
278 None
279 }
280 };
281 let dirty_internal_page = writer.get_mut_dirty(dirty_page_id)?;
283
284 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
285 if let Some((old_id, new_id)) = current_replacement_info {
286 dirty_internal_node.replace_last_child_id(old_id, new_id)?;
287 if verbose {
288 println!(
289 "Replaced {old_id:?} with {new_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
290 );
291 }
292 } else if verbose {
293 println!("Nothing to replace in {dirty_page_id:?}")
294 }
295 } else {
296 return Err(DCBError::DatabaseCorrupted(
297 "Expected EventInternal node".to_string(),
298 ));
299 }
300
301 if let Some((promoted_key, promoted_page_id)) = split_info {
302 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
303 dirty_internal_node
305 .append_promoted_key_and_page_id(promoted_key, promoted_page_id)?;
306
307 if verbose {
308 println!(
309 "Appended promoted key {promoted_key:?} and child {promoted_page_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
310 );
311 }
312 } else {
313 return Err(DCBError::DatabaseCorrupted(
314 "Expected EventInternal node".to_string(),
315 ));
316 }
317 }
318
319 if dirty_internal_page.calc_serialized_size() > mvcc.page_size {
322 if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
323 if verbose {
324 println!("Splitting internal {dirty_page_id:?}...");
325 }
326 if dirty_internal_node.keys.len() < 3 || dirty_internal_node.child_ids.len() < 4 {
329 return Err(DCBError::DatabaseCorrupted(
330 "Cannot split internal node with too few keys/children".to_string(),
331 ));
332 }
333
334 let (promoted_key, new_keys, new_child_ids) = dirty_internal_node.split_off()?;
336
337 assert_eq!(
339 dirty_internal_node.keys.len() + 1,
340 dirty_internal_node.child_ids.len()
341 );
342
343 let new_internal_node = EventInternalNode {
344 keys: new_keys,
345 child_ids: new_child_ids,
346 };
347
348 assert_eq!(
350 new_internal_node.keys.len() + 1,
351 new_internal_node.child_ids.len()
352 );
353
354 let new_internal_page_id = writer.alloc_page_id();
356 let new_internal_page =
357 Page::new(new_internal_page_id, Node::EventInternal(new_internal_node));
358 if verbose {
359 println!(
360 "Created internal {:?}: {:?}",
361 new_internal_page_id, new_internal_page.node
362 );
363 }
364 writer.insert_dirty(new_internal_page)?;
365
366 split_info = Some((promoted_key, new_internal_page_id));
367 } else {
368 return Err(DCBError::DatabaseCorrupted(
369 "Expected EventInternal node".to_string(),
370 ));
371 }
372 } else {
373 split_info = None;
374 }
375 current_replacement_info = parent_replacement_info;
376 }
377
378 if let Some((old_id, new_id)) = current_replacement_info {
379 if writer.events_tree_root_id == old_id {
380 writer.events_tree_root_id = new_id;
381 if verbose {
382 println!("Replaced root {old_id:?} with {new_id:?}");
383 }
384 } else {
385 return Err(DCBError::RootIDMismatch(old_id.0, new_id.0));
386 }
387 }
388
389 if let Some((promoted_key, promoted_page_id)) = split_info {
390 let new_internal_node = EventInternalNode {
392 keys: vec![promoted_key],
393 child_ids: vec![writer.events_tree_root_id, promoted_page_id],
394 };
395
396 let new_root_page_id = writer.alloc_page_id();
397 let new_root_page = Page::new(new_root_page_id, Node::EventInternal(new_internal_node));
398 if verbose {
399 println!(
400 "Created new internal root {:?}: {:?}",
401 new_root_page_id, new_root_page.node
402 );
403 }
404 writer.insert_dirty(new_root_page)?;
405
406 writer.events_tree_root_id = new_root_page_id;
407 }
408
409 Ok(())
410}
411
412pub fn event_tree_lookup(
413 mvcc: &Mvcc,
414 dirty: &HashMap<PageID, Page>,
415 events_tree_root_id: PageID,
416 position: Position,
417) -> DCBResult<EventRecord> {
418 let mut current_page_id: PageID = events_tree_root_id;
419 loop {
420 let page = if let Some(p) = dirty.get(¤t_page_id) {
422 p.clone()
423 } else {
424 mvcc.read_page(current_page_id)?
425 };
426 match &page.node {
427 Node::EventInternal(internal) => {
428 let idx = match internal.keys.binary_search(&position) {
430 Ok(i) => i + 1,
431 Err(i) => i,
432 };
433 if idx >= internal.child_ids.len() {
434 return Err(DCBError::DatabaseCorrupted(
435 "Child index out of bounds in event tree".to_string(),
436 ));
437 }
438 current_page_id = internal.child_ids[idx];
439 }
440 Node::EventLeaf(leaf) => {
441 return match leaf.keys.binary_search(&position) {
442 Ok(i) => {
443 let rec = materialize_event_value(mvcc, dirty, &leaf.values[i])?;
444 Ok(rec)
445 }
446 Err(_) => Err(DCBError::DatabaseCorrupted(format!(
447 "Event at position {position:?} not found",
448 ))),
449 };
450 }
451 _ => {
452 return Err(DCBError::DatabaseCorrupted(format!(
453 "Expected EventInternal or EventLeaf node in event tree, got {}",
454 page.node.type_name()
455 )));
456 }
457 }
458 }
459}
460
461pub struct EventIterator<'a> {
462 pub mvcc: &'a Mvcc,
463 pub dirty: &'a HashMap<PageID, Page>,
464 pub stack: Vec<(PageID, Option<usize>)>,
465 pub page_cache: HashMap<PageID, Page>,
466 pub start: Option<Position>, pub backwards: bool,
468}
469
470impl<'a> EventIterator<'a> {
471 pub fn new(
472 mvcc: &'a Mvcc,
473 dirty: &'a HashMap<PageID, Page>,
474 events_tree_root_id: PageID,
475 start: Option<Position>,
476 backwards: bool,
477 ) -> Self {
478 let next_position = (events_tree_root_id, None);
479 Self {
480 mvcc,
481 dirty,
482 stack: vec![next_position],
483 page_cache: HashMap::new(),
484 start,
485 backwards,
486 }
487 }
488
489 pub fn next_batch(&mut self, batch_size: u32) -> DCBResult<Vec<(Position, EventRecord)>> {
490 let mut result: Vec<(Position, EventRecord)> = Vec::with_capacity(batch_size as usize);
491 if batch_size == 0 {
492 return Ok(result);
493 }
494 if self.backwards && self.start == Some(Position(0)) {
495 return Ok(result);
496 }
497 while result.len() < batch_size as usize {
498 let Some((page_id, mut stacked_idx)) = self.stack.pop() else {
499 break; };
501
502 let mut remove_page = false;
504 let mut push_revisit: Option<(PageID, Option<usize>)> = None;
505 let mut push_child: Option<(PageID, Option<usize>)> = None; let mut emit_event: Option<(Position, EventRecord)> = None;
507
508 {
509 let page_ref: &Page = if let Some(p) = self.dirty.get(&page_id) {
511 p
512 } else if let Some(p) = self.page_cache.get(&page_id) {
513 p
514 } else {
515 let page = self.mvcc.read_page(page_id)?;
516 self.page_cache.insert(page_id, page);
517 self.page_cache
518 .get(&page_id)
519 .expect("page should be in cache")
520 };
521
522 match &page_ref.node {
523 Node::EventInternal(internal) => {
524 if stacked_idx.is_none() && !internal.keys.is_empty() {
526 stacked_idx = match &self.start {
532 Some(from) => match internal.keys.binary_search(from) {
533 Ok(i) => Some(i + 1),
534 Err(i) => Some(i),
535 },
536 None => {
537 if !self.backwards {
538 Some(0)
539 } else {
540 Some(internal.child_ids.len() - 1)
541 }
542 }
543 };
544 }
545
546 if let Some(child_ids_idx) = stacked_idx {
547 push_child = Some((internal.child_ids[child_ids_idx], None));
551 if !self.backwards {
553 if child_ids_idx + 1 < internal.child_ids.len() {
554 push_revisit = Some((page_id, Some(child_ids_idx + 1)));
557 } else {
558 remove_page = true;
560 }
562 } else if child_ids_idx > 0 {
563 push_revisit = Some((page_id, Some(child_ids_idx - 1)));
566 } else {
567 remove_page = true;
569 }
571 } else {
572 remove_page = true
574 };
575 }
576 Node::EventLeaf(leaf) => {
577 if stacked_idx.is_none() {
579 let values_len = leaf.values.len();
582
583 stacked_idx = if values_len > 0 {
584 match &self.start {
585 Some(from) => match leaf.keys.binary_search(from) {
586 Ok(i) => Some(i),
587 Err(i) => {
588 if !self.backwards {
589 Some(i)
590 } else {
591 Some(i - 1)
592 }
593 }
594 },
595 None => {
596 if !self.backwards {
597 Some(0)
598 } else {
599 Some(values_len - 1)
600 }
601 }
602 }
603 } else {
604 None
605 }
606 }
607
608 if let Some(values_idx) = stacked_idx {
609 if values_idx < leaf.values.len() {
611 let event_position = leaf.keys[values_idx];
612 let event_record = materialize_event_value(
613 self.mvcc,
614 self.dirty,
615 &leaf.values[values_idx],
616 )?;
617 emit_event = Some((event_position, event_record));
619
620 if !self.backwards {
621 if values_idx + 1 < leaf.values.len() {
622 push_revisit = Some((page_id, Some(values_idx + 1)));
624 } else {
626 remove_page = true;
628 }
630 } else if values_idx > 0 {
631 push_revisit = Some((page_id, Some(values_idx - 1)));
633 } else {
635 remove_page = true;
637 }
639 } else {
640 remove_page = true;
643 }
644 } else {
645 remove_page = true;
647 }
648 }
649 _ => {
650 return Err(DCBError::DatabaseCorrupted(format!(
651 "Expected EventInternal or EventLeaf node in event tree, got {}",
652 page_ref.node.type_name()
653 )));
654 }
655 }
656 }
657
658 if let Some(revisit) = push_revisit {
660 self.stack.push(revisit);
662 }
663 if let Some((child_id, child_start_idx)) = push_child {
664 self.stack.push((child_id, child_start_idx));
665 }
666 if let Some((event_position, event_record)) = emit_event {
667 result.push((event_position, event_record));
668 }
669 if remove_page {
670 self.page_cache.remove(&page_id);
671 }
672 }
673 Ok(result)
674 }
675}
676
677#[cfg(test)]
678mod tests {
679 use super::*;
680 use crate::node::Node;
681 use rand::random;
682 use serial_test::serial;
683 use tempfile::tempdir;
684
685 static VERBOSE: bool = false;
686
687 fn construct_db(page_size: usize) -> (tempfile::TempDir, Mvcc) {
689 let temp_dir = tempdir().unwrap();
690 let db_path = temp_dir.path().join("mvcc-test.db");
691 let db = Mvcc::new(&db_path, page_size, VERBOSE).unwrap();
692 (temp_dir, db)
693 }
694
695 #[test]
696 #[serial]
697 fn test_append_event_to_empty_leaf_root() {
698 let (_temp_dir, db) = construct_db(64);
700
701 let mut writer = db.writer().unwrap();
703
704 let position = writer.issue_position();
706
707 let record = EventRecord {
709 event_type: "UserCreated".to_string(),
710 data: vec![1, 2, 3, 4],
711 tags: vec!["users".to_string(), "creation".to_string()],
712 uuid: None,
713 };
714
715 event_tree_append(&db, &mut writer, record.clone(), position).unwrap();
717
718 let new_root_id = writer.events_tree_root_id;
720 assert!(writer.dirty.contains_key(&new_root_id));
721 let page = writer.dirty.get(&new_root_id).unwrap();
722 match &page.node {
723 Node::EventLeaf(node) => {
724 assert_eq!(vec![position], node.keys);
725 assert_eq!(
726 vec![crate::events_tree_nodes::EventValue::Inline(record.clone())],
727 node.values
728 );
729 }
730 _ => panic!("Expected EventLeaf node"),
731 }
732
733 db.commit(&mut writer).unwrap();
735
736 let (_header_page_id, header) = db.get_latest_header().unwrap();
738 let persisted_page = db.read_page(header.events_tree_root_id).unwrap();
739 match &persisted_page.node {
740 Node::EventLeaf(node) => {
741 assert_eq!(vec![position], node.keys);
742 assert_eq!(
743 vec![crate::events_tree_nodes::EventValue::Inline(record)],
744 node.values
745 );
746 }
747 _ => panic!("Expected EventLeaf node after commit"),
748 }
749 }
750
751 #[test]
752 #[serial]
753 fn test_insert_events_until_split_leaf_one_writer() {
754 let (_temp_dir, db) = construct_db(256);
756
757 let mut writer = db.writer().unwrap();
759
760 let mut has_split_leaf = false;
761 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
762
763 while !has_split_leaf {
765 let position = writer.issue_position();
767 let record = EventRecord {
768 event_type: "UserCreated".to_string(),
769 data: (0..8).map(|_| random::<u8>()).collect(),
770 tags: vec!["users".to_string(), "creation".to_string()],
771 uuid: None,
772 };
773 appended.push((position, record.clone()));
774
775 event_tree_append(&db, &mut writer, record, position).unwrap();
777
778 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
780 match &root_page.node {
781 Node::EventInternal(_) => {
782 has_split_leaf = true;
783 }
784 _ => {}
785 }
786 }
787
788 let mut copy_inserted = appended.clone();
790
791 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
793 let root_node = match &root_page.node {
794 Node::EventInternal(node) => node,
795 _ => panic!("Expected EventInternal node"),
796 };
797
798 for (i, &child_id) in root_node.child_ids.iter().enumerate() {
800 let child_page = writer.dirty.get(&child_id).unwrap();
801 assert_eq!(child_id, child_page.page_id);
802
803 let child_node = match &child_page.node {
804 Node::EventLeaf(node) => node,
805 _ => panic!("Expected EventLeaf node"),
806 };
807
808 if i > 0 {
810 assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
811 }
812
813 for (k, &key) in child_node.keys.iter().enumerate() {
815 let record = &child_node.values[k];
816 let (appended_position, appended_record) = copy_inserted.remove(0);
817 assert_eq!(appended_position, key);
818 assert_eq!(appended_record, record.clone());
819 }
820 }
821 }
822
823 #[test]
824 #[serial]
825 fn test_insert_events_until_split_leaf_many_writers() {
826 let (_temp_dir, db) = construct_db(256);
828
829 let mut has_split_leaf = false;
830 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
831
832 while !has_split_leaf {
834 let mut writer = db.writer().unwrap();
836
837 let position = writer.issue_position();
839 let record = EventRecord {
840 event_type: "UserCreated".to_string(),
841 data: (0..8).map(|_| random::<u8>()).collect(),
842 tags: vec!["users".to_string(), "creation".to_string()],
843 uuid: None,
844 };
845 appended.push((position, record.clone()));
846
847 event_tree_append(&db, &mut writer, record, position).unwrap();
849
850 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
852 match &root_page.node {
853 Node::EventInternal(_) => {
854 has_split_leaf = true;
855 }
856 _ => {}
857 }
858
859 db.commit(&mut writer).unwrap();
860 }
861
862 let mut copy_inserted = appended.clone();
864
865 let writer = db.writer().unwrap();
867
868 let root_page = db.read_page(writer.events_tree_root_id).unwrap();
870 let root_node = match &root_page.node {
871 Node::EventInternal(node) => node,
872 _ => panic!("Expected EventInternal node"),
873 };
874
875 for (i, &child_id) in root_node.child_ids.iter().enumerate() {
877 let child_page = db.read_page(child_id).unwrap();
878 assert_eq!(child_id, child_page.page_id);
879
880 let child_node = match &child_page.node {
881 Node::EventLeaf(node) => node,
882 _ => panic!("Expected EventLeaf node"),
883 };
884
885 if i > 0 {
887 assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
888 }
889
890 for (k, &key) in child_node.keys.iter().enumerate() {
892 let record = &child_node.values[k];
893 let (appended_position, appended_record) = copy_inserted.remove(0);
894 assert_eq!(appended_position, key);
895 assert_eq!(appended_record, record.clone());
896 }
897 }
898 }
899
900 #[test]
901 #[serial]
902 fn test_insert_events_until_split_internal_one_writer() {
903 let (_temp_dir, db) = construct_db(256);
905
906 let mut writer = db.writer().unwrap();
908
909 let mut has_split_internal = false;
910 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
911
912 while !has_split_internal {
914 let position = writer.issue_position();
916 let record = EventRecord {
917 event_type: "UserCreated".to_string(),
918 data: (0..8).map(|_| random::<u8>()).collect(),
919 tags: vec!["users".to_string(), "creation".to_string()],
920 uuid: None,
921 };
922 appended.push((position, record.clone()));
923
924 event_tree_append(&db, &mut writer, record, position).unwrap();
926
927 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
929 match &root_page.node {
930 Node::EventInternal(root_node) => {
931 if !root_node.child_ids.is_empty() {
933 let child_id = root_node.child_ids[0];
934 if let Some(child_page) = writer.dirty.get(&child_id) {
935 match &child_page.node {
936 Node::EventInternal(_) => {
937 has_split_internal = true;
938 }
939 _ => {}
940 }
941 }
942 }
943 }
944 _ => {}
945 }
946 }
947
948 let mut copy_inserted = appended.clone();
950
951 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
953 let root_node = match &root_page.node {
954 Node::EventInternal(node) => node,
955 _ => panic!("Expected EventInternal node"),
956 };
957
958 for &child_id in root_node.child_ids.iter() {
960 let child_page = writer.dirty.get(&child_id).unwrap();
961 assert_eq!(child_id, child_page.page_id);
962
963 let child_node = match &child_page.node {
964 Node::EventInternal(node) => node,
965 _ => panic!("Expected EventInternal node"),
966 };
967
968 for &grand_child_id in child_node.child_ids.iter() {
969 let grand_child_page = writer.dirty.get(&grand_child_id).unwrap();
970 assert_eq!(grand_child_id, grand_child_page.page_id);
971
972 let grand_child_node = match &grand_child_page.node {
973 Node::EventLeaf(node) => node,
974 _ => panic!("Expected EventLeaf node"),
975 };
976
977 for (k, &key) in grand_child_node.keys.iter().enumerate() {
979 let record = &grand_child_node.values[k];
980 let (appended_position, appended_record) = copy_inserted.remove(0);
981 assert_eq!(appended_position, key);
982 assert_eq!(appended_record, record.clone());
983 }
984 }
985 }
986 }
987
988 #[test]
989 #[serial]
990 fn test_insert_events_until_split_internal_many_writers() {
991 let (_temp_dir, db) = construct_db(512);
993
994 let mut has_split_internal = false;
995 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
996
997 while !has_split_internal {
999 let mut writer = db.writer().unwrap();
1001
1002 let position = writer.issue_position();
1004 let record = EventRecord {
1005 event_type: "UserCreated".to_string(),
1006 data: (0..8).map(|_| random::<u8>()).collect(),
1007 tags: vec!["users".to_string(), "creation".to_string()],
1008 uuid: None,
1009 };
1010 appended.push((position, record.clone()));
1011
1012 event_tree_append(&db, &mut writer, record, position).unwrap();
1014
1015 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1017 match &root_page.node {
1018 Node::EventInternal(root_node) => {
1019 if !root_node.child_ids.is_empty() {
1021 let child_id = root_node.child_ids[0];
1022 if let Some(child_page) = writer.dirty.get(&child_id) {
1023 match &child_page.node {
1024 Node::EventInternal(_) => {
1025 has_split_internal = true;
1026 }
1027 _ => {}
1028 }
1029 }
1030 }
1031 }
1032 _ => {}
1033 }
1034 db.commit(&mut writer).unwrap();
1035 }
1036
1037 let mut copy_inserted = appended.clone();
1039
1040 let writer = db.writer().unwrap();
1042
1043 let root_page = db.read_page(writer.events_tree_root_id).unwrap();
1045 let root_node = match &root_page.node {
1046 Node::EventInternal(node) => node,
1047 _ => panic!("Expected EventInternal node"),
1048 };
1049
1050 for &child_id in root_node.child_ids.iter() {
1052 let child_page = db.read_page(child_id).unwrap();
1053 assert_eq!(child_id, child_page.page_id);
1054
1055 let child_node = match &child_page.node {
1056 Node::EventInternal(node) => node,
1057 _ => panic!("Expected EventInternal node"),
1058 };
1059
1060 for &grand_child_id in child_node.child_ids.iter() {
1061 let grand_child_page = db.read_page(grand_child_id).unwrap();
1062 assert_eq!(grand_child_id, grand_child_page.page_id);
1063
1064 let grand_child_node = match &grand_child_page.node {
1065 Node::EventLeaf(node) => node,
1066 _ => panic!("Expected EventLeaf node"),
1067 };
1068
1069 for (k, &key) in grand_child_node.keys.iter().enumerate() {
1071 let record = &grand_child_node.values[k];
1072 let (appended_position, appended_record) = copy_inserted.remove(0);
1073 assert_eq!(appended_position, key);
1075 assert_eq!(appended_record, record.clone());
1076 }
1077 }
1078 }
1079 assert_eq!(0, copy_inserted.len());
1080 }
1081
1082 #[test]
1083 #[serial]
1084 fn test_read_events_all() {
1085 let (_temp_dir, db) = construct_db(512);
1087
1088 let mut has_split_internal = false;
1089 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1090
1091 while !has_split_internal {
1093 let mut writer = db.writer().unwrap();
1095
1096 let position = writer.issue_position();
1098 let record = EventRecord {
1099 event_type: "UserCreated".to_string(),
1100 data: (0..8).map(|_| random::<u8>()).collect(),
1101 tags: vec!["users".to_string(), "creation".to_string()],
1102 uuid: None,
1103 };
1104 appended.push((position, record.clone()));
1105
1106 event_tree_append(&db, &mut writer, record, position).unwrap();
1108
1109 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1111 match &root_page.node {
1112 Node::EventInternal(root_node) => {
1113 if !root_node.child_ids.is_empty() {
1115 let child_id = root_node.child_ids[0];
1116 if let Some(child_page) = writer.dirty.get(&child_id) {
1117 match &child_page.node {
1118 Node::EventInternal(_) => {
1119 has_split_internal = true;
1120 }
1121 _ => {}
1122 }
1123 }
1124 }
1125 }
1126 _ => {}
1127 }
1128 db.commit(&mut writer).unwrap();
1129 }
1130
1131 let copy_inserted = appended.clone();
1133
1134 let reader = db.reader().unwrap();
1136 let events_tree_root_id = reader.events_tree_root_id;
1137 let reader_tsn = reader.tsn;
1138
1139 let dirty = HashMap::new();
1140 let mut events_iterator = EventIterator::new(&db, &dirty, events_tree_root_id, None, false);
1141
1142 {
1144 assert!(
1146 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1147 "TSN should remain registered until reader is dropped"
1148 );
1149 }
1150
1151 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1153 loop {
1154 let batch = events_iterator.next_batch(3).unwrap();
1155 if batch.is_empty() {
1156 break;
1157 }
1158 scanned.extend(batch);
1159
1160 assert!(
1163 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1164 "TSN should remain registered until reader is dropped"
1165 );
1166 }
1167
1168 assert_eq!(copy_inserted.len(), scanned.len());
1169 for (i, expected) in copy_inserted.iter().enumerate() {
1170 assert_eq!(expected.0, scanned[i].0);
1171 assert_eq!(expected.1, scanned[i].1);
1172 }
1173
1174 let dirty = HashMap::new();
1176 for (pos, expected_rec) in copy_inserted.iter() {
1177 let found = event_tree_lookup(&db, &dirty, events_tree_root_id, *pos).unwrap();
1178 assert_eq!(expected_rec, &found);
1179 }
1180
1181 assert!(
1183 events_iterator.page_cache.is_empty(),
1184 "EventIterator page_cache should be empty after full scan"
1185 );
1186
1187 {
1189 assert!(
1191 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1192 "TSN should remain registered until reader is dropped"
1193 );
1194 }
1195
1196 drop(reader);
1198 {
1199 assert!(
1201 db.reader_tsns.iter().all(|r| *r.value() != reader_tsn),
1202 "TSN should be removed after reader is dropped"
1203 );
1204 assert_eq!(0, db.reader_tsns.len());
1205 }
1206 }
1207
1208 #[test]
1209 #[serial]
1210 fn test_read_events_from_forwards() {
1211 let (_temp_dir, db) = construct_db(128);
1213
1214 let mut has_split_internal = false;
1215 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1216
1217 while !has_split_internal {
1219 let mut writer = db.writer().unwrap();
1221
1222 let position = writer.issue_position();
1224 let record = EventRecord {
1225 event_type: "UserCreated".to_string(),
1226 data: (0..8).map(|_| random::<u8>()).collect(),
1227 tags: vec!["users".to_string(), "creation".to_string()],
1228 uuid: None,
1229 };
1230 appended.push((position, record.clone()));
1231
1232 event_tree_append(&db, &mut writer, record, position).unwrap();
1234
1235 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1237 match &root_page.node {
1238 Node::EventInternal(root_node) => {
1239 if !root_node.child_ids.is_empty() {
1241 let child_id = root_node.child_ids[0];
1242 if let Some(child_page) = writer.dirty.get(&child_id) {
1243 match &child_page.node {
1244 Node::EventInternal(_) => {
1245 has_split_internal = true;
1246 }
1247 _ => {}
1248 }
1249 }
1250 }
1251 }
1252 _ => {}
1253 }
1254 db.commit(&mut writer).unwrap();
1255 }
1256
1257 let appended_len = appended.len();
1258 for i in 0..appended_len + 2 {
1262 let from = Position(i as u64);
1263
1264 let reader = db.reader().unwrap();
1266 let events_tree_root_id = reader.events_tree_root_id;
1267 let reader_tsn = reader.tsn;
1268 let dirty = HashMap::new();
1269 let mut events_iterator =
1270 EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), false);
1271
1272 {
1274 assert!(
1276 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1277 "TSN should remain registered until reader is dropped"
1278 );
1279 }
1280
1281 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1283 loop {
1284 let batch = events_iterator.next_batch(3).unwrap();
1285 if batch.is_empty() {
1286 break;
1287 }
1288 scanned.extend(batch);
1289
1290 assert!(
1293 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1294 "TSN should remain registered until reader is dropped"
1295 );
1296 }
1297
1298 let num_to_skip = i.max(1) - 1;
1300 let expected: Vec<(Position, EventRecord)> =
1301 appended.clone().into_iter().skip(num_to_skip).collect();
1302 assert_eq!(expected.len(), scanned.len());
1303 for (i, exp) in expected.iter().enumerate() {
1305 assert_eq!(exp.0, scanned[i].0);
1306 assert_eq!(exp.1, scanned[i].1);
1307 }
1308
1309 assert!(
1311 events_iterator.page_cache.is_empty(),
1312 "EventIterator page_cache should be empty after filtered scan"
1313 );
1314
1315 drop(reader);
1317 {
1318 assert!(
1320 db.reader_tsns.iter().all(|r| *r.value() != reader_tsn),
1321 "TSN should be removed after reader is dropped"
1322 );
1323 }
1324 }
1325 }
1326
1327 #[test]
1328 #[serial]
1329 fn test_read_events_from_backwards() {
1330 let (_temp_dir, db) = construct_db(128);
1332
1333 let mut has_split_internal = false;
1334 let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1335
1336 while !has_split_internal {
1338 let mut writer = db.writer().unwrap();
1340
1341 let position = writer.issue_position();
1343 let record = EventRecord {
1344 event_type: "UserCreated".to_string(),
1345 data: (0..8).map(|_| random::<u8>()).collect(),
1346 tags: vec!["users".to_string(), "creation".to_string()],
1347 uuid: None,
1348 };
1349 appended.push((position, record.clone()));
1350
1351 event_tree_append(&db, &mut writer, record, position).unwrap();
1353
1354 let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1356 match &root_page.node {
1357 Node::EventInternal(root_node) => {
1358 if !root_node.child_ids.is_empty() {
1360 let child_id = root_node.child_ids[0];
1361 if let Some(child_page) = writer.dirty.get(&child_id) {
1362 match &child_page.node {
1363 Node::EventInternal(_) => {
1364 has_split_internal = true;
1365 }
1366 _ => {}
1367 }
1368 }
1369 }
1370 }
1371 _ => {}
1372 }
1373 db.commit(&mut writer).unwrap();
1374 }
1375
1376 let appended_len = appended.len();
1377 for i in 0..appended_len + 2 {
1381 let from = Position(i as u64);
1382
1383 let reader = db.reader().unwrap();
1385 let events_tree_root_id = reader.events_tree_root_id;
1386 let reader_tsn = reader.tsn;
1387 let dirty = HashMap::new();
1388 let mut events_iterator =
1389 EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), true);
1390
1391 {
1393 assert!(
1395 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1396 "TSN should remain registered until reader is dropped"
1397 );
1398 }
1399
1400 let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1402 loop {
1403 let batch = events_iterator.next_batch(3).unwrap();
1404 if batch.is_empty() {
1405 break;
1406 }
1407 scanned.extend(batch);
1408
1409 assert!(
1412 db.reader_tsns.iter().any(|r| *r.value() == reader_tsn),
1413 "TSN should remain registered until reader is dropped"
1414 );
1415 }
1416
1417 let mut expected: Vec<(Position, EventRecord)> = appended.clone();
1419 expected.truncate(i);
1420 expected.reverse();
1421 assert_eq!(expected.len(), scanned.len());
1422 for (i, exp) in expected.iter().enumerate() {
1424 assert_eq!(exp.0, scanned[i].0);
1425 assert_eq!(exp.1, scanned[i].1);
1426 }
1427
1428 assert!(
1430 events_iterator.page_cache.is_empty(),
1431 "EventIterator page_cache should be empty after filtered scan"
1432 );
1433
1434 drop(reader);
1436 {
1437 assert!(
1439 db.reader_tsns.iter().all(|r| *r.value() != reader_tsn),
1440 "TSN should be removed after reader is dropped"
1441 );
1442 }
1443 }
1444 }
1445
1446 #[test]
1447 #[serial]
1448 fn test_large_event_data_exact_page_size() {
1449 let (_tmp, db) = construct_db(512);
1450 let mut writer = db.writer().unwrap();
1452 let pos = writer.issue_position();
1453 let data = vec![0xAB; 512];
1454 let event = EventRecord {
1455 event_type: "Big".into(),
1456 data: data.clone(),
1457 tags: vec![],
1458 uuid: None,
1459 };
1460 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1461 db.commit(&mut writer).unwrap();
1462
1463 let reader = db.reader().unwrap();
1465 let dirty = HashMap::new();
1466 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1467 assert_eq!(event, got);
1468
1469 let (_hdr_id, header) = db.get_latest_header().unwrap();
1471 let root = db.read_page(header.events_tree_root_id).unwrap();
1472 match root.node {
1473 Node::EventInternal(internal) => {
1474 let leaf_id = *internal.child_ids.last().unwrap();
1476 let leaf_page = db.read_page(leaf_id).unwrap();
1477 match leaf_page.node {
1478 Node::EventLeaf(leaf) => match &leaf.values[0] {
1479 EventValue::Overflow { data_len, .. } => {
1480 assert_eq!(*data_len as usize, data.len())
1481 }
1482 _ => panic!("Expected Overflow for large event"),
1483 },
1484 _ => panic!("Expected EventLeaf child"),
1485 }
1486 }
1487 Node::EventLeaf(leaf) => match &leaf.values[0] {
1488 EventValue::Overflow { data_len, .. } => assert_eq!(*data_len as usize, data.len()),
1489 _ => panic!("Expected Overflow for large event"),
1490 },
1491 _ => panic!("Unexpected root node type"),
1492 }
1493 }
1494
1495 #[test]
1496 #[serial]
1497 fn test_large_event_data_four_times_page_size() {
1498 let (_tmp, db) = construct_db(512);
1499 let mut writer = db.writer().unwrap();
1500 let pos = writer.issue_position();
1501 let data = vec![0xCD; 512 * 4];
1502 let event = EventRecord {
1503 event_type: "Bigger".into(),
1504 data: data.clone(),
1505 tags: vec![],
1506 uuid: None,
1507 };
1508 event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1509 db.commit(&mut writer).unwrap();
1510
1511 let reader = db.reader().unwrap();
1513 let dirty = HashMap::new();
1514 let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1515 assert_eq!(event, got);
1516
1517 let (_hdr_id, header) = db.get_latest_header().unwrap();
1519 let root = db.read_page(header.events_tree_root_id).unwrap();
1520 let check_leaf = |leaf: &EventLeafNode| match &leaf.values[0] {
1521 EventValue::Overflow { data_len, .. } => assert_eq!(*data_len as usize, data.len()),
1522 _ => panic!("Expected Overflow for very large event"),
1523 };
1524 match root.node {
1525 Node::EventInternal(internal) => {
1526 let leaf_id = *internal.child_ids.last().unwrap();
1527 let leaf_page = db.read_page(leaf_id).unwrap();
1528 match leaf_page.node {
1529 Node::EventLeaf(leaf) => check_leaf(&leaf),
1530 _ => panic!("Expected leaf"),
1531 }
1532 }
1533 Node::EventLeaf(leaf) => check_leaf(&leaf),
1534 _ => panic!("Unexpected root node type"),
1535 }
1536 }
1537
1538 }