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}