umadb_core/
db.rs

1use std::path::Path;
2
3use crate::common::{PageID, Position};
4use crate::events_tree::{EventIterator, event_tree_append, event_tree_lookup};
5use crate::events_tree_nodes::EventRecord;
6use crate::mvcc::{Mvcc, Writer};
7use crate::page::Page;
8use crate::tags_tree::{TagsTreeIterator, tags_tree_insert};
9use crate::tags_tree_nodes::TagHash;
10use itertools::Itertools;
11use std::collections::{HashMap, HashSet, VecDeque};
12use std::sync::Arc;
13use umadb_dcb::{
14    DCBAppendCondition, DCBError, DCBEvent, DCBEventStoreSync, DCBQuery, DCBReadResponseSync,
15    DCBResult, DCBSequencedEvent,
16};
17use uuid::Uuid;
18
19pub static DEFAULT_PAGE_SIZE: usize = 4096;
20pub const DEFAULT_DB_FILENAME: &str = "uma.db";
21
22/// EventStore implementing the DCBEventStoreSync interface
23pub struct UmaDB {
24    mvcc: Arc<Mvcc>,
25}
26
27impl UmaDB {
28    /// Create a new EventStore at the given directory or file path.
29    /// If a directory path is provided, a file named "uma.db" will be used inside it.
30    pub fn new<P: AsRef<Path>>(path: P) -> DCBResult<Self> {
31        let p = path.as_ref();
32        let file_path = if p.is_dir() {
33            p.join(DEFAULT_DB_FILENAME)
34        } else {
35            p.to_path_buf()
36        };
37        let mvcc = Mvcc::new(&file_path, DEFAULT_PAGE_SIZE, false)?;
38        Ok(Self {
39            mvcc: Arc::new(mvcc),
40        })
41    }
42
43    pub fn from_arc(mvcc: Arc<Mvcc>) -> Self {
44        Self { mvcc }
45    }
46
47    /// Appends a batch of (events, condition) using a single writer/transaction.
48    /// For each item, behaves like append():
49    /// - If condition is Some and matches any events (considering uncommitted writes), returns Err(IntegrityError) for that item and continues.
50    /// - If events is empty, returns Ok(0) for that item and continues.
51    /// - Otherwise performs unconditional append and records Ok(last_position) for that item.
52    ///
53    /// At the end, commits the writer once. If commit fails, returns the commit error and discards per-item results.
54    pub fn append_batch(
55        &self,
56        items: Vec<(Vec<DCBEvent>, Option<DCBAppendCondition>)>,
57        force_sequential_read: bool,
58    ) -> DCBResult<Vec<DCBResult<u64>>> {
59        // println!("Processing batch of {} items", items.len());
60
61        let mvcc = &self.mvcc;
62        let mut writer = mvcc.writer()?;
63        let mut results: Vec<DCBResult<u64>> = Vec::with_capacity(items.len());
64
65        for (events, condition) in items.into_iter() {
66            // Check condition using read_conditional (limit 1), starting after the provided position
67            if let Some(cond) = condition {
68                let from = cond.after.map(|after| Position(after + 1));
69                let read_result1 = read_conditional(
70                    mvcc,
71                    &writer.dirty,
72                    writer.events_tree_root_id,
73                    writer.tags_tree_root_id,
74                    cond.fail_if_events_match.clone(),
75                    from,
76                    false,
77                    Some(1),
78                    force_sequential_read,
79                );
80                match read_result1 {
81                    Ok(found_vec) => {
82                        // Read didn't error...
83                        if let Some(matched) = found_vec.first() {
84                            // Found one event... consider if the request is idempotent...
85                            match is_request_idempotent(
86                                mvcc,
87                                &writer.dirty,
88                                writer.events_tree_root_id,
89                                writer.tags_tree_root_id,
90                                &events,
91                                cond.fail_if_events_match.clone(),
92                                from,
93                            ) {
94                                Ok(Some(last_recorded_position)) => {
95                                    results.push(Ok(last_recorded_position));
96                                }
97                                Ok(None) => {
98                                    // Propagate an integrity error for this item but continue with others
99                                    let msg = format!(
100                                        "condition: {:?} matched: {:?}, ",
101                                        cond.clone(),
102                                        matched,
103                                    );
104                                    results.push(Err(DCBError::IntegrityError(msg)));
105                                }
106                                Err(err) => {
107                                    // Propagate the error for this item but continue with others
108                                    results.push(Err(err));
109                                }
110                            }
111                            continue;
112                        }
113                    }
114                    Err(e) => {
115                        // Propagate the read error for this item but continue with others
116                        results.push(Err(e));
117                        continue;
118                    }
119                }
120            }
121
122            if events.is_empty() {
123                results.push(Ok(0));
124                continue;
125            }
126
127            // Append unconditionally
128            match unconditional_append(mvcc, &mut writer, events) {
129                Ok(last) => results.push(Ok(last)),
130                Err(e) => {
131                    // Record error for this item and continue
132                    results.push(Err(e));
133                }
134            }
135        }
136
137        // Single commit at the end of the batch
138        mvcc.commit(&mut writer)?;
139        Ok(results)
140    }
141}
142
143impl DCBEventStoreSync for UmaDB {
144    fn read(
145        &self,
146        query: Option<DCBQuery>,
147        start: Option<u64>,
148        backwards: bool,
149        limit: Option<u32>,
150        _subscribe: bool,
151    ) -> DCBResult<Box<dyn DCBReadResponseSync + 'static>> {
152        let mvcc = &self.mvcc;
153        let reader = mvcc.reader()?;
154
155        // Compute last committed position for unlimited head
156        let last_committed_position = reader.next_position.0.saturating_sub(1);
157
158        // Build query and after
159        let q = query.unwrap_or(DCBQuery { items: vec![] });
160        let from = start.map(Position);
161
162        // Delegate to read_conditional
163        let events = read_conditional(
164            mvcc,
165            &HashMap::new(),
166            reader.events_tree_root_id,
167            reader.tags_tree_root_id,
168            q,
169            from,
170            backwards,
171            limit,
172            false,
173        )?;
174
175        // Compute head according to semantics
176        let head = if limit.is_none() {
177            if last_committed_position == 0 {
178                None
179            } else {
180                Some(last_committed_position)
181            }
182        } else {
183            events.last().map(|e| e.position)
184        };
185
186        Ok(Box::new(ReadResponse {
187            events: VecDeque::from(events),
188            head,
189        }))
190    }
191
192    fn head(&self) -> DCBResult<Option<u64>> {
193        let db = &self.mvcc;
194        let (_, header) = db.get_latest_header()?;
195        let last = header.next_position.0.saturating_sub(1);
196        if last == 0 { Ok(None) } else { Ok(Some(last)) }
197    }
198
199    fn append(
200        &self,
201        events: Vec<DCBEvent>,
202        condition: Option<DCBAppendCondition>,
203    ) -> DCBResult<u64> {
204        // Preserve existing fast-path: if no events, do nothing and return 0 (avoid opening/committing a writer)
205        if events.is_empty() {
206            return Ok(0);
207        }
208        // Delegate to append_batch with a single item to reuse unified batching logic
209        let mut results = self.append_batch(vec![(events, condition)], false)?;
210        debug_assert_eq!(results.len(), 1);
211        match results.remove(0) {
212            Ok(pos) => Ok(pos),
213            Err(e) => Err(e),
214        }
215    }
216}
217
218struct ReadResponse {
219    events: VecDeque<DCBSequencedEvent>,
220    head: Option<u64>,
221}
222
223impl Iterator for ReadResponse {
224    type Item = DCBResult<DCBSequencedEvent>;
225    fn next(&mut self) -> Option<Self::Item> {
226        self.events.pop_front().map(Ok)
227    }
228}
229
230impl DCBReadResponseSync for ReadResponse {
231    fn head(&mut self) -> DCBResult<Option<u64>> {
232        Ok(self.head)
233    }
234    fn collect_with_head(&mut self) -> DCBResult<(Vec<DCBSequencedEvent>, Option<u64>)> {
235        let events = self.events.drain(..).collect();
236        Ok((events, self.head))
237    }
238    fn next_batch(&mut self) -> DCBResult<Vec<DCBSequencedEvent>> {
239        let batch = self.events.drain(..).collect();
240        Ok(batch)
241    }
242}
243
244/// Append events unconditionally to the database.
245///
246/// For each event, this will:
247/// - issue a position from the writer
248/// - append an EventRecord to the event tree
249/// - insert the position for each tag into the tags tree
250///
251/// Caller is responsible for committing the writer.
252pub fn unconditional_append(
253    mvcc: &Mvcc,
254    writer: &mut Writer,
255    events: Vec<DCBEvent>,
256) -> DCBResult<u64> {
257    let mut last_pos_u64: u64 = 0;
258
259    for ev in events.into_iter() {
260        let position = writer.issue_position();
261        last_pos_u64 = position.0;
262        // Index tags before moving an event record into event_tree_append
263        for tag in ev.tags.iter() {
264            let tag_hash: TagHash = tag_to_hash(tag);
265            tags_tree_insert(mvcc, writer, tag_hash, position)?;
266        }
267        let record = EventRecord {
268            event_type: ev.event_type,
269            data: ev.data,
270            tags: ev.tags,
271            uuid: ev.uuid,
272        };
273        event_tree_append(mvcc, writer, record, position)?;
274    }
275
276    Ok(last_pos_u64)
277}
278
279/// Read events using the tags index by merging per-tag iterators, grouping by position,
280/// filtering by tag and type matches, and then looking up the event record.
281pub fn read_conditional(
282    mvcc: &Mvcc,
283    dirty: &HashMap<PageID, Page>,
284    events_tree_root_id: PageID,
285    tags_tree_root_id: PageID,
286    query: DCBQuery,
287    start: Option<Position>,
288    backwards: bool,
289    limit: Option<u32>,
290    force_sequential_read: bool,
291) -> DCBResult<Vec<DCBSequencedEvent>> {
292    const SCAN_BATCH_SIZE: u32 = 256;
293    // Special case: explicit zero limit
294    if let Some(0) = limit {
295        return Ok(Vec::new());
296    }
297
298    // If no items, return all events with after/limit respected via sequential scan
299    if query.items.is_empty() {
300        let mut iter = EventIterator::new(mvcc, dirty, events_tree_root_id, start, backwards);
301        let mut out: Vec<DCBSequencedEvent> = Vec::new();
302        'outer_all: loop {
303            let batch = iter.next_batch(limit.unwrap_or(SCAN_BATCH_SIZE))?;
304            if batch.is_empty() {
305                break;
306            }
307            for (pos, rec) in batch.into_iter() {
308                out.push(DCBSequencedEvent {
309                    position: pos.0,
310                    event: DCBEvent {
311                        event_type: rec.event_type,
312                        data: rec.data,
313                        tags: rec.tags,
314                        uuid: rec.uuid,
315                    },
316                });
317                if let Some(lim) = limit
318                    && out.len() >= lim as usize
319                {
320                    break 'outer_all;
321                }
322            }
323        }
324        return Ok(out);
325    }
326
327    // All query items must have at least one tag to use the tag index path.
328    let all_items_have_tags = query.items.iter().all(|it| !it.tags.is_empty());
329    if !all_items_have_tags || force_sequential_read {
330        // Fallback: sequentially scan all events and apply the same matching logic
331        let mut iter = EventIterator::new(mvcc, dirty, events_tree_root_id, start, backwards);
332        let mut out: Vec<DCBSequencedEvent> = Vec::new();
333        let matches_item = |rec: &EventRecord| -> bool {
334            for item in &query.items {
335                let type_ok =
336                    item.types.is_empty() || item.types.iter().any(|t| t == &rec.event_type);
337                if !type_ok {
338                    continue;
339                }
340                let tags_ok = item.tags.iter().all(|t| rec.tags.iter().any(|et| et == t));
341                if type_ok && tags_ok {
342                    return true;
343                }
344            }
345            false
346        };
347        'outer_fallback: loop {
348            let batch = iter.next_batch(SCAN_BATCH_SIZE)?;
349            if batch.is_empty() {
350                break;
351            }
352            for (pos, rec) in batch.into_iter() {
353                if matches_item(&rec) {
354                    out.push(DCBSequencedEvent {
355                        position: pos.0,
356                        event: DCBEvent {
357                            event_type: rec.event_type,
358                            data: rec.data,
359                            tags: rec.tags,
360                            uuid: rec.uuid,
361                        },
362                    });
363                    if let Some(lim) = limit
364                        && out.len() >= lim as usize
365                    {
366                        break 'outer_fallback;
367                    }
368                }
369            }
370        }
371        return Ok(out);
372    }
373
374    // Invert query: tag -> list of query item indices that require this tag
375    let mut tag_qiis: HashMap<String, Vec<usize>> = HashMap::with_capacity(query.items.len() * 2);
376    let mut qi_tags: Vec<HashSet<String>> = Vec::with_capacity(query.items.len());
377
378    for (qiid, item) in query.items.iter().enumerate() {
379        qi_tags.push(item.tags.iter().cloned().collect());
380        for tag in &item.tags {
381            tag_qiis.entry(tag.clone()).or_default().push(qiid);
382        }
383    }
384
385    // Prepare per-tag iterators yielding (position, tag, qiids)
386    struct PositionTagQiidIterator<I>
387    where
388        I: Iterator<Item = Position>,
389    {
390        inner: I,
391        tag: String,
392        qiids: Vec<usize>,
393    }
394    impl<I> PositionTagQiidIterator<I>
395    where
396        I: Iterator<Item = Position>,
397    {
398        fn new(inner: I, tag: String, qiids: Vec<usize>) -> Self {
399            Self { inner, tag, qiids }
400        }
401    }
402    impl<I> Iterator for PositionTagQiidIterator<I>
403    where
404        I: Iterator<Item = Position>,
405    {
406        type Item = (Position, String, Vec<usize>);
407        fn next(&mut self) -> Option<Self::Item> {
408            self.inner
409                .next()
410                .map(|p| (p, self.tag.clone(), self.qiids.clone()))
411        }
412    }
413
414    let mut tag_iters: Vec<PositionTagQiidIterator<_>> = Vec::new();
415    for (tag, qiids) in tag_qiis.iter() {
416        let tag_hash: TagHash = tag_to_hash(tag);
417        let positions_iter =
418            TagsTreeIterator::new(mvcc, dirty, tags_tree_root_id, tag_hash, start, backwards); // yields positions for tag
419        tag_iters.push(PositionTagQiidIterator::new(
420            positions_iter,
421            tag.clone(),
422            qiids.clone(),
423        ));
424    }
425
426    // Merge iterators ordered by position
427    let merged = tag_iters
428        .into_iter()
429        .kmerge_by(|a, b| if !backwards { a.0 < b.0 } else { a.0 > b.0 });
430
431    // Group by position, collecting tags and qiids
432    struct GroupByPositionIterator<I>
433    where
434        I: Iterator<Item = (Position, String, Vec<usize>)>,
435    {
436        inner: I,
437        current_pos: Option<Position>,
438        tags: HashSet<String>,
439        qiis: HashSet<usize>,
440        finished: bool,
441    }
442    impl<I> GroupByPositionIterator<I>
443    where
444        I: Iterator<Item = (Position, String, Vec<usize>)>,
445    {
446        fn new(inner: I) -> Self {
447            Self {
448                inner,
449                current_pos: None,
450                tags: HashSet::new(),
451                qiis: HashSet::new(),
452                finished: false,
453            }
454        }
455    }
456    impl<I> Iterator for GroupByPositionIterator<I>
457    where
458        I: Iterator<Item = (Position, String, Vec<usize>)>,
459    {
460        type Item = (Position, HashSet<String>, HashSet<usize>);
461        fn next(&mut self) -> Option<Self::Item> {
462            if self.finished {
463                return None;
464            }
465            for (pos, tag, qiids) in self.inner.by_ref() {
466                if self.current_pos.is_none() {
467                    self.current_pos = Some(pos);
468                } else if self.current_pos.unwrap() != pos {
469                    let out_pos = self.current_pos.unwrap();
470                    let out_tags = std::mem::take(&mut self.tags);
471                    let out_qiis = std::mem::take(&mut self.qiis);
472                    self.current_pos = Some(pos);
473                    self.tags.insert(tag);
474                    for q in qiids {
475                        self.qiis.insert(q);
476                    }
477                    return Some((out_pos, out_tags, out_qiis));
478                }
479                self.tags.insert(tag);
480                for q in qiids {
481                    self.qiis.insert(q);
482                }
483            }
484            if let Some(p) = self.current_pos.take() {
485                self.finished = true;
486                let out_tags = std::mem::take(&mut self.tags);
487                let out_qiis = std::mem::take(&mut self.qiis);
488                return Some((p, out_tags, out_qiis));
489            }
490            None
491        }
492    }
493
494    let mut out: Vec<DCBSequencedEvent> = Vec::new();
495    for (pos, tags_present, qiis_present) in GroupByPositionIterator::new(merged) {
496        // Find any query item whose required tag set is subset of tags_present
497        let matching_qiis: Vec<usize> = qiis_present
498            .iter()
499            .copied()
500            .filter(|&qii| qi_tags[qii].is_subset(&tags_present))
501            .collect();
502        if matching_qiis.is_empty() {
503            continue;
504        }
505
506        // Lookup the event record at position
507        let rec = event_tree_lookup(mvcc, dirty, events_tree_root_id, pos)?;
508
509        // Check type and actual tag matching against any of the matching items to avoid hash-collision false positives
510        let mut match_ok = false;
511        'matchcheck: for qii in matching_qiis.iter().copied() {
512            let item = &query.items[qii];
513            // Type must match (or be unspecified)
514            let type_ok = item.types.is_empty() || item.types.iter().any(|t| t == &rec.event_type);
515            if !type_ok {
516                continue;
517            }
518            // Verify actual event tags contain all item tags (guards against tag-hash collisions)
519            let tags_ok = item.tags.iter().all(|t| rec.tags.iter().any(|et| et == t));
520            if tags_ok {
521                match_ok = true;
522                break 'matchcheck;
523            }
524        }
525        if !match_ok {
526            continue;
527        }
528
529        out.push(DCBSequencedEvent {
530            position: pos.0,
531            event: DCBEvent {
532                event_type: rec.event_type,
533                data: rec.data,
534                tags: rec.tags,
535                uuid: rec.uuid,
536            },
537        });
538        if let Some(lim) = limit
539            && out.len() >= lim as usize
540        {
541            break;
542        }
543    }
544
545    Ok(out)
546}
547/// Compute a TagHash ([u8; 8]) from a tag string using a stable 64-bit hash.
548#[inline(always)]
549pub fn tag_to_hash(tag: &str) -> TagHash {
550    const SALT: [u8; 4] = [0x9E, 0x37, 0x79, 0xB9];
551    // Build a 64-bit value by combining two crc32 hashes for stability and simplicity.
552    let mut hasher1 = crc32fast::Hasher::new();
553    hasher1.update(tag.as_bytes());
554    let a = hasher1.finalize();
555
556    let mut hasher2 = crc32fast::Hasher::new();
557    // Note: Benchmark (benches/tag_hash_bench.rs) shows two update() calls
558    // are consistently faster than concatenating bytes+salt into a buffer
559    // and calling update() once, because concatenation requires allocation
560    // and copying. Keeping the two calls avoids extra work and is at least
561    // as fast across sizes from 0..8192 bytes.
562    hasher2.update(tag.as_bytes());
563    hasher2.update(&SALT);
564    let b = hasher2.finalize();
565
566    let value = ((a as u64) << 32) | (b as u64);
567    value.to_le_bytes()
568}
569
570pub fn is_request_idempotent(
571    mvcc: &Arc<Mvcc>,
572    dirty: &HashMap<PageID, Page>,
573    events_tree_root_id: PageID,
574    tags_tree_root_id: PageID,
575    events: &Vec<DCBEvent>,
576    fail_if_events_match: DCBQuery,
577    start: Option<Position>,
578) -> DCBResult<Option<u64>> {
579    // Check events for event IDs. If all have events IDs then
580    // call read_conditional again with limit=event.len() and then
581    // see if all events have matching UUIDs.
582    let submitted_events_len = events.len();
583    let mut submitted_event_ids: Vec<Option<Uuid>> = vec![];
584    for submitted_event in events {
585        if submitted_event.uuid.is_some() {
586            submitted_event_ids.push(submitted_event.uuid);
587        }
588    }
589    if submitted_events_len == submitted_event_ids.len()
590        && submitted_events_len as u64 <= u32::MAX as u64
591    {
592        // All events have UUIDs and there are less than the max size of limit.
593        let read_result = read_conditional(
594            mvcc,
595            dirty,
596            events_tree_root_id,
597            tags_tree_root_id,
598            fail_if_events_match,
599            start,
600            false,
601            Some(submitted_events_len as u32),
602            false,
603        );
604        match read_result {
605            Ok(found_events) => {
606                let mut found_event_ids: Vec<Option<Uuid>> = vec![];
607                let found_events_len = found_events.len();
608                if found_events_len == submitted_events_len {
609                    let last_found_event = &found_events[found_events_len - 1];
610                    let last_found_event_position = last_found_event.position;
611                    for found_event in found_events {
612                        found_event_ids.push(found_event.event.uuid);
613                    }
614                    if found_event_ids == submitted_event_ids {
615                        // It's an idempotent request.
616                        return Ok(Some(last_found_event_position));
617                        // results.push(Ok(last_found_event_position));
618                        // return true
619                    }
620                }
621            }
622            Err(e) => {
623                // Propagate read error for this item but continue with others
624                return Err(e);
625                // results.push(Err(e));
626                // return true;
627            }
628        }
629    }
630    Ok(None)
631}
632
633#[cfg(test)]
634mod tests {
635    use super::*;
636    use crate::page::Page;
637    use serial_test::serial;
638    use std::collections::HashMap;
639    use tempfile::tempdir;
640    use umadb_dcb::{
641        DCBAppendCondition, DCBError, DCBEvent, DCBEventStoreSync, DCBQuery, DCBQueryItem,
642    };
643    use uuid::Uuid;
644
645    // Backward-compatible wrapper for tests: call new read_conditional with an empty dirty map
646    fn read_conditional(
647        mvcc: &Mvcc,
648        events_tree_root_id: PageID,
649        tags_tree_root_id: PageID,
650        query: DCBQuery,
651        start: Option<Position>,
652        backwards: bool,
653        limit: Option<u32>,
654    ) -> DCBResult<Vec<DCBSequencedEvent>> {
655        super::read_conditional(
656            mvcc,
657            &HashMap::<PageID, Page>::new(),
658            events_tree_root_id,
659            tags_tree_root_id,
660            query,
661            start,
662            backwards,
663            limit,
664            false,
665        )
666    }
667
668    static VERBOSE: bool = false;
669
670    // Helper to produce a deterministic set of 10 events with shared tags and unique types
671    fn standard_events() -> Vec<DCBEvent> {
672        let shared_tags = vec![
673            "alpha".to_string(),
674            "beta".to_string(),
675            "gamma".to_string(),
676            "delta".to_string(),
677            "epsilon".to_string(),
678        ];
679        let mut input: Vec<DCBEvent> = Vec::new();
680        for i in 0..10u8 {
681            let t1 = shared_tags[(i % 5) as usize].clone();
682            let t2 = shared_tags[((i + 2) % 5) as usize].clone();
683            input.push(DCBEvent {
684                event_type: format!("Type{}", i),
685                data: vec![i, i + 1, i + 2],
686                tags: vec![t1, t2],
687                uuid: None,
688            });
689        }
690        input
691    }
692
693    // Create DB with the standard events; keep temp dir alive by returning it
694    fn setup_db_with_standard_events() -> (tempfile::TempDir, Mvcc, Vec<DCBEvent>) {
695        let temp_dir = tempdir().unwrap();
696        let db_path = temp_dir.path().join("mvcc-api-test.db");
697        let db = Mvcc::new(db_path.as_ref(), DEFAULT_PAGE_SIZE, VERBOSE).unwrap();
698        let input = standard_events();
699        let mut writer = db.writer().unwrap();
700        let last = unconditional_append(&db, &mut writer, input.clone()).unwrap();
701        db.commit(&mut writer).unwrap();
702        // Verify last equals committed head
703        let (_, header) = db.get_latest_header().unwrap();
704        let head = header.next_position.0.saturating_sub(1);
705        assert_eq!(last, head);
706        (temp_dir, db, input)
707    }
708
709    #[test]
710    #[serial]
711    fn empty_query_after_and_limit() {
712        let (_tmp, mut mvcc, input) = setup_db_with_standard_events();
713
714        // after = 0 -> all
715        let reader = mvcc.reader().unwrap();
716
717        let all = read_conditional(
718            &mut mvcc,
719            reader.events_tree_root_id,
720            reader.tags_tree_root_id,
721            DCBQuery { items: vec![] },
722            Some(Position(1)),
723            false,
724            None,
725        )
726        .unwrap();
727        assert_eq!(all.len(), input.len());
728        assert!(all.windows(2).all(|w| w[0].position < w[1].position));
729
730        // after = first -> tail
731        let first = all[0].position;
732        let tail = read_conditional(
733            &mut mvcc,
734            reader.events_tree_root_id,
735            reader.tags_tree_root_id,
736            DCBQuery { items: vec![] },
737            Some(Position(first + 1)),
738            false,
739            None,
740        )
741        .unwrap();
742        assert_eq!(tail.len(), input.len() - 1);
743
744        // after = last -> empty
745        let last = all.last().unwrap().position;
746        let none = read_conditional(
747            &mut mvcc,
748            reader.events_tree_root_id,
749            reader.tags_tree_root_id,
750            DCBQuery { items: vec![] },
751            Some(Position(last + 1)),
752            false,
753            None,
754        )
755        .unwrap();
756        assert!(none.is_empty());
757
758        // limits
759        let lim0 = read_conditional(
760            &mut mvcc,
761            reader.events_tree_root_id,
762            reader.tags_tree_root_id,
763            DCBQuery { items: vec![] },
764            Some(Position(1)),
765            false,
766            Some(0),
767        )
768        .unwrap();
769        assert!(lim0.is_empty());
770        let lim3 = read_conditional(
771            &mut mvcc,
772            reader.events_tree_root_id,
773            reader.tags_tree_root_id,
774            DCBQuery { items: vec![] },
775            Some(Position(1)),
776            false,
777            Some(3),
778        )
779        .unwrap();
780        assert_eq!(lim3.len(), 3);
781        let lim20 = read_conditional(
782            &mut mvcc,
783            reader.events_tree_root_id,
784            reader.tags_tree_root_id,
785            DCBQuery { items: vec![] },
786            Some(Position(1)),
787            false,
788            Some(20),
789        )
790        .unwrap();
791        assert_eq!(lim20.len(), input.len());
792    }
793
794    #[test]
795    #[serial]
796    fn tags_only_single_tag_after_and_limit() {
797        let (_tmp, mut db, _input) = setup_db_with_standard_events();
798        let qi = DCBQuery {
799            items: vec![DCBQueryItem {
800                types: vec![],
801                tags: vec!["alpha".to_string()],
802            }],
803        };
804        let reader = db.reader().unwrap();
805        let res = read_conditional(
806            &mut db,
807            reader.events_tree_root_id,
808            reader.tags_tree_root_id,
809            qi.clone(),
810            Some(Position(1)),
811            false,
812            None,
813        )
814        .unwrap();
815        assert_eq!(res.len(), 4);
816        assert!(
817            res.iter()
818                .all(|e| e.event.tags.iter().any(|t| t == "alpha"))
819        );
820        assert!(res.windows(2).all(|w| w[0].position < w[1].position));
821
822        // after combinations
823        let positions: Vec<u64> = res.iter().map(|e| e.position).collect();
824        let after_first = read_conditional(
825            &mut db,
826            reader.events_tree_root_id,
827            reader.tags_tree_root_id,
828            qi.clone(),
829            Some(Position(positions[0] + 1)),
830            false,
831            None,
832        )
833        .unwrap();
834        assert_eq!(after_first.len(), positions.len() - 1);
835        let after_last = read_conditional(
836            &mut db,
837            reader.events_tree_root_id,
838            reader.tags_tree_root_id,
839            qi.clone(),
840            Some(Position(*positions.last().unwrap() + 1)),
841            false,
842            None,
843        )
844        .unwrap();
845        assert!(after_last.is_empty());
846
847        // limits
848        let lim0 = read_conditional(
849            &mut db,
850            reader.events_tree_root_id,
851            reader.tags_tree_root_id,
852            qi.clone(),
853            Some(Position(1)),
854            false,
855            Some(0),
856        )
857        .unwrap();
858        assert!(lim0.is_empty());
859        let lim1 = read_conditional(
860            &mut db,
861            reader.events_tree_root_id,
862            reader.tags_tree_root_id,
863            qi.clone(),
864            Some(Position(1)),
865            false,
866            Some(1),
867        )
868        .unwrap();
869        assert_eq!(lim1.len(), 1);
870        let lim10 = read_conditional(
871            &mut db,
872            reader.events_tree_root_id,
873            reader.tags_tree_root_id,
874            qi,
875            Some(Position(1)),
876            false,
877            Some(10),
878        )
879        .unwrap();
880        assert_eq!(lim10.len(), 4);
881    }
882
883    #[test]
884    #[serial]
885    fn tags_only_multi_tag_and() {
886        let (_tmp, mut db, _input) = setup_db_with_standard_events();
887        let qi = DCBQuery {
888            items: vec![DCBQueryItem {
889                types: vec![],
890                tags: vec!["alpha".to_string(), "gamma".to_string()],
891            }],
892        };
893        let reader = db.reader().unwrap();
894        let res = read_conditional(
895            &mut db,
896            reader.events_tree_root_id,
897            reader.tags_tree_root_id,
898            qi,
899            Some(Position(1)),
900            false,
901            None,
902        )
903        .unwrap();
904        assert_eq!(res.len(), 2);
905        assert!(
906            res.iter()
907                .all(|e| e.event.tags.iter().any(|t| t == "alpha"))
908        );
909        assert!(
910            res.iter()
911                .all(|e| e.event.tags.iter().any(|t| t == "gamma"))
912        );
913    }
914
915    #[test]
916    #[serial]
917    fn types_plus_tags_index_path() {
918        let (_tmp, mut db, _input) = setup_db_with_standard_events();
919        let qi = DCBQuery {
920            items: vec![DCBQueryItem {
921                types: vec!["Type0".to_string()],
922                tags: vec!["alpha".to_string()],
923            }],
924        };
925        let reader = db.reader().unwrap();
926        let res = read_conditional(
927            &mut db,
928            reader.events_tree_root_id,
929            reader.tags_tree_root_id,
930            qi,
931            Some(Position(1)),
932            false,
933            None,
934        )
935        .unwrap();
936        assert_eq!(res.len(), 1);
937        assert_eq!(res[0].event.event_type, "Type0");
938        assert!(res[0].event.tags.iter().any(|t| t == "alpha"));
939    }
940
941    #[test]
942    #[serial]
943    fn or_semantics_and_deduplication() {
944        let (_tmp, mut db, _input) = setup_db_with_standard_events();
945        let alpha_only = DCBQuery {
946            items: vec![DCBQueryItem {
947                types: vec![],
948                tags: vec!["alpha".to_string()],
949            }],
950        };
951        let reader = db.reader().unwrap();
952        let alpha_positions: Vec<u64> = read_conditional(
953            &mut db,
954            reader.events_tree_root_id,
955            reader.tags_tree_root_id,
956            alpha_only.clone(),
957            Some(Position(1)),
958            false,
959            None,
960        )
961        .unwrap()
962        .into_iter()
963        .map(|e| e.position)
964        .collect();
965
966        // Overlapping items: alpha OR (alpha AND gamma) should deduplicate
967        let query = DCBQuery {
968            items: vec![
969                DCBQueryItem {
970                    types: vec![],
971                    tags: vec!["alpha".to_string()],
972                },
973                DCBQueryItem {
974                    types: vec![],
975                    tags: vec!["alpha".to_string(), "gamma".to_string()],
976                },
977            ],
978        };
979        let res = read_conditional(
980            &mut db,
981            reader.events_tree_root_id,
982            reader.tags_tree_root_id,
983            query,
984            Some(Position(1)),
985            false,
986            None,
987        )
988        .unwrap();
989        let res_positions: Vec<u64> = res.into_iter().map(|e| e.position).collect();
990        assert_eq!(res_positions, alpha_positions);
991    }
992
993    #[test]
994    #[serial]
995    fn fallback_types_only_after_and_limit() {
996        let temp_dir = tempdir().unwrap();
997        let db_path = temp_dir.path().join("mvcc-fallback-types-only.db");
998        let mut db = Mvcc::new(db_path.as_ref(), DEFAULT_PAGE_SIZE, VERBOSE).unwrap();
999
1000        // Use a smaller custom set to make counts obvious
1001        let events = vec![
1002            DCBEvent {
1003                event_type: "TypeA".to_string(),
1004                data: vec![1],
1005                tags: vec!["x".to_string()],
1006                uuid: None,
1007            },
1008            DCBEvent {
1009                event_type: "TypeB".to_string(),
1010                data: vec![2],
1011                tags: vec!["y".to_string()],
1012                uuid: None,
1013            },
1014            DCBEvent {
1015                event_type: "TypeA".to_string(),
1016                data: vec![3],
1017                tags: vec!["z".to_string()],
1018                uuid: None,
1019            },
1020        ];
1021        let mut writer = db.writer().unwrap();
1022        let last = unconditional_append(&db, &mut writer, events).unwrap();
1023        db.commit(&mut writer).unwrap();
1024        let (_, header) = db.get_latest_header().unwrap();
1025        let head = header.next_position.0.saturating_sub(1);
1026        assert_eq!(last, head);
1027
1028        // Query item with no tags => forces fallback path; select TypeA only
1029        let qi = DCBQuery {
1030            items: vec![DCBQueryItem {
1031                types: vec!["TypeA".to_string()],
1032                tags: vec![],
1033            }],
1034        };
1035        let reader = db.reader().unwrap();
1036        let res = read_conditional(
1037            &mut db,
1038            reader.events_tree_root_id,
1039            reader.tags_tree_root_id,
1040            qi.clone(),
1041            Some(Position(1)),
1042            false,
1043            None,
1044        )
1045        .unwrap();
1046        assert_eq!(res.len(), 2);
1047        assert!(res.iter().all(|e| e.event.event_type == "TypeA"));
1048
1049        // After skip first matching
1050        let first_pos = res[0].position;
1051        let res_after = read_conditional(
1052            &mut db,
1053            reader.events_tree_root_id,
1054            reader.tags_tree_root_id,
1055            qi.clone(),
1056            Some(Position(first_pos + 1)),
1057            false,
1058            None,
1059        )
1060        .unwrap();
1061        assert_eq!(res_after.len(), 1);
1062
1063        // Limit 1
1064        let res_lim1 = read_conditional(
1065            &mut db,
1066            reader.events_tree_root_id,
1067            reader.tags_tree_root_id,
1068            qi,
1069            Some(Position(1)),
1070            false,
1071            Some(1),
1072        )
1073        .unwrap();
1074        assert_eq!(res_lim1.len(), 1);
1075    }
1076
1077    #[test]
1078    #[serial]
1079    fn fallback_empty_item_matches_all() {
1080        let (_tmp, mut db, input) = setup_db_with_standard_events();
1081        // An empty item (no types, no tags) should match all events via fallback path
1082        let qi = DCBQuery {
1083            items: vec![DCBQueryItem {
1084                types: vec![],
1085                tags: vec![],
1086            }],
1087        };
1088
1089        let reader = db.reader().unwrap();
1090        let all = read_conditional(
1091            &mut db,
1092            reader.events_tree_root_id,
1093            reader.tags_tree_root_id,
1094            qi.clone(),
1095            Some(Position(1)),
1096            false,
1097            None,
1098        )
1099        .unwrap();
1100        assert_eq!(all.len(), input.len());
1101
1102        // After and limit still apply
1103        let first = all[1].position;
1104        let tail = read_conditional(
1105            &mut db,
1106            reader.events_tree_root_id,
1107            reader.tags_tree_root_id,
1108            qi.clone(),
1109            Some(Position(first)),
1110            false,
1111            None,
1112        )
1113        .unwrap();
1114        assert_eq!(tail.len(), input.len() - 1);
1115        let lim5 = read_conditional(
1116            &mut db,
1117            reader.events_tree_root_id,
1118            reader.tags_tree_root_id,
1119            qi,
1120            Some(Position(1)),
1121            false,
1122            Some(5),
1123        )
1124        .unwrap();
1125        assert_eq!(lim5.len(), 5);
1126    }
1127
1128    #[test]
1129    #[serial]
1130    fn test_event_store() {
1131        let temp_dir = tempdir().unwrap();
1132        let store = UmaDB::new(temp_dir.path()).unwrap();
1133
1134        // Head is None on empty store
1135        assert_eq!(None, store.head().unwrap());
1136
1137        // Append a couple of events
1138        let events = vec![
1139            DCBEvent {
1140                event_type: "TypeA".to_string(),
1141                data: vec![1],
1142                tags: vec!["foo".to_string()],
1143                uuid: None,
1144            },
1145            DCBEvent {
1146                event_type: "TypeB".to_string(),
1147                data: vec![2],
1148                tags: vec!["bar".to_string(), "foo".to_string()],
1149                uuid: None,
1150            },
1151        ];
1152        let last = store.append(events.clone(), None).unwrap();
1153        assert!(last > 0);
1154        assert_eq!(store.head().unwrap(), Some(last));
1155
1156        // Read all
1157        let mut resp = store.read(None, None, false, None, false).unwrap();
1158        let (all, head) = resp.collect_with_head().unwrap();
1159        assert_eq!(head, Some(last));
1160        assert_eq!(all.len(), 2);
1161        assert_eq!(all[0].event.event_type, "TypeA");
1162        assert_eq!(all[1].event.event_type, "TypeB");
1163
1164        // Limit semantics: only first event returned and head equals that position
1165        let mut resp_lim1 = store.read(None, None, false, Some(1), false).unwrap();
1166        let (only_one, head_lim1) = resp_lim1.collect_with_head().unwrap();
1167        assert_eq!(only_one.len(), 1);
1168        assert_eq!(only_one[0].event.event_type, "TypeA");
1169        assert_eq!(head_lim1, Some(only_one[0].position));
1170
1171        // Tag-filtered read ("foo")
1172        let query = DCBQuery {
1173            items: vec![DCBQueryItem {
1174                types: vec![],
1175                tags: vec!["foo".to_string()],
1176            }],
1177        };
1178        let mut resp2 = store.read(Some(query), None, false, None, false).unwrap();
1179        let out2 = resp2.next_batch().unwrap();
1180        assert_eq!(out2.len(), 2);
1181        assert!(out2.iter().all(|e| e.event.tags.iter().any(|t| t == "foo")));
1182
1183        // From semantics: skip the first event
1184        let first_pos = all[0].position + 1;
1185        let mut resp3 = store
1186            .read(None, Some(first_pos), false, None, false)
1187            .unwrap();
1188        let out3 = resp3.next_batch().unwrap();
1189        assert_eq!(out3.len(), 1);
1190        assert_eq!(out3[0].event.event_type, "TypeB");
1191
1192        // Append with a condition that should PASS: query matches existing 'foo' but after = last
1193        let cond_pass = DCBAppendCondition {
1194            fail_if_events_match: DCBQuery {
1195                items: vec![DCBQueryItem {
1196                    types: vec![],
1197                    tags: vec!["foo".to_string()],
1198                }],
1199            },
1200            after: Some(last),
1201        };
1202        let ok_last = store
1203            .append(
1204                vec![DCBEvent {
1205                    event_type: "TypeC".to_string(),
1206                    data: vec![3],
1207                    tags: vec!["baz".to_string()],
1208                    uuid: None,
1209                }],
1210                Some(cond_pass),
1211            )
1212            .expect("append with passing condition should succeed");
1213        assert!(ok_last > last);
1214        assert_eq!(store.head().unwrap(), Some(ok_last));
1215
1216        // Append with a condition that should FAIL: same query but after = 0
1217        let cond_fail = DCBAppendCondition {
1218            fail_if_events_match: DCBQuery {
1219                items: vec![DCBQueryItem {
1220                    types: vec![],
1221                    tags: vec!["foo".to_string()],
1222                }],
1223            },
1224            after: Some(0),
1225        };
1226        let before_head = store.head().unwrap();
1227        let res = store.append(
1228            vec![DCBEvent {
1229                event_type: "TypeD".to_string(),
1230                data: vec![4],
1231                tags: vec!["qux".to_string()],
1232                uuid: None,
1233            }],
1234            Some(cond_fail),
1235        );
1236        match res {
1237            Err(DCBError::IntegrityError(_)) => {}
1238            other => panic!("Expected IntegrityError, got {:?}", other),
1239        }
1240        // Ensure head unchanged after failed append
1241        assert_eq!(store.head().unwrap(), before_head);
1242    }
1243
1244    #[test]
1245    fn test_append_batch_mixed_conditions() {
1246        let temp_dir = tempdir().unwrap();
1247        let store = UmaDB::new(temp_dir.path()).unwrap();
1248
1249        let e1 = DCBEvent {
1250            event_type: "A".into(),
1251            data: b"1".to_vec(),
1252            tags: vec!["t1".into()],
1253            uuid: None,
1254        };
1255        let e2 = DCBEvent {
1256            event_type: "B".into(),
1257            data: b"2".to_vec(),
1258            tags: vec!["t2".into()],
1259            uuid: None,
1260        };
1261        let e3 = DCBEvent {
1262            event_type: "C".into(),
1263            data: b"3".to_vec(),
1264            tags: vec!["t3".into()],
1265            uuid: None,
1266        };
1267
1268        // Batch: first succeeds, second fails due to condition matching any event, third succeeds (after high position)
1269        let items = vec![
1270            (vec![e1.clone()], None),
1271            (
1272                vec![e2.clone()],
1273                Some(DCBAppendCondition {
1274                    fail_if_events_match: DCBQuery::default(),
1275                    after: None,
1276                }),
1277            ),
1278            (
1279                vec![e3.clone()],
1280                Some(DCBAppendCondition {
1281                    fail_if_events_match: DCBQuery::default(),
1282                    after: Some(10),
1283                }),
1284            ),
1285        ];
1286
1287        let results = store.append_batch(items, false).unwrap();
1288
1289        assert_eq!(results.len(), 3);
1290        // First item should succeed with last position 1
1291        match &results[0] {
1292            Ok(pos) => assert_eq!(*pos, 1),
1293            Err(e) => panic!("unexpected error for first item: {:?}", e),
1294        }
1295        // Second item should fail integrity
1296        match &results[1] {
1297            Ok(pos) => panic!("expected integrity error, got Ok({})", pos),
1298            Err(e) => assert!(matches!(e, DCBError::IntegrityError(_))),
1299        }
1300        // Third item should succeed with last position 2 (since second didn't append)
1301        match &results[2] {
1302            Ok(pos) => assert_eq!(*pos, 2),
1303            Err(e) => panic!("unexpected error for third item: {:?}", e),
1304        }
1305
1306        // Verify committed state: only e1 and e3 should be present, head is 2
1307        let (events, head) = store.read_with_head(None, None, false, None).unwrap();
1308        assert_eq!(events.len(), 2);
1309        assert_eq!(events[0].event.data, e1.data);
1310        assert_eq!(events[1].event.data, e3.data);
1311        assert_eq!(head, Some(2));
1312    }
1313
1314    #[test]
1315    fn test_append_batch_dirty_visibility_with_tags() {
1316        let temp_dir = tempdir().unwrap();
1317        let store = UmaDB::new(temp_dir.path()).unwrap();
1318
1319        // Item1 introduces tag "x"; Item2 condition queries tag "x" and must see it via dirty tags tree; Item3 uses after to ignore it
1320        let e1 = DCBEvent {
1321            event_type: "T".into(),
1322            data: b"one".to_vec(),
1323            tags: vec!["x".into()],
1324            uuid: None,
1325        };
1326        let e2 = DCBEvent {
1327            event_type: "T".into(),
1328            data: b"two".to_vec(),
1329            tags: vec!["y".into()],
1330            uuid: None,
1331        };
1332        let e3 = DCBEvent {
1333            event_type: "T".into(),
1334            data: b"three".to_vec(),
1335            tags: vec!["z".into()],
1336            uuid: None,
1337        };
1338
1339        let query_tag_x = DCBQuery {
1340            items: vec![DCBQueryItem {
1341                types: vec![],
1342                tags: vec!["x".into()],
1343            }],
1344        };
1345
1346        let items = vec![
1347            // 1) Append e1 (tag x)
1348            (vec![e1.clone()], None),
1349            // 2) Attempt append e2, but fail if any events with tag x exist after None (i.e., from the start); should fail due to e1 in dirty pages
1350            (
1351                vec![e2.clone()],
1352                Some(DCBAppendCondition {
1353                    fail_if_events_match: query_tag_x.clone(),
1354                    after: None,
1355                }),
1356            ),
1357            // 3) Append e3 with condition that ignores position 1 by using after=Some(1); should pass
1358            (
1359                vec![e3.clone()],
1360                Some(DCBAppendCondition {
1361                    fail_if_events_match: query_tag_x.clone(),
1362                    after: Some(1),
1363                }),
1364            ),
1365        ];
1366
1367        let results = store.append_batch(items, false).unwrap();
1368
1369        assert_eq!(results.len(), 3);
1370        match &results[0] {
1371            Ok(pos) => assert_eq!(*pos, 1),
1372            Err(e) => panic!("unexpected error for first item: {:?}", e),
1373        }
1374        match &results[1] {
1375            Ok(pos) => panic!("expected integrity error, got Ok({})", pos),
1376            Err(e) => assert!(matches!(e, DCBError::IntegrityError(_))),
1377        }
1378        match &results[2] {
1379            Ok(pos) => assert_eq!(*pos, 2),
1380            Err(e) => panic!("unexpected error for third item: {:?}", e),
1381        }
1382
1383        // Verify committed state and tag index behavior
1384        let (events, head) = store.read_with_head(None, None, false, None).unwrap();
1385        assert_eq!(events.len(), 2);
1386        assert_eq!(events[0].event.data, e1.data);
1387        assert_eq!(events[1].event.data, e3.data);
1388        assert_eq!(head, Some(2));
1389
1390        // Query by tag x returns only the first event
1391        let (tagx_events, _) = store
1392            .read_with_head(Some(query_tag_x.clone()), None, false, None)
1393            .unwrap();
1394        assert_eq!(tagx_events.len(), 1);
1395        assert_eq!(tagx_events[0].event.data, e1.data);
1396    }
1397
1398    #[test]
1399    fn test_append_batch_dirty_visibility_with_types_small_and_big_overflow() {
1400        let temp_dir = tempdir().unwrap();
1401        let store = UmaDB::new(temp_dir.path()).unwrap();
1402
1403        // Prepare events
1404        let small = DCBEvent {
1405            event_type: "S".into(),
1406            data: b"sm".to_vec(),
1407            tags: vec!["tS".into()],
1408            uuid: None,
1409        };
1410        // Large data to ensure it spills into event overflow pages
1411        let big_data_len = DEFAULT_PAGE_SIZE * 3; // 3 pages worth to be safe
1412        let big = DCBEvent {
1413            event_type: "B".into(),
1414            data: vec![0xAB; big_data_len],
1415            tags: vec!["tB".into()],
1416            uuid: None,
1417        };
1418        let filler1 = DCBEvent {
1419            event_type: "X".into(),
1420            data: b"x".to_vec(),
1421            tags: vec![],
1422            uuid: None,
1423        };
1424        let filler2 = DCBEvent {
1425            event_type: "Y".into(),
1426            data: b"y".to_vec(),
1427            tags: vec![],
1428            uuid: None,
1429        };
1430        let final_ok = DCBEvent {
1431            event_type: "C".into(),
1432            data: b"c".to_vec(),
1433            tags: vec![],
1434            uuid: None,
1435        };
1436
1437        // Queries by type only (no tags) to force fallback path over events tree (which reads from dirty pages)
1438        let q_type_s = DCBQuery {
1439            items: vec![DCBQueryItem {
1440                types: vec!["S".into()],
1441                tags: vec![],
1442            }],
1443        };
1444        let q_type_b = DCBQuery {
1445            items: vec![DCBQueryItem {
1446                types: vec!["B".into()],
1447                tags: vec![],
1448            }],
1449        };
1450
1451        let items = vec![
1452            // 1) Append small S
1453            (vec![small.clone()], None),
1454            // 2) Should fail because type S exists in dirty pages (after None)
1455            (
1456                vec![filler1.clone()],
1457                Some(DCBAppendCondition {
1458                    fail_if_events_match: q_type_s.clone(),
1459                    after: None,
1460                }),
1461            ),
1462            // 3) Append big B (overflow)
1463            (vec![big.clone()], None),
1464            // 4) Should fail because type B exists in dirty pages (after None)
1465            (
1466                vec![filler2.clone()],
1467                Some(DCBAppendCondition {
1468                    fail_if_events_match: q_type_b.clone(),
1469                    after: None,
1470                }),
1471            ),
1472            // 5) Should succeed because after=Some(2) ignores positions <= 2 (small at 1, big at 2)
1473            (
1474                vec![final_ok.clone()],
1475                Some(DCBAppendCondition {
1476                    fail_if_events_match: q_type_b.clone(),
1477                    after: Some(2),
1478                }),
1479            ),
1480        ];
1481
1482        let results = store.append_batch(items, false).unwrap();
1483        assert_eq!(results.len(), 5);
1484        match &results[0] {
1485            Ok(pos) => assert_eq!(*pos, 1),
1486            other => panic!("unexpected for item0: {:?}", other),
1487        }
1488        match &results[1] {
1489            Err(DCBError::IntegrityError(_)) => {}
1490            other => panic!("expected IntegrityError for item1, got {:?}", other),
1491        }
1492        match &results[2] {
1493            Ok(pos) => assert_eq!(*pos, 2),
1494            other => panic!("unexpected for item2: {:?}", other),
1495        }
1496        match &results[3] {
1497            Err(DCBError::IntegrityError(_)) => {}
1498            other => panic!("expected IntegrityError for item3, got {:?}", other),
1499        }
1500        match &results[4] {
1501            Ok(pos) => assert_eq!(*pos, 3),
1502            other => panic!("unexpected for item4: {:?}", other),
1503        }
1504
1505        // Verify committed state: we should have small (pos1), big (pos2), final_ok (pos3)
1506        let (events, head) = store.read_with_head(None, None, false, None).unwrap();
1507        assert_eq!(events.len(), 3);
1508        assert_eq!(events[0].event.event_type, small.event_type);
1509        assert_eq!(events[1].event.event_type, big.event_type);
1510        assert_eq!(events[2].event.event_type, final_ok.event_type);
1511        assert_eq!(head, Some(3));
1512
1513        // Check type queries and large data integrity
1514        let (small_by_type, _) = store
1515            .read_with_head(Some(q_type_s.clone()), None, false, None)
1516            .unwrap();
1517        assert_eq!(small_by_type.len(), 1);
1518        assert_eq!(small_by_type[0].event.event_type, "S");
1519
1520        let (big_by_type, _) = store
1521            .read_with_head(Some(q_type_b.clone()), None, false, None)
1522            .unwrap();
1523        assert_eq!(big_by_type.len(), 1);
1524        assert_eq!(big_by_type[0].event.event_type, "B");
1525        assert_eq!(big_by_type[0].event.data.len(), big_data_len);
1526        assert!(big_by_type[0].event.data.iter().all(|&b| b == 0xAB));
1527    }
1528
1529    #[test]
1530    fn test_append_batch_dirty_visibility_with_tags_and_types_small_and_big_overflow() {
1531        let temp_dir = tempdir().unwrap();
1532        let store = UmaDB::new(temp_dir.path()).unwrap();
1533
1534        // Small inline event: type "S" with tag "x"
1535        let small = DCBEvent {
1536            event_type: "S".into(),
1537            data: b"sm".to_vec(),
1538            tags: vec!["x".into()],
1539            uuid: None,
1540        };
1541        // Big overflow event: type "B" with tag "y" and large payload to exercise overflow pages
1542        let big_data_len = DEFAULT_PAGE_SIZE * 3; // ensure multiple overflow pages
1543        let big = DCBEvent {
1544            event_type: "B".into(),
1545            data: vec![0xCD; big_data_len],
1546            tags: vec!["y".into()],
1547            uuid: None,
1548        };
1549        // Fillers that will be conditioned out
1550        let filler1 = DCBEvent {
1551            event_type: "X".into(),
1552            data: b"x".to_vec(),
1553            tags: vec![],
1554            uuid: None,
1555        };
1556        let filler2 = DCBEvent {
1557            event_type: "Y".into(),
1558            data: b"y".to_vec(),
1559            tags: vec![],
1560            uuid: None,
1561        };
1562        let final_ok = DCBEvent {
1563            event_type: "C".into(),
1564            data: b"c".to_vec(),
1565            tags: vec![],
1566            uuid: None,
1567        };
1568
1569        // Conditions combining tags and types so the tags index is used and the type filter applies after lookup
1570        let q_s_and_x = DCBQuery {
1571            items: vec![DCBQueryItem {
1572                types: vec!["S".into()],
1573                tags: vec!["x".into()],
1574            }],
1575        };
1576        let q_b_and_y = DCBQuery {
1577            items: vec![DCBQueryItem {
1578                types: vec!["B".into()],
1579                tags: vec!["y".into()],
1580            }],
1581        };
1582
1583        let items = vec![
1584            // 1) Append small S@x
1585            (vec![small.clone()], None),
1586            // 2) Should fail because S@x exists in dirty pages (tags path + type filter)
1587            (
1588                vec![filler1.clone()],
1589                Some(DCBAppendCondition {
1590                    fail_if_events_match: q_s_and_x.clone(),
1591                    after: None,
1592                }),
1593            ),
1594            // 3) Append big B@y (overflow)
1595            (vec![big.clone()], None),
1596            // 4) Should fail because B@y exists in dirty pages (tags path + type filter and overflow read)
1597            (
1598                vec![filler2.clone()],
1599                Some(DCBAppendCondition {
1600                    fail_if_events_match: q_b_and_y.clone(),
1601                    after: None,
1602                }),
1603            ),
1604            // 5) Should succeed because after=Some(2) ignores positions <= 2 (small at 1, big at 2)
1605            (
1606                vec![final_ok.clone()],
1607                Some(DCBAppendCondition {
1608                    fail_if_events_match: q_b_and_y.clone(),
1609                    after: Some(2),
1610                }),
1611            ),
1612        ];
1613
1614        let results = store.append_batch(items, false).unwrap();
1615        assert_eq!(results.len(), 5);
1616        match &results[0] {
1617            Ok(pos) => assert_eq!(*pos, 1),
1618            other => panic!("unexpected for item0: {:?}", other),
1619        }
1620        match &results[1] {
1621            Err(DCBError::IntegrityError(_)) => {}
1622            other => panic!("expected IntegrityError for item1, got {:?}", other),
1623        }
1624        match &results[2] {
1625            Ok(pos) => assert_eq!(*pos, 2),
1626            other => panic!("unexpected for item2: {:?}", other),
1627        }
1628        match &results[3] {
1629            Err(DCBError::IntegrityError(_)) => {}
1630            other => panic!("expected IntegrityError for item3, got {:?}", other),
1631        }
1632        match &results[4] {
1633            Ok(pos) => assert_eq!(*pos, 3),
1634            other => panic!("unexpected for item4: {:?}", other),
1635        }
1636
1637        // Verify committed state and order
1638        let (events, head) = store.read_with_head(None, None, false, None).unwrap();
1639        assert_eq!(events.len(), 3);
1640        assert_eq!(events[0].event.event_type, small.event_type);
1641        assert_eq!(events[1].event.event_type, big.event_type);
1642        assert_eq!(events[2].event.event_type, final_ok.event_type);
1643        assert_eq!(head, Some(3));
1644
1645        // Query by combined type+tag should return exactly one for each
1646        let (small_combined, _) = store
1647            .read_with_head(Some(q_s_and_x.clone()), None, false, None)
1648            .unwrap();
1649        assert_eq!(small_combined.len(), 1);
1650        assert_eq!(small_combined[0].event.event_type, "S");
1651        assert!(small_combined[0].event.tags.iter().any(|t| t == "x"));
1652
1653        let (big_combined, _) = store
1654            .read_with_head(Some(q_b_and_y.clone()), None, false, None)
1655            .unwrap();
1656        assert_eq!(big_combined.len(), 1);
1657        assert_eq!(big_combined[0].event.event_type, "B");
1658        assert!(big_combined[0].event.tags.iter().any(|t| t == "y"));
1659        assert_eq!(big_combined[0].event.data.len(), big_data_len);
1660        assert!(big_combined[0].event.data.iter().all(|&b| b == 0xCD));
1661    }
1662
1663    #[test]
1664    fn test_append_event_with_uuid_is_maintained_and_activated_append_idempotency() {
1665        let temp_dir = tempdir().unwrap();
1666        let store = UmaDB::new(temp_dir.path()).unwrap();
1667
1668        let condition1 = Some(DCBAppendCondition {
1669            fail_if_events_match: DCBQuery { items: vec![] },
1670            after: None,
1671        });
1672
1673        let event1 = DCBEvent {
1674            event_type: "type1".to_string(),
1675            data: b"data1".to_vec(),
1676            tags: vec!["tag1".to_string()],
1677            uuid: Some(Uuid::new_v4()),
1678        };
1679
1680        let mut commit_position1 = store
1681            .append(vec![event1.clone()], condition1.clone())
1682            .unwrap();
1683        assert_eq!(1, commit_position1);
1684
1685        let (result, head) = store.read_with_head(None, None, false, None).unwrap();
1686        assert_eq!(1, result.len());
1687        assert_eq!(Some(1), head);
1688        assert_eq!(event1.uuid, result[0].event.uuid);
1689
1690        // Test idempotency - retry the same append operation.
1691        commit_position1 = store
1692            .append(vec![event1.clone()], condition1.clone())
1693            .unwrap();
1694
1695        // Check the response is the same as before.
1696        assert_eq!(1, commit_position1);
1697
1698        // Check we still have only one sequenced event.
1699        let (result, head) = store.read_with_head(None, None, false, None).unwrap();
1700        assert_eq!(1, result.len());
1701        assert_eq!(Some(1), head);
1702        assert_eq!(event1.uuid, result[0].event.uuid);
1703
1704        // Append another event.
1705        let event2 = DCBEvent {
1706            event_type: "type2".to_string(),
1707            data: b"data2".to_vec(),
1708            tags: vec!["tag2".to_string()],
1709            uuid: Some(Uuid::new_v4()),
1710        };
1711
1712        let mut commit_position2 = store.append(vec![event2.clone()], None).unwrap();
1713        assert_eq!(2, commit_position2);
1714
1715        // Check we have two sequenced events.
1716        let (result, head) = store.read_with_head(None, None, false, None).unwrap();
1717        assert_eq!(2, result.len());
1718        assert_eq!(Some(2), head);
1719        assert_eq!(event1.uuid, result[0].event.uuid);
1720        assert_eq!(event2.uuid, result[1].event.uuid);
1721
1722        // Test idempotency - retry the same append operation.
1723        commit_position1 = store
1724            .append(vec![event1.clone()], condition1.clone())
1725            .unwrap();
1726
1727        // Check the response is the same as before.
1728        assert_eq!(1, commit_position1);
1729
1730        // Test idempotency - try an operation with event1 and event2.
1731        commit_position2 = store
1732            .append(vec![event1.clone(), event2.clone()], condition1.clone())
1733            .unwrap();
1734
1735        // Check the response is the same as before.
1736        assert_eq!(2, commit_position2);
1737
1738        // Check we still have two sequenced events.
1739        let (result, head) = store.read_with_head(None, None, false, None).unwrap();
1740        assert_eq!(2, result.len());
1741        assert_eq!(Some(2), head);
1742        assert_eq!(event1.uuid, result[0].event.uuid);
1743        assert_eq!(event2.uuid, result[1].event.uuid);
1744
1745        // Try with event2 and condition1 - should get an error.
1746        let result = store.append(vec![event2.clone()], condition1.clone());
1747        assert!(matches!(result, Err(DCBError::IntegrityError(_))));
1748
1749        // Try with two events in different order - should get an error.
1750        let result = store.append(vec![event2.clone(), event1.clone()], condition1.clone());
1751        assert!(matches!(result, Err(DCBError::IntegrityError(_))));
1752    }
1753
1754    #[test]
1755    #[serial]
1756    fn empty_query_backwards_from_and_limit() {
1757        let (_tmp, mvcc, _input) = setup_db_with_standard_events();
1758        let reader = mvcc.reader().unwrap();
1759
1760        // Forwards: all events starting from position 1
1761        let fwd = read_conditional(
1762            &mvcc,
1763            reader.events_tree_root_id,
1764            reader.tags_tree_root_id,
1765            DCBQuery { items: vec![] },
1766            Some(Position(1)),
1767            false,
1768            None,
1769        )
1770        .unwrap();
1771        let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
1772        assert!(!fwd_pos.is_empty());
1773        assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
1774
1775        // Backwards: all events (from=None) in descending order
1776        let back_all = read_conditional(
1777            &mvcc,
1778            reader.events_tree_root_id,
1779            reader.tags_tree_root_id,
1780            DCBQuery { items: vec![] },
1781            None,
1782            true,
1783            None,
1784        )
1785        .unwrap();
1786        let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
1787        let mut fwd_rev = fwd_pos.clone();
1788        fwd_rev.reverse();
1789        assert_eq!(fwd_rev, back_all_pos);
1790
1791        // Backwards with from=last should still return all (<= last)
1792        let last = *fwd_pos.last().unwrap();
1793        let back_from_last = read_conditional(
1794            &mvcc,
1795            reader.events_tree_root_id,
1796            reader.tags_tree_root_id,
1797            DCBQuery { items: vec![] },
1798            Some(Position(last)),
1799            true,
1800            None,
1801        )
1802        .unwrap();
1803        let back_from_last_pos: Vec<u64> = back_from_last.iter().map(|e| e.position).collect();
1804        assert_eq!(back_from_last_pos, fwd_rev);
1805
1806        // Backwards with from=last-1 should drop the very last element
1807        let back_from_before_last = read_conditional(
1808            &mvcc,
1809            reader.events_tree_root_id,
1810            reader.tags_tree_root_id,
1811            DCBQuery { items: vec![] },
1812            Some(Position(last - 1)),
1813            true,
1814            None,
1815        )
1816        .unwrap();
1817        let back_from_before_last_pos: Vec<u64> =
1818            back_from_before_last.iter().map(|e| e.position).collect();
1819        assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
1820
1821        // Limit in backwards order: first 3 of the reversed forward vector
1822        let back_lim3 = read_conditional(
1823            &mvcc,
1824            reader.events_tree_root_id,
1825            reader.tags_tree_root_id,
1826            DCBQuery { items: vec![] },
1827            None,
1828            true,
1829            Some(3),
1830        )
1831        .unwrap();
1832        let back_lim3_pos: Vec<u64> = back_lim3.iter().map(|e| e.position).collect();
1833        assert_eq!(back_lim3_pos, fwd_rev[..3.min(fwd_rev.len())].to_vec());
1834    }
1835
1836    #[test]
1837    #[serial]
1838    fn tags_only_single_tag_backwards() {
1839        let (_tmp, mvcc, _input) = setup_db_with_standard_events();
1840        let reader = mvcc.reader().unwrap();
1841        let qi = DCBQuery {
1842            items: vec![DCBQueryItem {
1843                types: vec![],
1844                tags: vec!["alpha".to_string()],
1845            }],
1846        };
1847
1848        let fwd = read_conditional(
1849            &mvcc,
1850            reader.events_tree_root_id,
1851            reader.tags_tree_root_id,
1852            qi.clone(),
1853            Some(Position(1)),
1854            false,
1855            None,
1856        )
1857        .unwrap();
1858        let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
1859        assert!(!fwd_pos.is_empty());
1860        assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
1861
1862        let back_all = read_conditional(
1863            &mvcc,
1864            reader.events_tree_root_id,
1865            reader.tags_tree_root_id,
1866            qi.clone(),
1867            None,
1868            true,
1869            None,
1870        )
1871        .unwrap();
1872        let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
1873        let mut fwd_rev = fwd_pos.clone();
1874        fwd_rev.reverse();
1875        assert_eq!(back_all_pos, fwd_rev);
1876
1877        // from = just before the last matching position should drop the newest one in backwards order
1878        let last = *fwd_pos.last().unwrap();
1879        let back_from_before_last = read_conditional(
1880            &mvcc,
1881            reader.events_tree_root_id,
1882            reader.tags_tree_root_id,
1883            qi,
1884            Some(Position(last - 1)),
1885            true,
1886            None,
1887        )
1888        .unwrap();
1889        let back_from_before_last_pos: Vec<u64> =
1890            back_from_before_last.iter().map(|e| e.position).collect();
1891        assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
1892    }
1893
1894    #[test]
1895    #[serial]
1896    fn tags_only_multi_tag_and_backwards() {
1897        let (_tmp, db, _input) = setup_db_with_standard_events();
1898        let reader = db.reader().unwrap();
1899        let qi = DCBQuery {
1900            items: vec![DCBQueryItem {
1901                types: vec![],
1902                tags: vec!["alpha".to_string(), "gamma".to_string()],
1903            }],
1904        };
1905
1906        // Forwards baseline
1907        let fwd = read_conditional(
1908            &db,
1909            reader.events_tree_root_id,
1910            reader.tags_tree_root_id,
1911            qi.clone(),
1912            Some(Position(1)),
1913            false,
1914            None,
1915        )
1916        .unwrap();
1917        let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
1918        assert!(!fwd_pos.is_empty());
1919        assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
1920        // All should include both tags
1921        assert!(
1922            fwd.iter()
1923                .all(|e| e.event.tags.iter().any(|t| t == "alpha"))
1924        );
1925        assert!(
1926            fwd.iter()
1927                .all(|e| e.event.tags.iter().any(|t| t == "gamma"))
1928        );
1929
1930        // Backwards from=None should equal reverse of forwards
1931        let back_all = read_conditional(
1932            &db,
1933            reader.events_tree_root_id,
1934            reader.tags_tree_root_id,
1935            qi.clone(),
1936            None,
1937            true,
1938            None,
1939        )
1940        .unwrap();
1941        let mut fwd_rev = fwd_pos.clone();
1942        fwd_rev.reverse();
1943        let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
1944        assert_eq!(back_all_pos, fwd_rev);
1945
1946        // Backwards from=last should still return full reverse (<= last)
1947        let last = *fwd_pos.last().unwrap();
1948        let back_from_last = read_conditional(
1949            &db,
1950            reader.events_tree_root_id,
1951            reader.tags_tree_root_id,
1952            qi.clone(),
1953            Some(Position(last)),
1954            true,
1955            None,
1956        )
1957        .unwrap();
1958        let back_from_last_pos: Vec<u64> = back_from_last.iter().map(|e| e.position).collect();
1959        assert_eq!(back_from_last_pos, fwd_rev);
1960
1961        // Backwards from just before last should drop newest
1962        let back_from_before_last = read_conditional(
1963            &db,
1964            reader.events_tree_root_id,
1965            reader.tags_tree_root_id,
1966            qi.clone(),
1967            Some(Position(last - 1)),
1968            true,
1969            None,
1970        )
1971        .unwrap();
1972        let back_from_before_last_pos: Vec<u64> =
1973            back_from_before_last.iter().map(|e| e.position).collect();
1974        assert_eq!(back_from_before_last_pos, fwd_rev[1..].to_vec());
1975
1976        // Backwards with limit
1977        let back_lim2 = read_conditional(
1978            &db,
1979            reader.events_tree_root_id,
1980            reader.tags_tree_root_id,
1981            qi,
1982            None,
1983            true,
1984            Some(2),
1985        )
1986        .unwrap();
1987        let back_lim2_pos: Vec<u64> = back_lim2.iter().map(|e| e.position).collect();
1988        assert_eq!(back_lim2_pos, fwd_rev[..2.min(fwd_rev.len())].to_vec());
1989    }
1990
1991    #[test]
1992    #[serial]
1993    fn tags_multi_item_two_tags_each_backwards() {
1994        let (_tmp, db, _input) = setup_db_with_standard_events();
1995        let reader = db.reader().unwrap();
1996        // Two items: (alpha & gamma) OR (beta & delta)
1997        let qi = DCBQuery {
1998            items: vec![
1999                DCBQueryItem {
2000                    types: vec![],
2001                    tags: vec!["alpha".to_string(), "gamma".to_string()],
2002                },
2003                DCBQueryItem {
2004                    types: vec![],
2005                    tags: vec!["beta".to_string(), "delta".to_string()],
2006                },
2007            ],
2008        };
2009
2010        // Forwards baseline
2011        let fwd = read_conditional(
2012            &db,
2013            reader.events_tree_root_id,
2014            reader.tags_tree_root_id,
2015            qi.clone(),
2016            Some(Position(1)),
2017            false,
2018            None,
2019        )
2020        .unwrap();
2021        let fwd_pos: Vec<u64> = fwd.iter().map(|e| e.position).collect();
2022        assert!(!fwd_pos.is_empty());
2023        assert!(fwd_pos.windows(2).all(|w| w[0] < w[1]));
2024        // Each event must satisfy one of the items fully
2025        assert!(fwd.iter().all(|e| {
2026            let tags = &e.event.tags;
2027            let has = |a: &str| tags.iter().any(|t| t == a);
2028            (has("alpha") && has("gamma")) || (has("beta") && has("delta"))
2029        }));
2030
2031        // Backwards with None should be reverse
2032        let back_all = read_conditional(
2033            &db,
2034            reader.events_tree_root_id,
2035            reader.tags_tree_root_id,
2036            qi.clone(),
2037            None,
2038            true,
2039            None,
2040        )
2041        .unwrap();
2042        let mut fwd_rev = fwd_pos.clone();
2043        fwd_rev.reverse();
2044        let back_all_pos: Vec<u64> = back_all.iter().map(|e| e.position).collect();
2045        assert_eq!(back_all_pos, fwd_rev);
2046
2047        // Backwards with limit 1 should return newest matching
2048        let back_lim1 = read_conditional(
2049            &db,
2050            reader.events_tree_root_id,
2051            reader.tags_tree_root_id,
2052            qi,
2053            None,
2054            true,
2055            Some(1),
2056        )
2057        .unwrap();
2058        assert_eq!(back_lim1.len(), 1);
2059        assert_eq!(back_lim1[0].position, *fwd_rev.first().unwrap());
2060    }
2061}