umadb_core/
events_tree.rs

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
12// Helpers for storing large event data across overflow pages
13fn write_overflow_chain(mvcc: &Mvcc, writer: &mut Writer, data: &[u8]) -> DCBResult<PageID> {
14    // Maximum payload per overflow page: page_size - header - next pointer (8 bytes)
15    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    // Split data into chunks from end to start so we can set next ids easily
22    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        // Store an empty chunk page to indicate zero-length data
31        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        // Prefer the dirty (unflushed) page if present; otherwise read from disk
62        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
112/// Append an event to the root event leaf page.
113///
114/// This function obtains a mutable reference to a dirty copy of the root event
115/// leaf page (using copy-on-write if necessary) and appends the provided
116/// Position to the keys and the EventRecord to the values.
117pub 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    // Get the current root page id for the event tree
129    let mut current_page_id: PageID = writer.events_tree_root_id;
130
131    // Traverse the tree to find a leaf node
132    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) = &current_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    // Decide inline vs overflow based on data length before mut-borrowing the page
158    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    // Make the leaf page dirty
172    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    // We may need to pop the last key/value for splitting; hold it after we drop the borrow
182    let mut popped: Option<(Position, EventValue)> = None;
183
184    // Get a mutable leaf node and append the data
185    {
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                // Check if the leaf needs splitting by estimating the serialized size
193                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    // Prepare for split propagation
221    let mut split_info: Option<(Position, PageID)> = None;
222
223    if let Some((last_key, mut last_value)) = popped {
224        // Build new leaf node; convert to overflow if needed to fit
225        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            // serialized_size = new_leaf_page.calc_serialized_size();
249        }
250        // if serialized_size > mvcc.page_size {
251        //     return Err(DCBError::DatabaseCorrupted(format!(
252        //         "Event too large even after overflow conversion (size: {serialized_size}, max: {})",
253        //         mvcc.page_size
254        //     )));
255        // }
256        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    // Propagate splits and replacements up the stack
270    let mut current_replacement_info = replacement_info;
271    while let Some(parent_page_id) = stack.pop() {
272        // Make the internal page dirty
273        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        // Get a mutable internal node....
282        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                // Add the promoted key and page ID
304                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        // Check if the internal page needs splitting
320
321        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                // Split the internal node
327                // Ensure we have at least 3 keys and 4 child IDs before splitting
328                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                // Move the right-most key to a new node. Promote the next right-most key.
335                let (promoted_key, new_keys, new_child_ids) = dirty_internal_node.split_off()?;
336
337                // Ensure old node maintain the B-tree invariant: n keys should have n+1 child pointers
338                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                // Ensure the new node also maintains the invariant
349                assert_eq!(
350                    new_internal_node.keys.len() + 1,
351                    new_internal_node.child_ids.len()
352                );
353
354                // Create a new internal page.
355                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        // Create a new root
391        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        // Prefer the dirty (unflushed) page if present; otherwise read from disk
421        let page = if let Some(p) = dirty.get(&current_page_id) {
422            p.clone()
423        } else {
424            mvcc.read_page(current_page_id)?
425        };
426        match &page.node {
427            Node::EventInternal(internal) => {
428                // Choose child based on upper bound of position in separator keys
429                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>, // inclusive position, better for binary search
467    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; // traversal finished
500            };
501
502            // Compute actions under a scoped immutable borrow, then mutate cache/stack afterwards.
503            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; // (child_id, stacked_keys_idx)
506            let mut emit_event: Option<(Position, EventRecord)> = None;
507
508            {
509                // Obtain the current page (from dirty, or page cache, or deserialize).
510                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                        // println!("Visit internal {page_id:?}");
525                        if stacked_idx.is_none() && !internal.keys.is_empty() {
526                            // println!(" - first visit");
527                            // println!(" - keys: {:?}", internal.keys.clone());
528                            // println!(" - child_ids: {:?}", internal.child_ids.clone());
529                            // println!(" - from: {:?}", self.from);
530
531                            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                            // println!(" - child ids index: {} / {}", child_ids_idx + 1, internal.child_ids.len());
548                            // println!(" - will visit child: {:?}", internal.child_ids[child_ids_idx]);
549                            // Push the chosen child.
550                            push_child = Some((internal.child_ids[child_ids_idx], None));
551                            // Do or don't revisit this internal node?
552                            if !self.backwards {
553                                if child_ids_idx + 1 < internal.child_ids.len() {
554                                    // Will revisit this internal node.
555                                    // println!(" - will revisit");
556                                    push_revisit = Some((page_id, Some(child_ids_idx + 1)));
557                                } else {
558                                    // Don't revisit this internal node.
559                                    remove_page = true;
560                                    // println!(" - will remove");
561                                }
562                            } else if child_ids_idx > 0 {
563                                // Will revisit this internal node.
564                                // println!(" - will revisit");
565                                push_revisit = Some((page_id, Some(child_ids_idx - 1)));
566                            } else {
567                                // Don't revisit this internal node.
568                                remove_page = true;
569                                // println!(" - will remove");
570                            }
571                        } else {
572                            // TODO: Clarify if this is always because internal node is empty?
573                            remove_page = true
574                        };
575                    }
576                    Node::EventLeaf(leaf) => {
577                        // println!("Visit leaf {page_id:?}");
578                        if stacked_idx.is_none() {
579                            // println!(" - first visit");
580                            // println!(" - keys: {:?}", leaf.keys.clone());
581                            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                            // println!(" - values index: {} / {}", values_idx + 1, leaf.values.len());
610                            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                                // println!(" - emit event position: {:?}", event_position.clone());
618                                emit_event = Some((event_position, event_record));
619
620                                if !self.backwards {
621                                    if values_idx + 1 < leaf.values.len() {
622                                        // Revisit this leaf.
623                                        push_revisit = Some((page_id, Some(values_idx + 1)));
624                                        // println!(" - not last value, will revisit");
625                                    } else {
626                                        // The last value.
627                                        remove_page = true;
628                                        // println!(" - last value, will remove");
629                                    }
630                                } else if values_idx > 0 {
631                                    // Revisit this leaf.
632                                    push_revisit = Some((page_id, Some(values_idx - 1)));
633                                    // println!(" - not last value, will revisit");
634                                } else {
635                                    // The last value.
636                                    remove_page = true;
637                                    // println!(" - last value, will remove");
638                                }
639                            } else {
640                                // No key greater or equal to 'from' in this leaf
641                                // println!(" - value index out of range, why wasn't this removed?");
642                                remove_page = true;
643                            }
644                        } else {
645                            // No leaf values.
646                            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            // Mutations after the borrow has ended
659            if let Some(revisit) = push_revisit {
660                // Revisit must be pushed first so that the child is processed next (LIFO)
661                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    // Helper function to create a test database with a specified page size
688    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        // Setup a temporary database
699        let (_temp_dir, db) = construct_db(64);
700
701        // Start a writer
702        let mut writer = db.writer().unwrap();
703
704        // Issue a new position
705        let position = writer.issue_position();
706
707        // Create an event record
708        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        // Call append_event
716        event_tree_append(&db, &mut writer, record.clone(), position).unwrap();
717
718        // Verify that the dirty root page contains the appended key/value
719        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        // Commit the writer and verify persistence
734        db.commit(&mut writer).unwrap();
735
736        // Read back the latest header and the persisted root event leaf page
737        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        // Setup a temporary database
755        let (_temp_dir, db) = construct_db(256);
756
757        // Start a writer
758        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        // Insert events until we split a leaf
764        while !has_split_leaf {
765            // Issue a new position and create a record
766            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            // Append the event
776            event_tree_append(&db, &mut writer, record, position).unwrap();
777
778            // Check if we've split the leaf
779            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        // Check keys and values of all pages
789        let mut copy_inserted = appended.clone();
790
791        // Get the root node
792        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        // Check each child of the root
799        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            // Check that the keys are properly ordered
809            if i > 0 {
810                assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
811            }
812
813            // Check each key and value in the child
814            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        // Setup a temporary database
827        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        // Insert events until we split a leaf
833        while !has_split_leaf {
834            // Start a writer
835            let mut writer = db.writer().unwrap();
836
837            // Issue a new position and create a record
838            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            // Append the event
848            event_tree_append(&db, &mut writer, record, position).unwrap();
849
850            // Check if we've split the leaf
851            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        // Check keys and values of all pages
863        let mut copy_inserted = appended.clone();
864
865        // Start a writer
866        let writer = db.writer().unwrap();
867
868        // Get the root node
869        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        // Check each child of the root
876        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            // Check that the keys are properly ordered
886            if i > 0 {
887                assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
888            }
889
890            // Check each key and value in the child
891            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        // Setup a temporary database
904        let (_temp_dir, db) = construct_db(256);
905
906        // Start a writer
907        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        // Insert events until we split a leaf
913        while !has_split_internal {
914            // Issue a new position and create a record
915            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            // Append the event
925            event_tree_append(&db, &mut writer, record, position).unwrap();
926
927            // Check if we've split an internal node
928            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
929            match &root_page.node {
930                Node::EventInternal(root_node) => {
931                    // Check if the first child is an internal node
932                    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        // Check keys and values of all pages
949        let mut copy_inserted = appended.clone();
950
951        // Get the root node
952        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        // Check each child of the root
959        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                // Check each key and value in the child
978                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        // Setup a temporary database
992        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        // Insert events until we split a root internal node
998        while !has_split_internal {
999            // Start a writer
1000            let mut writer = db.writer().unwrap();
1001
1002            // Issue a new position and create a record
1003            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            // Append the event
1013            event_tree_append(&db, &mut writer, record, position).unwrap();
1014
1015            // Check if the root is an internal node
1016            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1017            match &root_page.node {
1018                Node::EventInternal(root_node) => {
1019                    // Check if the first child is an internal node
1020                    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        // Check keys and values of all pages
1038        let mut copy_inserted = appended.clone();
1039
1040        // Start a writer
1041        let writer = db.writer().unwrap();
1042
1043        // Get the root node
1044        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        // Check each child of the root
1051        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                // Check each key and value in the child
1070                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                    // println!("Checking appended event: {appended_position:?} {appended_record:?}");
1074                    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        // Setup a temporary database
1086        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        // Insert events until we split a root internal node
1092        while !has_split_internal {
1093            // Start a writer
1094            let mut writer = db.writer().unwrap();
1095
1096            // Issue a new position and create a record
1097            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            // Append the event
1107            event_tree_append(&db, &mut writer, record, position).unwrap();
1108
1109            // Check if the root is an internal node
1110            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1111            match &root_page.node {
1112                Node::EventInternal(root_node) => {
1113                    // Check if the first child is an internal node
1114                    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        // Check keys and values of all pages
1132        let copy_inserted = appended.clone();
1133
1134        // Start a reader
1135        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        // Ensure the reader's tsn is registered while the iterator is alive
1143        {
1144            // DashMap: direct access without intermediate variable
1145            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        // Progressively iterate over events using batches
1152        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            // The reader should remain registered throughout iteration
1161            // DashMap: direct access without intermediate variable
1162            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        // Additionally, validate lookup_event for each appended position using the existing reader in the iterator
1175        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        // Ensure we did not accumulate pages in the iterator cache
1182        assert!(
1183            events_iterator.page_cache.is_empty(),
1184            "EventIterator page_cache should be empty after full scan"
1185        );
1186
1187        // While iterator is still alive, the reader should still be registered
1188        {
1189            // DashMap: direct access without intermediate variable
1190            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 the reader and ensure the reader tsn is removed
1197        drop(reader);
1198        {
1199            // DashMap: direct access without intermediate variable
1200            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        // Setup a temporary database
1212        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        // Insert events until we split a root internal node
1218        while !has_split_internal {
1219            // Start a writer
1220            let mut writer = db.writer().unwrap();
1221
1222            // Issue a new position and create a record
1223            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            // Append the event
1233            event_tree_append(&db, &mut writer, record, position).unwrap();
1234
1235            // Check if the root is an internal node
1236            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1237            match &root_page.node {
1238                Node::EventInternal(root_node) => {
1239                    // Check if the first child is an internal node
1240                    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        // println!("Appended number: {}", appended_len);
1259
1260        // Iterate through various 'from' positions.
1261        for i in 0..appended_len + 2 {
1262            let from = Position(i as u64);
1263
1264            // Start a reader
1265            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            // Ensure the reader's tsn is registered while the iterator is alive
1273            {
1274                // DashMap: direct access without intermediate variable
1275                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            // Progressively iterate over events using batches
1282            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                // The reader should remain registered throughout iteration
1291                // DashMap: direct access without intermediate variable
1292                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            // Expected are strictly from the chosen 'from' position
1299            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            // println!("Get expected number: {}", expected.len());
1304            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            // Ensure we did not accumulate pages in the iterator cache
1310            assert!(
1311                events_iterator.page_cache.is_empty(),
1312                "EventIterator page_cache should be empty after filtered scan"
1313            );
1314
1315            // Drop the reader and ensure the reader tsn is removed
1316            drop(reader);
1317            {
1318                // DashMap: direct access without intermediate variable
1319                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        // Setup a temporary database
1331        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        // Insert events until we split a root internal node
1337        while !has_split_internal {
1338            // Start a writer
1339            let mut writer = db.writer().unwrap();
1340
1341            // Issue a new position and create a record
1342            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            // Append the event
1352            event_tree_append(&db, &mut writer, record, position).unwrap();
1353
1354            // Check if the root is an internal node
1355            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1356            match &root_page.node {
1357                Node::EventInternal(root_node) => {
1358                    // Check if the first child is an internal node
1359                    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        // println!("Appended number: {}", appended_len);
1378
1379        // Iterate through various 'from' positions.
1380        for i in 0..appended_len + 2 {
1381            let from = Position(i as u64);
1382
1383            // Start a reader
1384            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            // Ensure the reader's tsn is registered while the iterator is alive
1392            {
1393                // DashMap: direct access without intermediate variable
1394                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            // Progressively iterate over events using batches
1401            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                // The reader should remain registered throughout iteration
1410                // DashMap: direct access without intermediate variable
1411                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            // Expected are strictly from the chosen 'from' position
1418            let mut expected: Vec<(Position, EventRecord)> = appended.clone();
1419            expected.truncate(i);
1420            expected.reverse();
1421            assert_eq!(expected.len(), scanned.len());
1422            // println!("Get expected number: {}", expected.len());
1423            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            // Ensure we did not accumulate pages in the iterator cache
1429            assert!(
1430                events_iterator.page_cache.is_empty(),
1431                "EventIterator page_cache should be empty after filtered scan"
1432            );
1433
1434            // Drop the reader and ensure the reader tsn is removed
1435            drop(reader);
1436            {
1437                // DashMap: direct access without intermediate variable
1438                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        // Append one large event with data exactly equal to page size
1451        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        // Lookup should return identical payload
1464        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        // Ensure an overflow page is used for storage
1470        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                // Our key should be in the last child
1475                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        // Lookup
1512        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        // Ensure overflow in leaf
1518        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    // #[test]
1539    // fn benchmark_append_and_lookup_varied_sizes() {
1540    //     // Benchmark-like test; prints durations for different sizes. Run with:
1541    //     // cargo test --lib mvcc_event_tree::tests::benchmark_append_and_lookup_varied_sizes -- --nocapture
1542    //     let sizes: [usize; 7] = [1, 10, 100, 1_000, 5_000, 10_000, 50_000];
1543    //     for &size in &sizes {
1544    //         let (_tmp, db) = construct_db(4096);
1545    //
1546    //         // Append phase
1547    //         let mut writer = db.writer().unwrap();
1548    //         let mut positions: Vec<Position> = Vec::with_capacity(size);
1549    //         let start_append = std::time::Instant::now();
1550    //         for n in 0..(size as u64) {
1551    //             let pos = writer.issue_position();
1552    //             let event = EventRecord {
1553    //                 event_type: "E".to_string(),
1554    //                 data: Vec::new(),
1555    //                 tags: Vec::new(),
1556    //             };
1557    //             std::hint::black_box(n);
1558    //             std::hint::black_box(&event);
1559    //             std::hint::black_box(pos);
1560    //             event_tree_append(&db, &mut writer, event, pos).unwrap();
1561    //             positions.push(pos);
1562    //         }
1563    //         let append_elapsed = start_append.elapsed();
1564    //         let start_commit = std::time::Instant::now();
1565    //         db.commit(&mut writer).unwrap();
1566    //         let commit_elapsed = start_commit.elapsed();
1567    //
1568    //         // Lookup phase
1569    //         let reader = db.reader().unwrap();
1570    //         let start_lookup = std::time::Instant::now();
1571    //         let dirty = HashMap::new();
1572    //         for &pos in &positions {
1573    //             let rec = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1574    //             std::hint::black_box(&rec);
1575    //         }
1576    //         let lookup_elapsed = start_lookup.elapsed();
1577    //
1578    //         let append_avg_us = (append_elapsed.as_secs_f64() * 1_000_000.0) / (size as f64);
1579    //         let commit_avg_us = commit_elapsed.as_secs_f64() * 1_000_000.0;
1580    //         let lookup_avg_us = (lookup_elapsed.as_secs_f64() * 1_000_000.0) / (size as f64);
1581    //
1582    //         println!(
1583    //             "mvcc_event_tree benchmark: size={}, append_us_per_call={:.3}, commit_us={:.3}, lookup_us_per_call={:.3}",
1584    //             size, append_avg_us, commit_avg_us, lookup_avg_us
1585    //         );
1586    //     }
1587    // }
1588}