1use std::cmp::Ordering;
2use std::io::Write;
3
4use khive_storage::event::{EventOrderKey, EventPageQuery, EventPageRow, EventPageWindow};
5use khive_storage::{Event, EventStore, StorageCapability, StorageError, StorageResult};
6use khive_types::{Details, EventKind, KhiveError};
7use sha2::{Digest, Sha256};
8use uuid::Uuid;
9
10use crate::{KhiveRuntime, Namespace, NamespaceToken, RuntimeError, RuntimeResult};
11
12const MAX_NAMESPACES: usize = 16;
13const MAX_EXCLUSIONS: usize = 32;
14const MAX_LIMIT: u32 = 1000;
15const MAX_CURSOR_BYTES: usize = 512;
16const MAX_AGGREGATE_BYTES: usize = 32 * 1024 * 1024;
17
18#[derive(Clone, Debug)]
19pub struct EventReadPageRequest {
20 pub since_us: i64,
21 pub until_us: Option<i64>,
22 pub kinds: Vec<EventKind>,
23 pub actors: Vec<String>,
24 pub namespaces: Option<Vec<String>>,
25 pub exclude_namespaces: Vec<String>,
26 pub limit: u32,
27 pub after: Option<String>,
28}
29
30#[derive(Clone, Debug)]
31pub struct EventReadPageResult {
32 pub events: Vec<Event>,
33 pub has_more: bool,
34 pub next_after: Option<String>,
35 pub since_us: i64,
36 pub until_us: i64,
37 pub namespaces: Vec<String>,
38 pub single_row_cursor: Option<String>,
41}
42
43pub fn page_budget_exceeded_error() -> RuntimeError {
46 RuntimeError::Khive(
47 KhiveError::invalid_input(
48 "event page exceeds its byte budget; retry with a smaller `limit`",
49 )
50 .with_details(Details::new([("reason", "page_budget_exceeded")])),
51 )
52}
53
54pub fn row_exceeds_budget_error(event_id: Uuid, resume_after: String) -> RuntimeError {
57 RuntimeError::Khive(
58 KhiveError::invalid_input(
59 "an event exceeds the event page byte budget and cannot be returned in a page; \
60 continue from `resume_after` to skip it, or read it by id",
61 )
62 .with_details(Details::new_owned([
63 ("reason", "row_exceeds_budget".to_owned()),
64 ("event_id", event_id.to_string()),
65 ("resume_after", resume_after),
66 ])),
67 )
68}
69
70impl KhiveRuntime {
71 pub async fn page_events(
74 &self,
75 token: &NamespaceToken,
76 mut request: EventReadPageRequest,
77 ) -> RuntimeResult<EventReadPageResult> {
78 if !(1..=MAX_LIMIT).contains(&request.limit) {
79 return Err(invalid("event page limit must be between 1 and 1000"));
80 }
81 let requested = request
82 .namespaces
83 .take()
84 .unwrap_or_else(|| vec![token.namespace().as_str().to_owned()]);
85 let candidates = normalize_names(requested, MAX_NAMESPACES)?
86 .into_iter()
87 .filter(|name| {
88 token
89 .visible_namespaces()
90 .iter()
91 .any(|ns| ns.as_str() == name.as_str())
92 })
93 .collect::<Vec<_>>();
94 let exclusions = normalize_names(request.exclude_namespaces, MAX_EXCLUSIONS)?;
95 request.kinds.sort_by_key(|kind| kind.name());
96 request.kinds.dedup();
97 request.actors.sort();
98 request.actors.dedup();
99 let cursor = request.after.as_deref().map(decode_cursor).transpose()?;
100 let until_us = match (&cursor, request.until_us) {
101 (Some(cursor), Some(until)) if cursor.until_us != until => {
102 return Err(invalid_cursor());
103 }
104 (Some(cursor), _) => cursor.until_us,
105 (None, Some(until)) => until,
106 (None, None) => chrono::Utc::now().timestamp_micros(),
107 };
108 if !valid_time(request.since_us) || !valid_time(until_us) || request.since_us >= until_us {
109 return Err(invalid("invalid event page time window"));
110 }
111 let binding = filter_binding(
112 token,
113 request.since_us,
114 until_us,
115 &request.kinds,
116 &request.actors,
117 &candidates,
118 &exclusions,
119 )?;
120 if cursor.as_ref().is_some_and(|cursor| {
121 cursor.binding != binding
122 || cursor.key.created_at_us < request.since_us
123 || cursor.key.created_at_us >= until_us
124 }) {
125 return Err(invalid_cursor());
126 }
127 let query = EventPageQuery {
128 since_us: request.since_us,
129 until_us,
130 kinds: request.kinds,
131 actors: request.actors,
132 exclude_namespaces: exclusions.clone(),
133 after: cursor.map(|cursor| cursor.key),
134 max_rows: request.limit + 1,
135 };
136 let mut rows = Vec::new();
137 let mut stops = Vec::new();
138 let mut budget = ByteBudget(MAX_AGGREGATE_BYTES);
139 for namespace in &candidates {
140 let scoped = token.with_namespace(
142 Namespace::parse(namespace).map_err(|_| invalid("invalid namespace"))?,
143 );
144 let mut window = self
145 .events(&scoped)?
146 .query_event_page(query.clone())
147 .await?;
148 validate_window(&query, Some(namespace), &window)?;
149 serde_json::to_writer(&mut budget, &window.rows)
150 .map_err(|_| page_budget_exceeded_error())?;
151 stops.extend(window.budget_stop.take());
152 rows.extend(window.rows);
153 }
154 rows.sort_by(|a, b| compare_keys(&a.order_key, &b.order_key));
155 reject_duplicate_keys(&rows)?;
156 reject_stop_collisions(&rows, &stops)?;
157 let limit = request.limit as usize;
158 let stop = stops.into_iter().min_by(compare_keys);
159 if let Some(stop) = &stop {
160 let servable =
161 rows.partition_point(|row| compare_keys(&row.order_key, stop) == Ordering::Less);
162 rows.truncate(servable);
163 if servable == 0 {
164 let event_id = physical_uuid(&stop.physical_id)
165 .ok_or_else(|| page_error("event page row invariant violated"))?;
166 return Err(row_exceeds_budget_error(
167 event_id,
168 encode_cursor(until_us, stop, &binding),
169 ));
170 }
171 if servable < limit {
172 return Err(page_budget_exceeded_error());
173 }
174 }
175 let has_more = stop.is_some() || rows.len() > limit;
176 rows.truncate(limit);
177 let next_after = if has_more {
178 rows.last()
179 .map(|row| encode_cursor(until_us, &row.order_key, &binding))
180 } else {
181 None
182 };
183 let single_row_cursor = match rows.as_slice() {
184 [row] => Some(encode_cursor(until_us, &row.order_key, &binding)),
185 _ => None,
186 };
187 Ok(EventReadPageResult {
188 events: rows.into_iter().map(|row| row.event).collect(),
189 has_more,
190 next_after,
191 since_us: request.since_us,
192 until_us,
193 namespaces: candidates
194 .into_iter()
195 .filter(|namespace| !exclusions.contains(namespace))
196 .collect(),
197 single_row_cursor,
198 })
199 }
200}
201
202fn normalize_names(names: Vec<String>, cap: usize) -> RuntimeResult<Vec<String>> {
203 if names.len() > cap {
204 return Err(invalid("event page namespace list exceeds its bound"));
205 }
206 let mut names = names
207 .into_iter()
208 .map(|name| {
209 Namespace::parse(&name)
210 .map(|ns| ns.as_str().to_owned())
211 .map_err(|_| invalid("invalid event page namespace"))
212 })
213 .collect::<RuntimeResult<Vec<_>>>()?;
214 names.sort();
215 names.dedup();
216 Ok(names)
217}
218
219fn valid_time(time: i64) -> bool {
220 chrono::DateTime::<chrono::Utc>::from_timestamp_micros(time).is_some()
221}
222
223fn invalid(message: &str) -> RuntimeError {
224 RuntimeError::InvalidInput(message.to_owned())
225}
226
227fn invalid_cursor() -> RuntimeError {
228 invalid("invalid event page cursor or changed query scope")
229}
230
231struct Cursor {
232 until_us: i64,
233 key: EventOrderKey,
234 binding: String,
235}
236
237fn decode_cursor(raw: &str) -> RuntimeResult<Cursor> {
238 if raw.len() > MAX_CURSOR_BYTES {
239 return Err(invalid_cursor());
240 }
241 let fields = raw.split(':').collect::<Vec<_>>();
242 if fields.len() != 5 || fields[0] != "ep1" {
243 return Err(invalid_cursor());
244 }
245 let parse_time = |value: &str| -> RuntimeResult<i64> {
246 let time = value.parse::<i64>().map_err(|_| invalid_cursor())?;
247 if time.to_string() != value || !valid_time(time) {
248 return Err(invalid_cursor());
249 }
250 Ok(time)
251 };
252 let until_us = parse_time(fields[1])?;
253 let created_at_us = parse_time(fields[2])?;
254 if !matches!(fields[3].len(), 64 | 72 | 76 | 90)
255 || !lower_hex(fields[3])
256 || fields[4].len() != 64
257 || !lower_hex(fields[4])
258 {
259 return Err(invalid_cursor());
260 }
261 let id_bytes = fields[3]
262 .as_bytes()
263 .chunks_exact(2)
264 .map(|pair| {
265 let pair = std::str::from_utf8(pair).map_err(|_| invalid_cursor())?;
266 u8::from_str_radix(pair, 16).map_err(|_| invalid_cursor())
267 })
268 .collect::<RuntimeResult<Vec<_>>>()?;
269 let physical_id = String::from_utf8(id_bytes).map_err(|_| invalid_cursor())?;
270 if physical_uuid(&physical_id).is_none() {
271 return Err(invalid_cursor());
272 }
273 let key = EventOrderKey {
274 created_at_us,
275 physical_id,
276 };
277 if encode_cursor(until_us, &key, fields[4]) != raw {
278 return Err(invalid_cursor());
279 }
280 Ok(Cursor {
281 until_us,
282 key,
283 binding: fields[4].to_owned(),
284 })
285}
286
287fn lower_hex(value: &str) -> bool {
288 value
289 .bytes()
290 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
291}
292
293fn hex(bytes: &[u8]) -> String {
294 const DIGITS: &[u8; 16] = b"0123456789abcdef";
295 let mut output = String::with_capacity(bytes.len() * 2);
296 for byte in bytes {
297 output.push(DIGITS[(byte >> 4) as usize] as char);
298 output.push(DIGITS[(byte & 15) as usize] as char);
299 }
300 output
301}
302
303fn encode_cursor(until_us: i64, key: &EventOrderKey, binding: &str) -> String {
304 format!(
305 "ep1:{until_us}:{}:{}:{binding}",
306 key.created_at_us,
307 hex(key.physical_id.as_bytes())
308 )
309}
310
311fn filter_binding(
312 token: &NamespaceToken,
313 since_us: i64,
314 until_us: i64,
315 kinds: &[EventKind],
316 actors: &[String],
317 namespaces: &[String],
318 exclusions: &[String],
319) -> RuntimeResult<String> {
320 let kinds = kinds.iter().map(|kind| kind.name()).collect::<Vec<_>>();
321 let bytes = serde_json::to_vec(&(
322 "event-page-v1",
323 &token.actor().kind,
324 &token.actor().id,
325 since_us,
326 until_us,
327 kinds,
328 actors,
329 namespaces,
330 exclusions,
331 ))
332 .map_err(|_| invalid("event page filter encoding failed"))?;
333 Ok(hex(&Sha256::digest(bytes)))
334}
335
336pub(crate) fn page_error(message: &str) -> StorageError {
337 StorageError::InvalidInput {
338 capability: StorageCapability::Events,
339 operation: "query_event_page".into(),
340 message: message.to_owned(),
341 }
342}
343
344pub(crate) fn compare_keys(a: &EventOrderKey, b: &EventOrderKey) -> Ordering {
345 a.created_at_us
346 .cmp(&b.created_at_us)
347 .then_with(|| a.physical_id.as_bytes().cmp(b.physical_id.as_bytes()))
348}
349
350fn physical_uuid(value: &str) -> Option<Uuid> {
351 matches!(value.len(), 32 | 36 | 38 | 45)
352 .then(|| Uuid::parse_str(value).ok())
353 .flatten()
354}
355
356pub(crate) fn validate_page_query(query: &EventPageQuery) -> StorageResult<()> {
357 if !(1..=crate::events_split::MAX_QUERY_EVENTS_PAGE_ROWS).contains(&query.max_rows)
358 || !valid_time(query.since_us)
359 || !valid_time(query.until_us)
360 || query.since_us >= query.until_us
361 || query.exclude_namespaces.len() > MAX_EXCLUSIONS
362 || query
363 .exclude_namespaces
364 .iter()
365 .any(|ns| Namespace::parse(ns).is_err())
366 || query.after.as_ref().is_some_and(|key| {
367 key.created_at_us < query.since_us
368 || key.created_at_us >= query.until_us
369 || physical_uuid(&key.physical_id).is_none()
370 })
371 {
372 return Err(page_error("invalid bounded event page query"));
373 }
374 Ok(())
375}
376
377pub(crate) fn validate_window(
378 query: &EventPageQuery,
379 namespace: Option<&str>,
380 window: &EventPageWindow,
381) -> StorageResult<()> {
382 if window.rows.len() > query.max_rows as usize {
383 return Err(page_error("event page row bound violated"));
384 }
385 let mut previous = query.after.as_ref();
386 for row in &window.rows {
387 let event = &row.event;
388 if row.order_key.created_at_us != event.created_at
389 || physical_uuid(&row.order_key.physical_id) != Some(event.id)
390 || namespace.is_some_and(|ns| ns != event.namespace)
391 || query.exclude_namespaces.contains(&event.namespace)
392 || (!query.kinds.is_empty() && !query.kinds.contains(&event.kind))
393 || (!query.actors.is_empty() && !query.actors.contains(&event.actor))
394 || event.created_at < query.since_us
395 || event.created_at >= query.until_us
396 || query
397 .after
398 .as_ref()
399 .is_some_and(|key| compare_keys(&row.order_key, key) != Ordering::Greater)
400 || previous.is_some_and(|key| compare_keys(&row.order_key, key) != Ordering::Greater)
401 {
402 return Err(page_error("event page row invariant violated"));
403 }
404 previous = Some(&row.order_key);
405 }
406 if let Some(stop) = &window.budget_stop {
407 if window.rows.len() >= query.max_rows as usize
408 || stop.created_at_us < query.since_us
409 || stop.created_at_us >= query.until_us
410 || physical_uuid(&stop.physical_id).is_none()
411 || previous.is_some_and(|key| compare_keys(stop, key) != Ordering::Greater)
412 {
413 return Err(page_error("event page row invariant violated"));
414 }
415 }
416 Ok(())
417}
418
419fn reject_duplicate_keys(rows: &[EventPageRow]) -> StorageResult<()> {
420 if rows
421 .windows(2)
422 .any(|pair| compare_keys(&pair[0].order_key, &pair[1].order_key) == Ordering::Equal)
423 {
424 return Err(page_error("event page duplicate ordering key"));
425 }
426 Ok(())
427}
428
429fn reject_stop_collisions(rows: &[EventPageRow], stops: &[EventOrderKey]) -> StorageResult<()> {
432 let collides = stops.iter().enumerate().any(|(index, stop)| {
433 stops[index + 1..]
434 .iter()
435 .any(|other| compare_keys(stop, other) == Ordering::Equal)
436 || rows
437 .iter()
438 .any(|row| compare_keys(&row.order_key, stop) == Ordering::Equal)
439 });
440 if collides {
441 return Err(page_error("event page duplicate ordering key"));
442 }
443 Ok(())
444}
445
446pub(crate) async fn split_page(
447 legacy: &dyn EventStore,
448 lane: &dyn EventStore,
449 query: EventPageQuery,
450) -> StorageResult<EventPageWindow> {
451 validate_page_query(&query)?;
452 let legacy = legacy.query_event_page(query.clone()).await?;
453 validate_window(&query, None, &legacy)?;
454 let lane = lane.query_event_page(query.clone()).await?;
455 validate_window(&query, None, &lane)?;
456 let stops = legacy
457 .budget_stop
458 .iter()
459 .chain(lane.budget_stop.iter())
460 .cloned()
461 .collect::<Vec<_>>();
462 let mut rows: Vec<EventPageRow> = legacy.rows;
463 rows.extend(lane.rows);
464 rows.sort_by(|a, b| compare_keys(&a.order_key, &b.order_key));
465 reject_stop_collisions(&rows, &stops)?;
466 let stop = stops.into_iter().min_by(compare_keys);
467 if let Some(stop) = &stop {
468 rows.truncate(
469 rows.partition_point(|row| compare_keys(&row.order_key, stop) == Ordering::Less),
470 );
471 }
472 reject_duplicate_keys(&rows)?;
473 let max_rows = query.max_rows as usize;
474 let budget_stop = if rows.len() >= max_rows { None } else { stop };
475 rows.truncate(max_rows);
476 Ok(EventPageWindow { rows, budget_stop })
477}
478
479struct ByteBudget(usize);
480
481impl Write for ByteBudget {
482 fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
483 if bytes.len() > self.0 {
484 return Err(std::io::Error::other("event page byte budget exceeded"));
485 }
486 self.0 -= bytes.len();
487 Ok(bytes.len())
488 }
489
490 fn flush(&mut self) -> std::io::Result<()> {
491 Ok(())
492 }
493}
494
495#[cfg(test)]
496#[path = "event_page_tests.rs"]
497mod tests;