Skip to main content

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    deserialize_metadata, metadata_serialized_size, serialize_metadata_into,
6    validate_event_record_for_append,
7};
8use crate::mvcc::{Mvcc, Writer};
9use crate::node::Node;
10use crate::page::{PAGE_HEADER_SIZE, Page};
11use std::collections::HashMap;
12use std::sync::Arc;
13use umadb_dcb::{DcbError, DcbResult};
14
15// Helpers for storing large event data across overflow pages
16fn write_overflow_chain(mvcc: &Mvcc, writer: &mut Writer, data: &[u8]) -> DcbResult<PageID> {
17    // Maximum payload per overflow page: page_size - header - next pointer (8 bytes)
18    let payload_cap = mvcc.page_size.saturating_sub(PAGE_HEADER_SIZE + 8);
19    if payload_cap == 0 {
20        return Err(DcbError::DatabaseCorrupted(
21            "Page size too small to store overflow data".to_string(),
22        ));
23    }
24    // Split data into chunks from end to start so we can set next ids easily
25    let chunk_count = data.len().div_ceil(payload_cap);
26    let mut chunks: Vec<&[u8]> = Vec::with_capacity(chunk_count);
27    let mut i = 0;
28    while i < data.len() {
29        let end = (i + payload_cap).min(data.len());
30        chunks.push(&data[i..end]);
31        i = end;
32    }
33    if chunks.is_empty() {
34        // Store an empty chunk page to indicate zero-length data
35        let page_id = writer.alloc_page_id();
36        let node = EventOverflowNode {
37            next: PageID(0),
38            data: Vec::new(),
39        };
40        let page = Page::new(page_id, Node::EventOverflow(node));
41        writer.insert_dirty(page)?;
42        return Ok(page_id);
43    }
44    let mut next_id = PageID(0);
45    for chunk in chunks.iter().rev() {
46        let page_id = writer.alloc_page_id();
47        let node = EventOverflowNode {
48            next: next_id,
49            data: (*chunk).to_vec(),
50        };
51        let page = Page::new(page_id, Node::EventOverflow(node));
52        writer.insert_dirty(page)?;
53        next_id = page_id;
54    }
55    Ok(next_id)
56}
57
58fn read_overflow_chain(
59    mvcc: &Mvcc,
60    dirty: &HashMap<PageID, Page>,
61    mut page_id: PageID,
62) -> DcbResult<Vec<u8>> {
63    let mut out: Vec<u8> = Vec::new();
64    while page_id.0 != 0 {
65        page_id = with_page(mvcc, dirty, page_id, |page| match &page.node {
66            Node::EventOverflow(node) => {
67                out.extend_from_slice(&node.data);
68                Ok(node.next)
69            }
70            _ => Err(DcbError::DatabaseCorrupted(
71                "Expected EventOverflow node".to_string(),
72            )),
73        })?;
74    }
75    Ok(out)
76}
77
78fn with_page<T, F>(
79    mvcc: &Mvcc,
80    dirty: &HashMap<PageID, Page>,
81    page_id: PageID,
82    f: F,
83) -> DcbResult<T>
84where
85    F: FnOnce(&Page) -> DcbResult<T>,
86{
87    if let Some(page) = dirty.get(&page_id) {
88        f(page)
89    } else {
90        let page = mvcc.read_page(page_id)?;
91        f(&page)
92    }
93}
94
95/// Writes a fresh overflow chain with serialized event metadata followed by the event data.
96///
97/// Metadata is placed first so that a future metadata-only read can stop after
98/// the leading page(s) without walking the potentially large event data.
99fn data_and_metadata_to_overflow_chain(
100    mvcc: &Mvcc,
101    writer: &mut Writer,
102    record_data: &Vec<u8>,
103    metadata: &Vec<(String, String)>,
104) -> DcbResult<(PageID, usize)> {
105    let metadata_len = if metadata.is_empty() {
106        0
107    } else {
108        metadata_serialized_size(metadata)
109    };
110
111    // 1. One allocation, perfectly sized for the final product
112    let mut combined = vec![0u8; metadata_len + record_data.len()];
113
114    if metadata_len > 0 {
115        // 2. Overwrite the first part with metadata
116        serialize_metadata_into(metadata, &mut combined[..metadata_len])?;
117    }
118
119    // 3. Overwrite the second part with record data (Fast memcpy)
120    combined[metadata_len..].copy_from_slice(record_data);
121
122    let root_id = write_overflow_chain(mvcc, writer, &combined)?;
123    Ok((root_id, metadata_len))
124}
125
126fn materialize_event_value(
127    mvcc: &Mvcc,
128    dirty: &HashMap<PageID, Page>,
129    value: &EventValue,
130) -> DcbResult<EventRecord> {
131    match value {
132        EventValue::Inline(rec) => Ok(rec.clone()),
133        EventValue::Overflow {
134            event_type,
135            data_len,
136            tags,
137            root_id,
138            uuid,
139            metadata_len,
140        } => {
141            let all_data = read_overflow_chain(mvcc, dirty, *root_id)?;
142            let all_data_len = all_data.len();
143            if (all_data_len as u64) != *data_len + *metadata_len {
144                return Err(DcbError::DatabaseCorrupted(format!(
145                    "Overflow data length mismatch, data: {data_len}, metadata: {metadata_len}"
146                )));
147            }
148            // The chain holds leading metadata bytes followed by event data.
149            let split = *metadata_len as usize;
150            let metadata = if *metadata_len > 0 {
151                deserialize_metadata(&all_data[..split])?
152            } else {
153                Vec::new()
154            };
155            let data = all_data[split..].to_vec();
156            Ok(EventRecord {
157                event_type: event_type.clone(),
158                data,
159                tags: tags.clone(),
160                uuid: *uuid,
161                metadata,
162            })
163        }
164    }
165}
166
167// TODO: Move this because it's now only used in tests.
168pub fn event_tree_append(
169    mvcc: &Mvcc,
170    writer: &mut Writer,
171    event_record: EventRecord,
172    position: Position,
173) -> DcbResult<()> {
174    let event_values_and_size_diffs =
175        validate_event_record_for_append(mvcc.page_size, event_record)?;
176    event_tree_append_event_value(mvcc, writer, event_values_and_size_diffs, position)
177}
178
179/// Append an event to the root event leaf page.
180///
181/// This function obtains a mutable reference to a dirty copy of the root event
182/// leaf page (using copy-on-write if necessary) and appends the provided
183/// Position to the keys and the EventRecord to the values.
184pub fn event_tree_append_event_value(
185    mvcc: &Mvcc,
186    writer: &mut Writer,
187    event_values_and_size_diffs: ((EventValue, usize), Option<(EventValue, usize)>),
188    position: Position,
189) -> DcbResult<()> {
190    let verbose = mvcc.verbose;
191    // if verbose {
192    //     println!("Appending event: {position:?} {inline_value:?}");
193    //     println!("Root is {:?}", writer.events_tree_root_id);
194    // }
195    // Get the current root page id for the event tree
196    let mut current_page_id: PageID = writer.events_tree_root_id;
197
198    // Traverse the tree to find a leaf node
199    let mut stack: Vec<PageID> = Vec::new();
200    loop {
201        let current_page_ref = writer.get_page_ref(mvcc, current_page_id)?;
202        if matches!(current_page_ref.node, Node::EventLeaf(_)) {
203            break;
204        }
205        if let Node::EventInternal(internal_node) = &current_page_ref.node {
206            if verbose {
207                println!("{:?} is internal node", current_page_ref.page_id);
208            }
209            stack.push(current_page_id);
210            current_page_id = *internal_node.child_ids.last().ok_or_else(|| {
211                DcbError::DatabaseCorrupted("Internal node has no children".to_string())
212            })?;
213        } else {
214            return Err(DcbError::DatabaseCorrupted(
215                "Expected EventInternal node".to_string(),
216            ));
217        }
218    }
219    if verbose {
220        println!("{current_page_id:?} is leaf node");
221    }
222
223    // Make the leaf page dirty
224    let dirty_page_id = { writer.get_dirty_page_id(current_page_id)? };
225    let replacement_info: Option<(PageID, PageID)> = {
226        if dirty_page_id != current_page_id {
227            Some((current_page_id, dirty_page_id))
228        } else {
229            None
230        }
231    };
232
233    let serialized_size = writer
234        .get_page_ref(mvcc, dirty_page_id)?
235        .calc_serialized_size();
236
237    // Decide what we need to do.
238    // - if we already know inline value is too large for any page, indicated by
239    //   the received overflow_value_and_size being Some(), we have to
240    //   consider what to do with the given overflow value. Either it fits in the
241    //   current page, or we must to split and append the overflow value to a new page.
242    // - otherwise, if the received over_value_and_size is None, then we know the given
243    //   inline_value can at least fit an empty page, so we have to decide whether it
244    //   will fit this page. If it does fit this page then we can append it.
245    // - otherwise, we need to decide between three options: (a) splitting the page
246    //   and appending the inline value to the new page, (b) making an overflow node
247    //   and seeing if it will fit this page, (c) splitting the page and appending an
248    //   overflow value to the new page. The decision is about avoiding sparse pages
249    //   as much as possible.
250    //
251    // Design:
252    // - if we don't have an overflow value, then append if inline value fits or split and append
253    // - otherwise, append if overflow value fits or split and append overflow value
254
255    let (event_value, size_diff) = match event_values_and_size_diffs {
256        ((EventValue::Inline(record), _), Some((mut overflow_value, overflow_size_diff))) => {
257            let (overflow_root_id, overflow_metadata_len) =
258                data_and_metadata_to_overflow_chain(mvcc, writer, &record.data, &record.metadata)?;
259            if let EventValue::Overflow {
260                root_id,
261                metadata_len,
262                ..
263            } = &mut overflow_value
264            {
265                *root_id = overflow_root_id;
266                *metadata_len = overflow_metadata_len as u64;
267            } else {
268                return Err(DcbError::InternalError(
269                    "Shouldn't get here when setting root_id on event overflow value".to_string(),
270                ));
271            }
272
273            (overflow_value, overflow_size_diff)
274        }
275        ((inline_value, inline_size_diff), None) => (inline_value, inline_size_diff),
276        _ => {
277            return Err(DcbError::InternalError(
278                "Shouldn't get here when matching event_values_and_size_diffs".to_string(),
279            ));
280        }
281    };
282
283    // We may need to split the page.
284    let mut popped: Option<EventValue> = None;
285
286    if serialized_size + size_diff <= mvcc.page_size {
287        // Get a mutable leaf node and append the event value and position.
288        let dirty_leaf_page = writer.get_mut_dirty(dirty_page_id)?;
289        match &mut dirty_leaf_page.node {
290            Node::EventLeaf(node) => {
291                if verbose {
292                    println!(
293                        "Pushing {:?} at {:?} onto {:?}",
294                        event_value, position, dirty_page_id,
295                    );
296                }
297                node.keys.push(position);
298                node.values.push(event_value);
299                // // Final check that the serialized size is okay.
300                // let serialized_size = dirty_leaf_page.calc_serialized_size();
301                // if serialized_size > mvcc.page_size {
302                //     return Err(DcbError::DatabaseCorrupted(format!(
303                //         "New event leaf page too large after appending value: {serialized_size}, max: {})",
304                //         mvcc.page_size
305                //     )));
306                // }
307            }
308            _ => {
309                return Err(DcbError::DatabaseCorrupted(
310                    "Expected EventLeaf node".to_string(),
311                ));
312            }
313        }
314    } else {
315        // Let's split and add this event value to a new page.
316        popped = Some(event_value);
317    }
318
319    // Prepare for split propagation
320    let mut split_info: Option<(Position, PageID)> = None;
321
322    if let Some(event_value) = popped {
323        // Build new leaf node page.
324        let new_leaf_page_id = writer.alloc_page_id();
325        let new_leaf_node = EventLeafNode {
326            keys: vec![position],
327            values: vec![event_value],
328        };
329        let new_leaf_page = Page::new(new_leaf_page_id, Node::EventLeaf(new_leaf_node.clone()));
330        // // Final check that the serialized size is okay.
331        // let serialized_size = new_leaf_page.calc_serialized_size();
332        // if serialized_size > mvcc.page_size {
333        //     return Err(DcbError::DatabaseCorrupted(format!(
334        //         "New event leaf page too large after appending value: {serialized_size}, max: {})",
335        //         mvcc.page_size
336        //     )));
337        // }
338        if verbose {
339            println!(
340                "Created new leaf {:?}: {:?}",
341                new_leaf_page_id, new_leaf_page.node
342            );
343        }
344        writer.insert_dirty(new_leaf_page)?;
345        if verbose {
346            println!("Promoting {position:?} and {new_leaf_page_id:?}");
347        }
348        split_info = Some((position, new_leaf_page_id));
349    }
350
351    // Propagate splits and replacements up the stack
352    let mut current_replacement_info = replacement_info;
353    while let Some(parent_page_id) = stack.pop() {
354        // Make the internal page dirty
355        let dirty_page_id = writer.get_dirty_page_id(parent_page_id)?;
356        let parent_replacement_info: Option<(PageID, PageID)> = {
357            if dirty_page_id != parent_page_id {
358                Some((parent_page_id, dirty_page_id))
359            } else {
360                None
361            }
362        };
363        // Get a mutable internal node....
364        let dirty_internal_page = writer.get_mut_dirty(dirty_page_id)?;
365
366        if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
367            if let Some((old_id, new_id)) = current_replacement_info {
368                dirty_internal_node.replace_last_child_id(old_id, new_id)?;
369                if verbose {
370                    println!(
371                        "Replaced {old_id:?} with {new_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
372                    );
373                }
374            } else if verbose {
375                println!("Nothing to replace in {dirty_page_id:?}")
376            }
377        } else {
378            return Err(DcbError::DatabaseCorrupted(
379                "Expected EventInternal node".to_string(),
380            ));
381        }
382
383        if let Some((promoted_key, promoted_page_id)) = split_info {
384            if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
385                // Add the promoted key and page ID
386                dirty_internal_node
387                    .append_promoted_key_and_page_id(promoted_key, promoted_page_id)?;
388
389                if verbose {
390                    println!(
391                        "Appended promoted key {promoted_key:?} and child {promoted_page_id:?} in {dirty_page_id:?}: {dirty_internal_node:?}"
392                    );
393                }
394            } else {
395                return Err(DcbError::DatabaseCorrupted(
396                    "Expected EventInternal node".to_string(),
397                ));
398            }
399        }
400
401        // Check if the internal page needs splitting
402
403        if dirty_internal_page.calc_serialized_size() > mvcc.page_size {
404            if let Node::EventInternal(dirty_internal_node) = &mut dirty_internal_page.node {
405                if verbose {
406                    println!("Splitting internal {dirty_page_id:?}...");
407                }
408                // Split the internal node
409                // Ensure we have at least 3 keys and 4 child IDs before splitting
410                if dirty_internal_node.keys.len() < 3 || dirty_internal_node.child_ids.len() < 4 {
411                    return Err(DcbError::DatabaseCorrupted(
412                        "Cannot split internal node with too few keys/children".to_string(),
413                    ));
414                }
415
416                // Move the right-most key to a new node. Promote the next right-most key.
417                let (promoted_key, new_keys, new_child_ids) = dirty_internal_node.split_off()?;
418
419                // Ensure old node maintain the B-tree invariant: n keys should have n+1 child pointers
420                assert_eq!(
421                    dirty_internal_node.keys.len() + 1,
422                    dirty_internal_node.child_ids.len()
423                );
424
425                let new_internal_node = EventInternalNode {
426                    keys: new_keys,
427                    child_ids: new_child_ids,
428                };
429
430                // Ensure the new node also maintains the invariant
431                assert_eq!(
432                    new_internal_node.keys.len() + 1,
433                    new_internal_node.child_ids.len()
434                );
435
436                // Create a new internal page.
437                let new_internal_page_id = writer.alloc_page_id();
438                let new_internal_page =
439                    Page::new(new_internal_page_id, Node::EventInternal(new_internal_node));
440                if verbose {
441                    println!(
442                        "Created internal {:?}: {:?}",
443                        new_internal_page_id, new_internal_page.node
444                    );
445                }
446                writer.insert_dirty(new_internal_page)?;
447
448                split_info = Some((promoted_key, new_internal_page_id));
449            } else {
450                return Err(DcbError::DatabaseCorrupted(
451                    "Expected EventInternal node".to_string(),
452                ));
453            }
454        } else {
455            split_info = None;
456        }
457        current_replacement_info = parent_replacement_info;
458    }
459
460    if let Some((old_id, new_id)) = current_replacement_info {
461        if writer.events_tree_root_id == old_id {
462            writer.events_tree_root_id = new_id;
463            if verbose {
464                println!("Replaced root {old_id:?} with {new_id:?}");
465            }
466        } else {
467            return Err(DcbError::RootIDMismatch(old_id.0, new_id.0));
468        }
469    }
470
471    if let Some((promoted_key, promoted_page_id)) = split_info {
472        // Create a new root
473        let new_internal_node = EventInternalNode {
474            keys: vec![promoted_key],
475            child_ids: vec![writer.events_tree_root_id, promoted_page_id],
476        };
477
478        let new_root_page_id = writer.alloc_page_id();
479        let new_root_page = Page::new(new_root_page_id, Node::EventInternal(new_internal_node));
480        if verbose {
481            println!(
482                "Created new internal root {:?}: {:?}",
483                new_root_page_id, new_root_page.node
484            );
485        }
486        writer.insert_dirty(new_root_page)?;
487
488        writer.events_tree_root_id = new_root_page_id;
489    }
490
491    Ok(())
492}
493
494pub fn event_tree_lookup(
495    mvcc: &Mvcc,
496    dirty: &HashMap<PageID, Page>,
497    events_tree_root_id: PageID,
498    position: Position,
499) -> DcbResult<EventRecord> {
500    enum LookupStep {
501        Next(PageID),
502        Found(EventRecord),
503    }
504
505    let mut current_page_id: PageID = events_tree_root_id;
506    loop {
507        let step = with_page(mvcc, dirty, current_page_id, |page| match &page.node {
508            Node::EventInternal(internal) => {
509                // Choose child based on upper bound of position in separator keys
510                let idx = match internal.keys.binary_search(&position) {
511                    Ok(i) => i + 1,
512                    Err(i) => i,
513                };
514                if idx >= internal.child_ids.len() {
515                    return Err(DcbError::DatabaseCorrupted(
516                        "Child index out of bounds in event tree".to_string(),
517                    ));
518                }
519                Ok(LookupStep::Next(internal.child_ids[idx]))
520            }
521            Node::EventLeaf(leaf) => match leaf.keys.binary_search(&position) {
522                Ok(i) => {
523                    let rec = materialize_event_value(mvcc, dirty, &leaf.values[i])?;
524                    Ok(LookupStep::Found(rec))
525                }
526                Err(_) => Err(DcbError::DatabaseCorrupted(format!(
527                    "Event at position {position:?} not found",
528                ))),
529            },
530            _ => Err(DcbError::DatabaseCorrupted(format!(
531                "Expected EventInternal or EventLeaf node in event tree, got {}",
532                page.node.type_name()
533            ))),
534        })?;
535
536        match step {
537            LookupStep::Next(next_page_id) => current_page_id = next_page_id,
538            LookupStep::Found(rec) => return Ok(rec),
539        }
540    }
541}
542
543pub struct EventIterator<'a> {
544    pub mvcc: &'a Mvcc,
545    pub dirty: &'a HashMap<PageID, Page>,
546    pub stack: Vec<(PageID, Option<usize>)>,
547    pub page_cache: HashMap<PageID, Arc<Page>>,
548    pub start: Option<Position>, // inclusive position, better for binary search
549    pub backwards: bool,
550}
551
552impl<'a> EventIterator<'a> {
553    pub fn new(
554        mvcc: &'a Mvcc,
555        dirty: &'a HashMap<PageID, Page>,
556        events_tree_root_id: PageID,
557        start: Option<Position>,
558        backwards: bool,
559    ) -> Self {
560        let next_position = (events_tree_root_id, None);
561        Self {
562            mvcc,
563            dirty,
564            stack: vec![next_position],
565            page_cache: HashMap::new(),
566            start,
567            backwards,
568        }
569    }
570
571    pub fn next_batch(
572        &mut self,
573        batch_size: u32,
574        cancel: Option<&std::sync::Arc<std::sync::atomic::AtomicBool>>,
575    ) -> DcbResult<Vec<(Position, EventRecord)>> {
576        let mut result: Vec<(Position, EventRecord)> = Vec::with_capacity(batch_size as usize);
577        if batch_size == 0 {
578            return Ok(result);
579        }
580        if self.backwards && self.start == Some(Position(0)) {
581            return Ok(result);
582        }
583        while result.len() < batch_size as usize {
584            if let Some(c) = cancel {
585                if c.load(std::sync::atomic::Ordering::Relaxed) {
586                    return Err(DcbError::CancelledByUser());
587                }
588            }
589            let Some((page_id, mut stacked_idx)) = self.stack.pop() else {
590                break; // traversal finished
591            };
592
593            // Compute actions under a scoped immutable borrow, then mutate cache/stack afterwards.
594            let mut remove_page = false;
595            let mut push_revisit: Option<(PageID, Option<usize>)> = None;
596            let mut push_child: Option<(PageID, Option<usize>)> = None; // (child_id, stacked_keys_idx)
597            let mut emit_event: Option<(Position, EventRecord)> = None;
598
599            {
600                // Obtain the current page (from dirty, or page cache, or deserialize).
601                let page_ref: &Page = if let Some(p) = self.dirty.get(&page_id) {
602                    p
603                } else if let Some(p) = self.page_cache.get(&page_id) {
604                    p.as_ref()
605                } else {
606                    let page_arc = self.mvcc.read_page(page_id)?;
607                    self.page_cache.insert(page_id, page_arc);
608                    self.page_cache
609                        .get(&page_id)
610                        .expect("page should be in cache")
611                        .as_ref()
612                };
613
614                match &page_ref.node {
615                    Node::EventInternal(internal) => {
616                        // println!("Visit internal {page_id:?}");
617                        if stacked_idx.is_none() && !internal.keys.is_empty() {
618                            // println!(" - first visit");
619                            // println!(" - keys: {:?}", internal.keys.clone());
620                            // println!(" - child_ids: {:?}", internal.child_ids.clone());
621                            // println!(" - from: {:?}", self.from);
622
623                            stacked_idx = match &self.start {
624                                Some(from) => match internal.keys.binary_search(from) {
625                                    Ok(i) => Some(i + 1),
626                                    Err(i) => Some(i),
627                                },
628                                None => {
629                                    if !self.backwards {
630                                        Some(0)
631                                    } else {
632                                        Some(internal.child_ids.len() - 1)
633                                    }
634                                }
635                            };
636                        }
637
638                        if let Some(child_ids_idx) = stacked_idx {
639                            // println!(" - child ids index: {} / {}", child_ids_idx + 1, internal.child_ids.len());
640                            // println!(" - will visit child: {:?}", internal.child_ids[child_ids_idx]);
641                            // Push the chosen child.
642                            push_child = Some((internal.child_ids[child_ids_idx], None));
643                            // Do or don't revisit this internal node?
644                            if !self.backwards {
645                                if child_ids_idx + 1 < internal.child_ids.len() {
646                                    // Will revisit this internal node.
647                                    // println!(" - will revisit");
648                                    push_revisit = Some((page_id, Some(child_ids_idx + 1)));
649                                } else {
650                                    // Don't revisit this internal node.
651                                    remove_page = true;
652                                    // println!(" - will remove");
653                                }
654                            } else if child_ids_idx > 0 {
655                                // Will revisit this internal node.
656                                // println!(" - will revisit");
657                                push_revisit = Some((page_id, Some(child_ids_idx - 1)));
658                            } else {
659                                // Don't revisit this internal node.
660                                remove_page = true;
661                                // println!(" - will remove");
662                            }
663                        } else {
664                            // Shouldn't get here unless an internal
665                            // node has no keys, which it should never do.
666                            remove_page = true
667                        };
668                    }
669                    Node::EventLeaf(leaf) => {
670                        // println!("Visit leaf {page_id:?}");
671                        if stacked_idx.is_none() {
672                            // println!(" - first visit");
673                            // println!(" - keys: {:?}", leaf.keys.clone());
674                            let values_len = leaf.values.len();
675
676                            stacked_idx = if values_len > 0 {
677                                match &self.start {
678                                    Some(from) => match leaf.keys.binary_search(from) {
679                                        Ok(i) => Some(i),
680                                        Err(i) => {
681                                            if !self.backwards {
682                                                Some(i)
683                                            } else {
684                                                Some(i - 1)
685                                            }
686                                        }
687                                    },
688                                    None => {
689                                        if !self.backwards {
690                                            Some(0)
691                                        } else {
692                                            Some(values_len - 1)
693                                        }
694                                    }
695                                }
696                            } else {
697                                None
698                            }
699                        }
700
701                        if let Some(values_idx) = stacked_idx {
702                            // println!(" - values index: {} / {}", values_idx + 1, leaf.values.len());
703                            if values_idx < leaf.values.len() {
704                                let event_position = leaf.keys[values_idx];
705                                let event_record = materialize_event_value(
706                                    self.mvcc,
707                                    self.dirty,
708                                    &leaf.values[values_idx],
709                                )?;
710                                // println!(" - emit event position: {:?}", event_position.clone());
711                                emit_event = Some((event_position, event_record));
712
713                                if !self.backwards {
714                                    if values_idx + 1 < leaf.values.len() {
715                                        // Revisit this leaf.
716                                        push_revisit = Some((page_id, Some(values_idx + 1)));
717                                        // println!(" - not last value, will revisit");
718                                    } else {
719                                        // The last value.
720                                        remove_page = true;
721                                        // println!(" - last value, will remove");
722                                    }
723                                } else if values_idx > 0 {
724                                    // Revisit this leaf.
725                                    push_revisit = Some((page_id, Some(values_idx - 1)));
726                                    // println!(" - not last value, will revisit");
727                                } else {
728                                    // The last value.
729                                    remove_page = true;
730                                    // println!(" - last value, will remove");
731                                }
732                            } else {
733                                // No key greater or equal to 'from' in this leaf
734                                // println!(" - value index out of range, why wasn't this removed?");
735                                remove_page = true;
736                            }
737                        } else {
738                            // No leaf values.
739                            remove_page = true;
740                        }
741                    }
742                    _ => {
743                        return Err(DcbError::DatabaseCorrupted(format!(
744                            "Expected EventInternal or EventLeaf node in event tree, got {}",
745                            page_ref.node.type_name()
746                        )));
747                    }
748                }
749            }
750
751            // Mutations after the borrow has ended
752            if let Some(revisit) = push_revisit {
753                // Revisit must be pushed first so that the child is processed next (LIFO)
754                self.stack.push(revisit);
755            }
756            if let Some((child_id, child_start_idx)) = push_child {
757                self.stack.push((child_id, child_start_idx));
758            }
759            if let Some((event_position, event_record)) = emit_event {
760                result.push((event_position, event_record));
761            }
762            if remove_page {
763                self.page_cache.remove(&page_id);
764            }
765        }
766        Ok(result)
767    }
768}
769
770#[cfg(test)]
771mod tests {
772    use super::*;
773    use crate::mvcc::{Mvcc, StorageOptions};
774    use crate::node::Node;
775    use rand::random;
776    use serial_test::serial;
777    use tempfile::tempdir;
778
779    static VERBOSE: bool = false;
780
781    // Helper function to create a test database with a specified page size
782    fn construct_db(page_size: usize) -> (tempfile::TempDir, Mvcc) {
783        let temp_dir = tempdir().unwrap();
784        let db_path = temp_dir.path().join("mvcc-test.db");
785        let db = Mvcc::new(
786            VERBOSE,
787            StorageOptions::default()
788                .db_path(db_path)
789                .page_size(page_size),
790        )
791        .unwrap();
792        (temp_dir, db)
793    }
794
795    #[test]
796    #[serial]
797    fn test_append_event_to_empty_leaf_root() {
798        // Setup a temporary database
799        let (_temp_dir, db) = construct_db(64);
800
801        // Start a writer
802        let mut writer = db.writer().unwrap();
803
804        // Issue a new position
805        let position = writer.issue_position();
806
807        // Create an event record
808        let record = EventRecord {
809            event_type: "UserCreated".to_string(),
810            data: vec![1, 2, 3, 4],
811            tags: vec!["users".to_string(), "creation".to_string()],
812            uuid: None,
813            metadata: Vec::new(),
814        };
815
816        // Call append_event
817        event_tree_append(&db, &mut writer, record.clone(), position).unwrap();
818
819        // Verify that the dirty root page contains the appended key/value
820        let new_root_id = writer.events_tree_root_id;
821        assert!(writer.dirty.contains_key(&new_root_id));
822        let page = writer.dirty.get(&new_root_id).unwrap();
823        match &page.node {
824            Node::EventLeaf(node) => {
825                assert_eq!(vec![position], node.keys);
826                assert_eq!(
827                    vec![crate::events_tree_nodes::EventValue::Inline(record.clone())],
828                    node.values
829                );
830            }
831            _ => panic!("Expected EventLeaf node"),
832        }
833
834        // Commit the writer and verify persistence
835        db.commit(&mut writer).unwrap();
836
837        // Read back the latest header and the persisted root event leaf page
838        let header_page = db.get_latest_header_page().unwrap();
839        let header = header_page.as_header_node().unwrap();
840        let persisted_page = db.read_page(header.events_tree_root_id).unwrap();
841        match &persisted_page.node {
842            Node::EventLeaf(node) => {
843                assert_eq!(vec![position], node.keys);
844                assert_eq!(
845                    vec![crate::events_tree_nodes::EventValue::Inline(record)],
846                    node.values
847                );
848            }
849            _ => panic!("Expected EventLeaf node after commit"),
850        }
851    }
852
853    #[test]
854    #[serial]
855    fn test_insert_events_until_split_leaf_one_writer() {
856        // Setup a temporary database
857        let (_temp_dir, db) = construct_db(256);
858
859        // Start a writer
860        let mut writer = db.writer().unwrap();
861
862        let mut has_split_leaf = false;
863        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
864
865        // Insert events until we split a leaf
866        while !has_split_leaf {
867            // Issue a new position and create a record
868            let position = writer.issue_position();
869            let record = EventRecord {
870                event_type: "UserCreated".to_string(),
871                data: (0..8).map(|_| random::<u8>()).collect(),
872                tags: vec!["users".to_string(), "creation".to_string()],
873                uuid: None,
874                metadata: Vec::new(),
875            };
876            appended.push((position, record.clone()));
877
878            // Append the event
879            event_tree_append(&db, &mut writer, record, position).unwrap();
880
881            // Check if we've split the leaf
882            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
883            match &root_page.node {
884                Node::EventInternal(_) => {
885                    has_split_leaf = true;
886                }
887                _ => {}
888            }
889        }
890
891        // Check keys and values of all pages
892        let mut copy_inserted = appended.clone();
893
894        // Get the root node
895        let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
896        let root_node = match &root_page.node {
897            Node::EventInternal(node) => node,
898            _ => panic!("Expected EventInternal node"),
899        };
900
901        // Check each child of the root
902        for (i, &child_id) in root_node.child_ids.iter().enumerate() {
903            let child_page = writer.dirty.get(&child_id).unwrap();
904            assert_eq!(child_id, child_page.page_id);
905
906            let child_node = match &child_page.node {
907                Node::EventLeaf(node) => node,
908                _ => panic!("Expected EventLeaf node"),
909            };
910
911            // Check that the keys are properly ordered
912            if i > 0 {
913                assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
914            }
915
916            // Check each key and value in the child
917            for (k, &key) in child_node.keys.iter().enumerate() {
918                let record = &child_node.values[k];
919                let (appended_position, appended_record) = copy_inserted.remove(0);
920                assert_eq!(appended_position, key);
921                assert_eq!(appended_record, record.clone());
922            }
923        }
924    }
925
926    #[test]
927    #[serial]
928    fn test_insert_events_until_split_leaf_many_writers() {
929        // Setup a temporary database
930        let (_temp_dir, db) = construct_db(256);
931
932        let mut has_split_leaf = false;
933        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
934
935        // Insert events until we split a leaf
936        while !has_split_leaf {
937            // Start a writer
938            let mut writer = db.writer().unwrap();
939
940            // Issue a new position and create a record
941            let position = writer.issue_position();
942            let record = EventRecord {
943                event_type: "UserCreated".to_string(),
944                data: (0..8).map(|_| random::<u8>()).collect(),
945                tags: vec!["users".to_string(), "creation".to_string()],
946                uuid: None,
947                metadata: Vec::new(),
948            };
949            appended.push((position, record.clone()));
950
951            // Append the event
952            event_tree_append(&db, &mut writer, record, position).unwrap();
953
954            // Check if we've split the leaf
955            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
956            match &root_page.node {
957                Node::EventInternal(_) => {
958                    has_split_leaf = true;
959                }
960                _ => {}
961            }
962
963            db.commit(&mut writer).unwrap();
964        }
965
966        // Check keys and values of all pages
967        let mut copy_inserted = appended.clone();
968
969        // Start a writer
970        let writer = db.writer().unwrap();
971
972        // Get the root node
973        let root_page = db.read_page(writer.events_tree_root_id).unwrap();
974        let root_node = match &root_page.node {
975            Node::EventInternal(node) => node,
976            _ => panic!("Expected EventInternal node"),
977        };
978
979        // Check each child of the root
980        for (i, &child_id) in root_node.child_ids.iter().enumerate() {
981            let child_page = db.read_page(child_id).unwrap();
982            assert_eq!(child_id, child_page.page_id);
983
984            let child_node = match &child_page.node {
985                Node::EventLeaf(node) => node,
986                _ => panic!("Expected EventLeaf node"),
987            };
988
989            // Check that the keys are properly ordered
990            if i > 0 {
991                assert_eq!(root_node.keys[i - 1], child_node.keys[0]);
992            }
993
994            // Check each key and value in the child
995            for (k, &key) in child_node.keys.iter().enumerate() {
996                let record = &child_node.values[k];
997                let (appended_position, appended_record) = copy_inserted.remove(0);
998                assert_eq!(appended_position, key);
999                assert_eq!(appended_record, record.clone());
1000            }
1001        }
1002    }
1003
1004    #[test]
1005    #[serial]
1006    fn test_insert_events_until_split_internal_one_writer() {
1007        // Setup a temporary database
1008        let (_temp_dir, db) = construct_db(256);
1009
1010        // Start a writer
1011        let mut writer = db.writer().unwrap();
1012
1013        let mut has_split_internal = false;
1014        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1015
1016        // Insert events until we split a leaf
1017        while !has_split_internal {
1018            // Issue a new position and create a record
1019            let position = writer.issue_position();
1020            let record = EventRecord {
1021                event_type: "UserCreated".to_string(),
1022                data: (0..8).map(|_| random::<u8>()).collect(),
1023                tags: vec!["users".to_string(), "creation".to_string()],
1024                uuid: None,
1025                metadata: Vec::new(),
1026            };
1027            appended.push((position, record.clone()));
1028
1029            // Append the event
1030            event_tree_append(&db, &mut writer, record, position).unwrap();
1031
1032            // Check if we've split an internal node
1033            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1034            match &root_page.node {
1035                Node::EventInternal(root_node) => {
1036                    // Check if the first child is an internal node
1037                    if !root_node.child_ids.is_empty() {
1038                        let child_id = root_node.child_ids[0];
1039                        if let Some(child_page) = writer.dirty.get(&child_id) {
1040                            match &child_page.node {
1041                                Node::EventInternal(_) => {
1042                                    has_split_internal = true;
1043                                }
1044                                _ => {}
1045                            }
1046                        }
1047                    }
1048                }
1049                _ => {}
1050            }
1051        }
1052
1053        // Check keys and values of all pages
1054        let mut copy_inserted = appended.clone();
1055
1056        // Get the root node
1057        let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1058        let root_node = match &root_page.node {
1059            Node::EventInternal(node) => node,
1060            _ => panic!("Expected EventInternal node"),
1061        };
1062
1063        // Check each child of the root
1064        for &child_id in root_node.child_ids.iter() {
1065            let child_page = writer.dirty.get(&child_id).unwrap();
1066            assert_eq!(child_id, child_page.page_id);
1067
1068            let child_node = match &child_page.node {
1069                Node::EventInternal(node) => node,
1070                _ => panic!("Expected EventInternal node"),
1071            };
1072
1073            for &grand_child_id in child_node.child_ids.iter() {
1074                let grand_child_page = writer.dirty.get(&grand_child_id).unwrap();
1075                assert_eq!(grand_child_id, grand_child_page.page_id);
1076
1077                let grand_child_node = match &grand_child_page.node {
1078                    Node::EventLeaf(node) => node,
1079                    _ => panic!("Expected EventLeaf node"),
1080                };
1081
1082                // Check each key and value in the child
1083                for (k, &key) in grand_child_node.keys.iter().enumerate() {
1084                    let record = &grand_child_node.values[k];
1085                    let (appended_position, appended_record) = copy_inserted.remove(0);
1086                    assert_eq!(appended_position, key);
1087                    assert_eq!(appended_record, record.clone());
1088                }
1089            }
1090        }
1091    }
1092
1093    #[test]
1094    #[serial]
1095    fn test_insert_events_until_split_internal_many_writers() {
1096        // Setup a temporary database
1097        let (_temp_dir, db) = construct_db(512);
1098
1099        let mut has_split_internal = false;
1100        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1101
1102        // Insert events until we split a root internal node
1103        while !has_split_internal {
1104            // Start a writer
1105            let mut writer = db.writer().unwrap();
1106
1107            // Issue a new position and create a record
1108            let position = writer.issue_position();
1109            let record = EventRecord {
1110                event_type: "UserCreated".to_string(),
1111                data: (0..8).map(|_| random::<u8>()).collect(),
1112                tags: vec!["users".to_string(), "creation".to_string()],
1113                uuid: None,
1114                metadata: Vec::new(),
1115            };
1116            appended.push((position, record.clone()));
1117
1118            // Append the event
1119            event_tree_append(&db, &mut writer, record, position).unwrap();
1120
1121            // Check if the root is an internal node
1122            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1123            match &root_page.node {
1124                Node::EventInternal(root_node) => {
1125                    // Check if the first child is an internal node
1126                    if !root_node.child_ids.is_empty() {
1127                        let child_id = root_node.child_ids[0];
1128                        if let Some(child_page) = writer.dirty.get(&child_id) {
1129                            match &child_page.node {
1130                                Node::EventInternal(_) => {
1131                                    has_split_internal = true;
1132                                }
1133                                _ => {}
1134                            }
1135                        }
1136                    }
1137                }
1138                _ => {}
1139            }
1140            db.commit(&mut writer).unwrap();
1141        }
1142
1143        // Check keys and values of all pages
1144        let mut copy_inserted = appended.clone();
1145
1146        // Start a writer
1147        let writer = db.writer().unwrap();
1148
1149        // Get the root node
1150        let root_page = db.read_page(writer.events_tree_root_id).unwrap();
1151        let root_node = match &root_page.node {
1152            Node::EventInternal(node) => node,
1153            _ => panic!("Expected EventInternal node"),
1154        };
1155
1156        // Check each child of the root
1157        for &child_id in root_node.child_ids.iter() {
1158            let child_page = db.read_page(child_id).unwrap();
1159            assert_eq!(child_id, child_page.page_id);
1160
1161            let child_node = match &child_page.node {
1162                Node::EventInternal(node) => node,
1163                _ => panic!("Expected EventInternal node"),
1164            };
1165
1166            for &grand_child_id in child_node.child_ids.iter() {
1167                let grand_child_page = db.read_page(grand_child_id).unwrap();
1168                assert_eq!(grand_child_id, grand_child_page.page_id);
1169
1170                let grand_child_node = match &grand_child_page.node {
1171                    Node::EventLeaf(node) => node,
1172                    _ => panic!("Expected EventLeaf node"),
1173                };
1174
1175                // Check each key and value in the child
1176                for (k, &key) in grand_child_node.keys.iter().enumerate() {
1177                    let record = &grand_child_node.values[k];
1178                    let (appended_position, appended_record) = copy_inserted.remove(0);
1179                    // println!("Checking appended event: {appended_position:?} {appended_record:?}");
1180                    assert_eq!(appended_position, key);
1181                    assert_eq!(appended_record, record.clone());
1182                }
1183            }
1184        }
1185        assert_eq!(0, copy_inserted.len());
1186    }
1187
1188    #[test]
1189    #[serial]
1190    fn test_read_events_all() {
1191        // Setup a temporary database
1192        let (_temp_dir, db) = construct_db(512);
1193
1194        let mut has_split_internal = false;
1195        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1196
1197        // Insert events until we split a root internal node
1198        while !has_split_internal {
1199            // Start a writer
1200            let mut writer = db.writer().unwrap();
1201
1202            // Issue a new position and create a record
1203            let position = writer.issue_position();
1204            let record = EventRecord {
1205                event_type: "UserCreated".to_string(),
1206                data: (0..8).map(|_| random::<u8>()).collect(),
1207                tags: vec!["users".to_string(), "creation".to_string()],
1208                uuid: None,
1209                metadata: Vec::new(),
1210            };
1211            appended.push((position, record.clone()));
1212
1213            // Append the event
1214            event_tree_append(&db, &mut writer, record, position).unwrap();
1215
1216            // Check if the root is an internal node
1217            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1218            match &root_page.node {
1219                Node::EventInternal(root_node) => {
1220                    // Check if the first child is an internal node
1221                    if !root_node.child_ids.is_empty() {
1222                        let child_id = root_node.child_ids[0];
1223                        if let Some(child_page) = writer.dirty.get(&child_id) {
1224                            match &child_page.node {
1225                                Node::EventInternal(_) => {
1226                                    has_split_internal = true;
1227                                }
1228                                _ => {}
1229                            }
1230                        }
1231                    }
1232                }
1233                _ => {}
1234            }
1235            db.commit(&mut writer).unwrap();
1236        }
1237
1238        // Check keys and values of all pages
1239        let copy_inserted = appended.clone();
1240
1241        // Start a reader
1242        let reader = db.reader().unwrap();
1243        let events_tree_root_id = reader.events_tree_root_id;
1244        let reader_tsn = reader.tsn;
1245
1246        let dirty = HashMap::new();
1247        let mut events_iterator = EventIterator::new(&db, &dirty, events_tree_root_id, None, false);
1248
1249        // Ensure the reader's tsn is registered while the iterator is alive
1250        {
1251            // Check registration state directly on the reader-TSN registry
1252            assert!(
1253                db.reader_tsns.contains_tsn(reader_tsn),
1254                "TSN should remain registered until reader is dropped"
1255            );
1256        }
1257
1258        // Progressively iterate over events using batches
1259        let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1260        loop {
1261            let batch = events_iterator.next_batch(3, None).unwrap();
1262            if batch.is_empty() {
1263                break;
1264            }
1265            scanned.extend(batch);
1266
1267            // The reader should remain registered throughout iteration
1268            // Check registration state directly on the reader-TSN registry
1269            assert!(
1270                db.reader_tsns.contains_tsn(reader_tsn),
1271                "TSN should remain registered until reader is dropped"
1272            );
1273        }
1274
1275        assert_eq!(copy_inserted.len(), scanned.len());
1276        for (i, expected) in copy_inserted.iter().enumerate() {
1277            assert_eq!(expected.0, scanned[i].0);
1278            assert_eq!(expected.1, scanned[i].1);
1279        }
1280
1281        // Additionally, validate lookup_event for each appended position using the existing reader in the iterator
1282        let dirty = HashMap::new();
1283        for (pos, expected_rec) in copy_inserted.iter() {
1284            let found = event_tree_lookup(&db, &dirty, events_tree_root_id, *pos).unwrap();
1285            assert_eq!(expected_rec, &found);
1286        }
1287
1288        // Ensure we did not accumulate pages in the iterator cache
1289        assert!(
1290            events_iterator.page_cache.is_empty(),
1291            "EventIterator page_cache should be empty after full scan"
1292        );
1293
1294        // While iterator is still alive, the reader should still be registered
1295        {
1296            // Check registration state directly on the reader-TSN registry
1297            assert!(
1298                db.reader_tsns.contains_tsn(reader_tsn),
1299                "TSN should remain registered until reader is dropped"
1300            );
1301        }
1302
1303        // Drop the reader and ensure the reader tsn is removed
1304        drop(reader);
1305        {
1306            // Check registration state directly on the reader-TSN registry
1307            assert!(
1308                !db.reader_tsns.contains_tsn(reader_tsn),
1309                "TSN should be removed after reader is dropped"
1310            );
1311            assert_eq!(0, db.reader_tsns.live_reader_count());
1312        }
1313    }
1314
1315    #[test]
1316    #[serial]
1317    fn test_read_events_from_forwards() {
1318        // Setup a temporary database
1319        let (_temp_dir, db) = construct_db(128);
1320
1321        let mut has_split_internal = false;
1322        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1323
1324        // Insert events until we split a root internal node
1325        while !has_split_internal {
1326            // Start a writer
1327            let mut writer = db.writer().unwrap();
1328
1329            // Issue a new position and create a record
1330            let position = writer.issue_position();
1331            let record = EventRecord {
1332                event_type: "UserCreated".to_string(),
1333                data: (0..8).map(|_| random::<u8>()).collect(),
1334                tags: vec!["users".to_string(), "creation".to_string()],
1335                uuid: None,
1336                metadata: Vec::new(),
1337            };
1338            appended.push((position, record.clone()));
1339
1340            // Append the event
1341            event_tree_append(&db, &mut writer, record, position).unwrap();
1342
1343            // Check if the root is an internal node
1344            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1345            match &root_page.node {
1346                Node::EventInternal(root_node) => {
1347                    // Check if the first child is an internal node
1348                    if !root_node.child_ids.is_empty() {
1349                        let child_id = root_node.child_ids[0];
1350                        if let Some(child_page) = writer.dirty.get(&child_id) {
1351                            match &child_page.node {
1352                                Node::EventInternal(_) => {
1353                                    has_split_internal = true;
1354                                }
1355                                _ => {}
1356                            }
1357                        }
1358                    }
1359                }
1360                _ => {}
1361            }
1362            db.commit(&mut writer).unwrap();
1363        }
1364
1365        let appended_len = appended.len();
1366        // println!("Appended number: {}", appended_len);
1367
1368        // Iterate through various 'from' positions.
1369        for i in 0..appended_len + 2 {
1370            let from = Position(i as u64);
1371
1372            // Start a reader
1373            let reader = db.reader().unwrap();
1374            let events_tree_root_id = reader.events_tree_root_id;
1375            let reader_tsn = reader.tsn;
1376            let dirty = HashMap::new();
1377            let mut events_iterator =
1378                EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), false);
1379
1380            // Ensure the reader's tsn is registered while the iterator is alive
1381            {
1382                // Check registration state directly on the reader-TSN registry
1383                assert!(
1384                    db.reader_tsns.contains_tsn(reader_tsn),
1385                    "TSN should remain registered until reader is dropped"
1386                );
1387            }
1388
1389            // Progressively iterate over events using batches
1390            let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1391            loop {
1392                let batch = events_iterator.next_batch(3, None).unwrap();
1393                if batch.is_empty() {
1394                    break;
1395                }
1396                scanned.extend(batch);
1397
1398                // The reader should remain registered throughout iteration
1399                // Check registration state directly on the reader-TSN registry
1400                assert!(
1401                    db.reader_tsns.contains_tsn(reader_tsn),
1402                    "TSN should remain registered until reader is dropped"
1403                );
1404            }
1405
1406            // Expected are strictly from the chosen 'from' position
1407            let num_to_skip = i.max(1) - 1;
1408            let expected: Vec<(Position, EventRecord)> =
1409                appended.clone().into_iter().skip(num_to_skip).collect();
1410            assert_eq!(expected.len(), scanned.len());
1411            // println!("Get expected number: {}", expected.len());
1412            for (i, exp) in expected.iter().enumerate() {
1413                assert_eq!(exp.0, scanned[i].0);
1414                assert_eq!(exp.1, scanned[i].1);
1415            }
1416
1417            // Ensure we did not accumulate pages in the iterator cache
1418            assert!(
1419                events_iterator.page_cache.is_empty(),
1420                "EventIterator page_cache should be empty after filtered scan"
1421            );
1422
1423            // Drop the reader and ensure the reader tsn is removed
1424            drop(reader);
1425            {
1426                // Check registration state directly on the reader-TSN registry
1427                assert!(
1428                    !db.reader_tsns.contains_tsn(reader_tsn),
1429                    "TSN should be removed after reader is dropped"
1430                );
1431            }
1432        }
1433    }
1434
1435    #[test]
1436    #[serial]
1437    fn test_read_events_from_backwards() {
1438        // Setup a temporary database
1439        let (_temp_dir, db) = construct_db(128);
1440
1441        let mut has_split_internal = false;
1442        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1443
1444        // Insert events until we split a root internal node
1445        while !has_split_internal {
1446            // Start a writer
1447            let mut writer = db.writer().unwrap();
1448
1449            // Issue a new position and create a record
1450            let position = writer.issue_position();
1451            let record = EventRecord {
1452                event_type: "UserCreated".to_string(),
1453                data: (0..8).map(|_| random::<u8>()).collect(),
1454                tags: vec!["users".to_string(), "creation".to_string()],
1455                uuid: None,
1456                metadata: Vec::new(),
1457            };
1458            appended.push((position, record.clone()));
1459
1460            // Append the event
1461            event_tree_append(&db, &mut writer, record, position).unwrap();
1462
1463            // Check if the root is an internal node
1464            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1465            match &root_page.node {
1466                Node::EventInternal(root_node) => {
1467                    // Check if the first child is an internal node
1468                    if !root_node.child_ids.is_empty() {
1469                        let child_id = root_node.child_ids[0];
1470                        if let Some(child_page) = writer.dirty.get(&child_id) {
1471                            match &child_page.node {
1472                                Node::EventInternal(_) => {
1473                                    has_split_internal = true;
1474                                }
1475                                _ => {}
1476                            }
1477                        }
1478                    }
1479                }
1480                _ => {}
1481            }
1482            db.commit(&mut writer).unwrap();
1483        }
1484
1485        let appended_len = appended.len();
1486        // println!("Appended number: {}", appended_len);
1487
1488        // Iterate through various 'from' positions.
1489        for i in 0..appended_len + 2 {
1490            let from = Position(i as u64);
1491
1492            // Start a reader
1493            let reader = db.reader().unwrap();
1494            let events_tree_root_id = reader.events_tree_root_id;
1495            let reader_tsn = reader.tsn;
1496            let dirty = HashMap::new();
1497            let mut events_iterator =
1498                EventIterator::new(&db, &dirty, events_tree_root_id, Some(from), true);
1499
1500            // Ensure the reader's tsn is registered while the iterator is alive
1501            {
1502                // Check registration state directly on the reader-TSN registry
1503                assert!(
1504                    db.reader_tsns.contains_tsn(reader_tsn),
1505                    "TSN should remain registered until reader is dropped"
1506                );
1507            }
1508
1509            // Progressively iterate over events using batches
1510            let mut scanned: Vec<(Position, EventRecord)> = Vec::new();
1511            loop {
1512                let batch = events_iterator.next_batch(3, None).unwrap();
1513                if batch.is_empty() {
1514                    break;
1515                }
1516                scanned.extend(batch);
1517
1518                // The reader should remain registered throughout iteration
1519                // Check registration state directly on the reader-TSN registry
1520                assert!(
1521                    db.reader_tsns.contains_tsn(reader_tsn),
1522                    "TSN should remain registered until reader is dropped"
1523                );
1524            }
1525
1526            // Expected are strictly from the chosen 'from' position
1527            let mut expected: Vec<(Position, EventRecord)> = appended.clone();
1528            expected.truncate(i);
1529            expected.reverse();
1530            assert_eq!(expected.len(), scanned.len());
1531            // println!("Get expected number: {}", expected.len());
1532            for (i, exp) in expected.iter().enumerate() {
1533                assert_eq!(exp.0, scanned[i].0);
1534                assert_eq!(exp.1, scanned[i].1);
1535            }
1536
1537            // Ensure we did not accumulate pages in the iterator cache
1538            assert!(
1539                events_iterator.page_cache.is_empty(),
1540                "EventIterator page_cache should be empty after filtered scan"
1541            );
1542
1543            // Drop the reader and ensure the reader tsn is removed
1544            drop(reader);
1545            {
1546                // Check registration state directly on the reader-TSN registry
1547                assert!(
1548                    !db.reader_tsns.contains_tsn(reader_tsn),
1549                    "TSN should be removed after reader is dropped"
1550                );
1551            }
1552        }
1553    }
1554
1555    #[test]
1556    #[serial]
1557    fn test_large_event_data_exact_page_size() {
1558        let (_tmp, db) = construct_db(512);
1559        // Append one large event with data exactly equal to page size
1560        let mut writer = db.writer().unwrap();
1561        let pos = writer.issue_position();
1562        let data = vec![0xAB; 512];
1563        let event = EventRecord {
1564            event_type: "Big".into(),
1565            data: data.clone(),
1566            tags: vec![],
1567            uuid: None,
1568            metadata: Vec::new(),
1569        };
1570        event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1571        db.commit(&mut writer).unwrap();
1572
1573        // Lookup should return identical payload
1574        let reader = db.reader().unwrap();
1575        let dirty = HashMap::new();
1576        let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1577        assert_eq!(event, got);
1578
1579        // Ensure an overflow page is used for storage
1580        let header_page = db.get_latest_header_page().unwrap();
1581        let header = header_page.as_header_node().unwrap();
1582        let root = db.read_page(header.events_tree_root_id).unwrap();
1583        match &root.node {
1584            Node::EventInternal(internal) => {
1585                // Our key should be in the last child
1586                let leaf_id = *internal.child_ids.last().unwrap();
1587                let leaf_page = db.read_page(leaf_id).unwrap();
1588                match &leaf_page.node {
1589                    Node::EventLeaf(leaf) => match &leaf.values[0] {
1590                        EventValue::Overflow { data_len, .. } => {
1591                            assert_eq!(*data_len as usize, data.len())
1592                        }
1593                        _ => panic!("Expected Overflow for large event"),
1594                    },
1595                    _ => panic!("Expected EventLeaf child"),
1596                }
1597            }
1598            Node::EventLeaf(leaf) => match &leaf.values[0] {
1599                EventValue::Overflow { data_len, .. } => {
1600                    assert_eq!(*data_len as usize, data.len())
1601                }
1602                _ => panic!("Expected Overflow for large event"),
1603            },
1604            _ => panic!("Unexpected root node type"),
1605        }
1606    }
1607
1608    #[test]
1609    #[serial]
1610    fn test_large_event_data_four_times_page_size() {
1611        let (_tmp, db) = construct_db(512);
1612        let mut writer = db.writer().unwrap();
1613        let pos = writer.issue_position();
1614        let data = vec![0xCD; 512 * 4];
1615        let event = EventRecord {
1616            event_type: "Bigger".into(),
1617            data: data.clone(),
1618            tags: vec![],
1619            uuid: None,
1620            metadata: Vec::new(),
1621        };
1622        event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1623        db.commit(&mut writer).unwrap();
1624
1625        // Lookup
1626        let reader = db.reader().unwrap();
1627        let dirty = HashMap::new();
1628        let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1629        assert_eq!(event, got);
1630
1631        // Ensure overflow in leaf
1632        let header_page = db.get_latest_header_page().unwrap();
1633        let header = header_page.as_header_node().unwrap();
1634        let root = db.read_page(header.events_tree_root_id).unwrap();
1635        let check_leaf = |leaf: &EventLeafNode| match &leaf.values[0] {
1636            EventValue::Overflow { data_len, .. } => {
1637                assert_eq!(*data_len as usize, data.len())
1638            }
1639            _ => panic!("Expected Overflow for very large event"),
1640        };
1641        match &root.node {
1642            Node::EventInternal(internal) => {
1643                let leaf_id = *internal.child_ids.last().unwrap();
1644                let leaf_page = db.read_page(leaf_id).unwrap();
1645                match &leaf_page.node {
1646                    Node::EventLeaf(leaf) => check_leaf(&leaf),
1647                    _ => panic!("Expected leaf"),
1648                }
1649            }
1650            Node::EventLeaf(leaf) => check_leaf(&leaf),
1651            _ => panic!("Unexpected root node type"),
1652        }
1653    }
1654
1655    #[test]
1656    #[serial]
1657    fn test_inline_event_metadata_roundtrip() {
1658        let (_tmp, db) = construct_db(4096);
1659        let mut writer = db.writer().unwrap();
1660        let pos = writer.issue_position();
1661
1662        let mut metadata = Vec::new();
1663        metadata.push(("source".to_string(), "web".to_string()));
1664        metadata.push(("correlation_id".to_string(), "abc-123".to_string()));
1665
1666        let event = EventRecord {
1667            event_type: "SmallWithMetadata".into(),
1668            data: vec![1, 2, 3, 4],
1669            tags: vec!["t".into()],
1670            uuid: None,
1671            metadata: metadata.clone(),
1672        };
1673        event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1674        db.commit(&mut writer).unwrap();
1675
1676        // Read back through the materialization path and confirm metadata survives.
1677        let reader = db.reader().unwrap();
1678        let dirty = HashMap::new();
1679        let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1680        assert_eq!(event, got);
1681        assert_eq!(metadata, got.metadata);
1682
1683        // Confirm the value was stored inline (not in an overflow chain).
1684        let header_page = db.get_latest_header_page().unwrap();
1685        let header = header_page.as_header_node().unwrap();
1686        let root = db.read_page(header.events_tree_root_id).unwrap();
1687        match &root.node {
1688            Node::EventLeaf(leaf) => match &leaf.values[0] {
1689                EventValue::Inline(rec) => assert_eq!(metadata, rec.metadata),
1690                _ => panic!("Expected Inline value"),
1691            },
1692            _ => panic!("Expected EventLeaf root for a single small event"),
1693        }
1694    }
1695
1696    #[test]
1697    #[serial]
1698    fn test_overflow_event_metadata_roundtrip() {
1699        let (_tmp, db) = construct_db(512);
1700        let mut writer = db.writer().unwrap();
1701        let pos = writer.issue_position();
1702
1703        let mut metadata = Vec::new();
1704        metadata.push(("source".to_string(), "bulk-import".to_string()));
1705        metadata.push(("schema".to_string(), "v3".to_string()));
1706
1707        // Data large enough to spill into overflow pages.
1708        let data = vec![0xAB; 512 * 4];
1709        let event = EventRecord {
1710            event_type: "BigWithMetadata".into(),
1711            data: data.clone(),
1712            tags: vec!["t".into()],
1713            uuid: Some(uuid::Uuid::new_v4()),
1714            metadata: metadata.clone(),
1715        };
1716        event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1717        db.commit(&mut writer).unwrap();
1718
1719        // Read back through the materialization path and confirm both data and
1720        // metadata survive the split/overflow chain.
1721        let reader = db.reader().unwrap();
1722        let dirty = HashMap::new();
1723        let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1724        assert_eq!(event, got);
1725        assert_eq!(data, got.data);
1726        assert_eq!(metadata, got.metadata);
1727
1728        // Confirm the value really was stored as an overflow value with a
1729        // non-zero metadata length.
1730        let header_page = db.get_latest_header_page().unwrap();
1731        let header = header_page.as_header_node().unwrap();
1732        let root = db.read_page(header.events_tree_root_id).unwrap();
1733        let check_leaf = |leaf: &EventLeafNode| match &leaf.values[0] {
1734            EventValue::Overflow {
1735                data_len,
1736                metadata_len,
1737                ..
1738            } => {
1739                assert_eq!(*data_len as usize, data.len());
1740                assert!(*metadata_len > 0, "expected non-zero metadata_len");
1741            }
1742            _ => panic!("Expected Overflow value for large event"),
1743        };
1744        match &root.node {
1745            Node::EventInternal(internal) => {
1746                let leaf_id = *internal.child_ids.last().unwrap();
1747                let leaf_page = db.read_page(leaf_id).unwrap();
1748                match &leaf_page.node {
1749                    Node::EventLeaf(leaf) => check_leaf(leaf),
1750                    _ => panic!("Expected EventLeaf child"),
1751                }
1752            }
1753            Node::EventLeaf(leaf) => check_leaf(leaf),
1754            _ => panic!("Unexpected root node type"),
1755        }
1756    }
1757
1758    #[test]
1759    #[serial]
1760    fn test_direct_overflow_event_metadata_roundtrip() {
1761        // Data larger than u16::MAX takes the direct inline->overflow path in
1762        // event_tree_append (rather than the leaf-split conversion path).
1763        let (_tmp, db) = construct_db(4096);
1764        let mut writer = db.writer().unwrap();
1765        let pos = writer.issue_position();
1766
1767        let mut metadata = Vec::new();
1768        metadata.push(("origin".to_string(), "direct-overflow".to_string()));
1769
1770        let data = vec![0xCD; (u16::MAX as usize) + 1024];
1771        let event = EventRecord {
1772            event_type: "HugeWithMetadata".into(),
1773            data: data.clone(),
1774            tags: vec!["t".into()],
1775            uuid: None,
1776            metadata: metadata.clone(),
1777        };
1778        event_tree_append(&db, &mut writer, event.clone(), pos).unwrap();
1779        db.commit(&mut writer).unwrap();
1780
1781        let reader = db.reader().unwrap();
1782        let dirty = HashMap::new();
1783        let got = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1784        assert_eq!(event, got);
1785        assert_eq!(data, got.data);
1786        assert_eq!(metadata, got.metadata);
1787    }
1788
1789    #[test]
1790    #[serial]
1791    fn test_append_rejects_oversized_metadata_key() {
1792        use crate::events_tree_nodes::MAX_METADATA_ENTRY_LEN;
1793
1794        let (_tmp, db) = construct_db(4096);
1795        let mut writer = db.writer().unwrap();
1796        let pos = writer.issue_position();
1797
1798        let mut metadata = Vec::new();
1799        metadata.push(("x".repeat(MAX_METADATA_ENTRY_LEN + 1), "source".to_string()));
1800        let event = EventRecord {
1801            event_type: "TooBigMetadata".into(),
1802            data: vec![1, 2, 3],
1803            tags: vec!["t".into()],
1804            uuid: None,
1805            metadata,
1806        };
1807
1808        // Append must fail cleanly with InvalidArgument rather than corrupting
1809        // the encoding.
1810        match event_tree_append(&db, &mut writer, event, pos) {
1811            Err(DcbError::InvalidArgument(_)) => {}
1812            other => panic!("Expected InvalidArgument, got {other:?}"),
1813        }
1814
1815        // Nothing should have been persisted: after committing, the position
1816        // must not resolve to a stored event.
1817        db.commit(&mut writer).unwrap();
1818        let reader = db.reader().unwrap();
1819        let dirty = HashMap::new();
1820        assert!(
1821            event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).is_err(),
1822            "rejected event must not be retrievable"
1823        );
1824    }
1825
1826    #[test]
1827    #[serial]
1828    fn test_append_rejects_oversized_metadata_value() {
1829        use crate::events_tree_nodes::MAX_METADATA_ENTRY_LEN;
1830
1831        let (_tmp, db) = construct_db(4096);
1832        let mut writer = db.writer().unwrap();
1833        let pos = writer.issue_position();
1834
1835        let mut metadata = Vec::new();
1836        metadata.push(("source".to_string(), "x".repeat(MAX_METADATA_ENTRY_LEN + 1)));
1837        let event = EventRecord {
1838            event_type: "TooBigMetadata".into(),
1839            data: vec![1, 2, 3],
1840            tags: vec!["t".into()],
1841            uuid: None,
1842            metadata,
1843        };
1844
1845        // Append must fail cleanly with InvalidArgument rather than corrupting
1846        // the encoding.
1847        match event_tree_append(&db, &mut writer, event, pos) {
1848            Err(DcbError::InvalidArgument(_)) => {}
1849            other => panic!("Expected InvalidArgument, got {other:?}"),
1850        }
1851
1852        // Nothing should have been persisted: after committing, the position
1853        // must not resolve to a stored event.
1854        db.commit(&mut writer).unwrap();
1855        let reader = db.reader().unwrap();
1856        let dirty = HashMap::new();
1857        assert!(
1858            event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).is_err(),
1859            "rejected event must not be retrievable"
1860        );
1861    }
1862
1863    #[test]
1864    #[serial]
1865    fn test_append_rejects_oversized_event_type() {
1866        use crate::events_tree_nodes::MAX_EVENT_TYPE_LEN;
1867
1868        let (_tmp, db) = construct_db(4096);
1869        let mut writer = db.writer().unwrap();
1870        let pos = writer.issue_position();
1871
1872        let event = EventRecord {
1873            event_type: "x".repeat(MAX_EVENT_TYPE_LEN + 1),
1874            data: vec![1, 2, 3],
1875            tags: vec!["t".into()],
1876            uuid: None,
1877            metadata: Vec::new(),
1878        };
1879
1880        match event_tree_append(&db, &mut writer, event, pos) {
1881            Err(DcbError::InvalidArgument(_)) => {}
1882            other => panic!("Expected InvalidArgument, got {other:?}"),
1883        }
1884    }
1885
1886    #[test]
1887    #[serial]
1888    fn test_append_rejects_oversized_tag() {
1889        use crate::events_tree_nodes::MAX_TAG_LEN;
1890
1891        let (_tmp, db) = construct_db(4096);
1892        let mut writer = db.writer().unwrap();
1893        let pos = writer.issue_position();
1894
1895        let event = EventRecord {
1896            event_type: "Type".into(),
1897            data: vec![1, 2, 3],
1898            tags: vec!["x".repeat(MAX_TAG_LEN + 1)],
1899            uuid: None,
1900            metadata: Vec::new(),
1901        };
1902
1903        match event_tree_append(&db, &mut writer, event, pos) {
1904            Err(DcbError::InvalidArgument(_)) => {}
1905            other => panic!("Expected InvalidArgument, got {other:?}"),
1906        }
1907    }
1908    // fn benchmark_append_and_lookup_varied_sizes() {
1909    //     // Benchmark-like test; prints durations for different sizes. Run with:
1910    //     // cargo test --lib mvcc_event_tree::tests::benchmark_append_and_lookup_varied_sizes -- --nocapture
1911    //     let sizes: [usize; 7] = [1, 10, 100, 1_000, 5_000, 10_000, 50_000];
1912    //     for &size in &sizes {
1913    //         let (_tmp, db) = construct_db(4096);
1914    //
1915    //         // Append phase
1916    //         let mut writer = db.writer().unwrap();
1917    //         let mut positions: Vec<Position> = Vec::with_capacity(size);
1918    //         let start_append = std::time::Instant::now();
1919    //         for n in 0..(size as u64) {
1920    //             let pos = writer.issue_position();
1921    //             let event = EventRecord {
1922    //                 event_type: "E".to_string(),
1923    //                 data: Vec::new(),
1924    //                 tags: Vec::new(),
1925    //             };
1926    //             std::hint::black_box(n);
1927    //             std::hint::black_box(&event);
1928    //             std::hint::black_box(pos);
1929    //             event_tree_append(&db, &mut writer, event, pos).unwrap();
1930    //             positions.push(pos);
1931    //         }
1932    //         let append_elapsed = start_append.elapsed();
1933    //         let start_commit = std::time::Instant::now();
1934    //         db.commit(&mut writer).unwrap();
1935    //         let commit_elapsed = start_commit.elapsed();
1936    //
1937    //         // Lookup phase
1938    //         let reader = db.reader().unwrap();
1939    //         let start_lookup = std::time::Instant::now();
1940    //         let dirty = HashMap::new();
1941    //         for &pos in &positions {
1942    //             let rec = event_tree_lookup(&db, &dirty, reader.events_tree_root_id, pos).unwrap();
1943    //             std::hint::black_box(&rec);
1944    //         }
1945    //         let lookup_elapsed = start_lookup.elapsed();
1946    //
1947    //         let append_avg_us = (append_elapsed.as_secs_f64() * 1_000_000.0) / (size as f64);
1948    //         let commit_avg_us = commit_elapsed.as_secs_f64() * 1_000_000.0;
1949    //         let lookup_avg_us = (lookup_elapsed.as_secs_f64() * 1_000_000.0) / (size as f64);
1950    //
1951    //         println!(
1952    //             "mvcc_event_tree benchmark: size={}, append_us_per_call={:.3}, commit_us={:.3}, lookup_us_per_call={:.3}",
1953    //             size, append_avg_us, commit_avg_us, lookup_avg_us
1954    //         );
1955    //     }
1956    // }
1957
1958    #[test]
1959    #[serial]
1960    fn test_read_events_all_backwards_with_internal_nodes() {
1961        // Setup a temporary database with a small page size to force internal nodes quickly
1962        let (_temp_dir, db) = construct_db(128);
1963
1964        let mut has_split_internal = false;
1965        let mut appended: Vec<(Position, EventRecord)> = Vec::new();
1966
1967        // 1. Insert events until the root is an internal node that points to other internal nodes
1968        while !has_split_internal {
1969            let mut writer = db.writer().unwrap();
1970            let position = writer.issue_position();
1971            let record = EventRecord {
1972                event_type: "TestEvent".to_string(),
1973                data: vec![1, 2, 3, 4],
1974                tags: vec!["test".to_string()],
1975                uuid: None,
1976                metadata: Vec::new(),
1977            };
1978            appended.push((position, record.clone()));
1979            event_tree_append(&db, &mut writer, record, position).unwrap();
1980
1981            let root_page = writer.dirty.get(&writer.events_tree_root_id).unwrap();
1982            if let Node::EventInternal(root_node) = &root_page.node {
1983                if !root_node.child_ids.is_empty() {
1984                    let child_id = root_node.child_ids[0];
1985                    if let Some(child_page) = writer.dirty.get(&child_id) {
1986                        if let Node::EventInternal(_) = &child_page.node {
1987                            has_split_internal = true;
1988                        }
1989                    }
1990                }
1991            }
1992            db.commit(&mut writer).unwrap();
1993        }
1994
1995        // 2. Iterate backwards from the end (start = None)
1996        let reader = db.reader().unwrap();
1997        let dirty = HashMap::new();
1998        let mut events_iterator = EventIterator::new(
1999            &db,
2000            &dirty,
2001            reader.events_tree_root_id,
2002            None, // This triggers the branch for line 561
2003            true, // Backwards
2004        );
2005
2006        let mut scanned = Vec::new();
2007        loop {
2008            let batch = events_iterator.next_batch(10, None).unwrap();
2009            if batch.is_empty() {
2010                break;
2011            }
2012            scanned.extend(batch);
2013        }
2014
2015        // 3. Verify all events are returned in reverse order
2016        let mut expected = appended.clone();
2017        expected.reverse();
2018
2019        assert_eq!(scanned.len(), expected.len());
2020        for i in 0..expected.len() {
2021            assert_eq!(scanned[i].0, expected[i].0);
2022        }
2023    }
2024
2025    #[test]
2026    #[serial]
2027    fn test_append_event_with_tag_larger_than_page_size() {
2028        // This test was created when I realized that even though
2029        // page data and metadata overflows, it may be that the
2030        // tags or the type is so large that it can't possibly
2031        // fit in one page... the tree should detect this.
2032        use crate::events_tree_nodes::MAX_TAG_LEN;
2033
2034        // Small page size to make it easy to exceed
2035        let page_size = 512;
2036        let (_tmp, db) = construct_db(page_size);
2037        let mut writer = db.writer().unwrap();
2038        let pos = writer.issue_position();
2039
2040        // Tag larger than page size but less than MAX_TAG_LEN
2041        // Let's make it 600 bytes.
2042        let large_tag = "t".repeat(600);
2043        assert!(large_tag.len() > page_size);
2044        assert!(large_tag.len() < MAX_TAG_LEN);
2045
2046        let event = EventRecord {
2047            event_type: "Type".into(),
2048            data: vec![1, 2, 3],
2049            tags: vec![large_tag],
2050            uuid: None,
2051            metadata: Vec::new(),
2052        };
2053
2054        match event_tree_append(&db, &mut writer, event, pos) {
2055            Err(DcbError::InvalidArgument(message)) => {
2056                assert!(
2057                    message.contains("event too large for page size"),
2058                    "unexpected InvalidArgument message: {message}"
2059                );
2060            }
2061            other => panic!("Expected DatabaseCorrupted, got {other:?}"),
2062        }
2063    }
2064}