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