Skip to main content

khive_runtime/
event_cursor_walk.rs

1use khive_storage::event::{Event, EventFilter};
2use khive_storage::types::PageRequest;
3use khive_storage::StorageError;
4
5use crate::events_split::MAX_QUERY_EVENTS_PAGE_ROWS as TRANSPORT_PAGE_ROWS;
6
7/// Causes left to each caller to classify and describe for its own surface.
8#[derive(Debug)]
9pub enum EventCursorWalkError {
10    Storage(StorageError),
11    MissingBoundary,
12    DenseTimestampTie { page_limit: u32 },
13}
14
15/// Visit at most `max_rows` owned events in the store's descending page order.
16///
17/// Re-read timestamp boundaries at offset zero and deduplicate their IDs; widen
18/// pages only up to the event transport cap. A wider tie returns a typed cause.
19/// Reads are independent against a live event plane, not a shared snapshot.
20/// The callback can have received a prefix before a later read or dense tie fails.
21pub async fn visit_events_cursor_walk<F: FnMut(Event)>(
22    store: &dyn khive_storage::event::EventStore,
23    base_filter: &EventFilter,
24    page_size: u32,
25    max_rows: u64,
26    mut visit: F,
27) -> Result<u64, EventCursorWalkError> {
28    let mut admitted = 0;
29    let mut cursor: Option<i64> = base_filter.before;
30    let mut boundary_at: Option<i64> = None;
31    let mut boundary_ids: std::collections::HashSet<uuid::Uuid> = std::collections::HashSet::new();
32    let mut fetch_limit = page_size.clamp(1, TRANSPORT_PAGE_ROWS);
33    while admitted < max_rows {
34        let mut filter = base_filter.clone();
35        filter.before = cursor;
36        let page = store
37            .query_events(
38                filter,
39                PageRequest {
40                    offset: 0,
41                    limit: fetch_limit,
42                },
43            )
44            .await
45            .map_err(EventCursorWalkError::Storage)?;
46        let fetched = page.items.len() as u64;
47        let fresh: Vec<Event> = page
48            .items
49            .into_iter()
50            .filter(|event| !boundary_ids.contains(&event.id))
51            .collect();
52        if fresh.is_empty() {
53            if fetched < u64::from(fetch_limit) {
54                // The store returned everything under the cursor and all of
55                // it was already collected: the window is exhausted.
56                break;
57            }
58            // A full page of already-collected boundary rows: the tie run at
59            // this microsecond fills the page. Widen and re-read — but only
60            // up to the transport cap, past which the daemon refuses the
61            // request.
62            if fetch_limit >= TRANSPORT_PAGE_ROWS {
63                // At the cap, distinguish a tie run that exactly fills the
64                // page (fully collected, pageable by stepping the strict
65                // bound to the boundary itself) from one wider than the cap
66                // (genuinely unpageable with a timestamp cursor). Every
67                // collected row is >= the boundary microsecond, so equality
68                // of the at-or-above count with the collected count proves
69                // the run is complete.
70                let boundary = boundary_at.ok_or(EventCursorWalkError::MissingBoundary)?;
71                let mut ge_boundary = base_filter.clone();
72                // `after` is a strict `created_at >` bound, so at-or-above
73                // the boundary is `> boundary - 1`. At `i64::MIN` every row
74                // already satisfies at-or-above; keep the base bound.
75                ge_boundary.after = boundary.checked_sub(1).or(base_filter.after);
76                let ge_total = store
77                    .count_events(ge_boundary)
78                    .await
79                    .map_err(EventCursorWalkError::Storage)?;
80                if ge_total == admitted {
81                    cursor = Some(boundary);
82                    continue;
83                }
84                return Err(EventCursorWalkError::DenseTimestampTie {
85                    page_limit: fetch_limit,
86                });
87            }
88            fetch_limit = fetch_limit.saturating_mul(2).min(TRANSPORT_PAGE_ROWS);
89            continue;
90        }
91        // Pages come back created_at DESC, so the last fresh row carries the
92        // new boundary microsecond.
93        let boundary = fresh
94            .last()
95            .map(|event| event.created_at)
96            .expect("fresh is non-empty");
97        if boundary_at != Some(boundary) {
98            boundary_ids.clear();
99            boundary_at = Some(boundary);
100        }
101        boundary_ids.extend(
102            fresh
103                .iter()
104                .filter(|event| event.created_at == boundary)
105                .map(|event| event.id),
106        );
107        // Preserve the old final truncate: only the remaining prefix is
108        // admitted, so surplus payloads never reach the aggregate callback.
109        let remaining = usize::try_from(max_rows - admitted).unwrap_or(usize::MAX);
110        for event in fresh.into_iter().take(remaining) {
111            visit(event);
112            admitted += 1;
113        }
114        // `i64::MAX` admits no exclusive bound above it: keep the cursor as
115        // is and re-read — dedup drops the re-admitted rows, and the
116        // at-the-cap completeness check above advances past the boundary (or
117        // reports the dense tie) once a page comes back all-duplicates.
118        cursor = if boundary == i64::MAX {
119            cursor
120        } else {
121            Some(boundary + 1)
122        };
123        if fetched < u64::from(fetch_limit) {
124            break;
125        }
126    }
127    Ok(admitted)
128}