Skip to main content

khive_runtime/
event_page.rs

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    /// Cursor positioned at the returned row when the page holds exactly one row.
39    /// Only used to build a `row_exceeds_budget` refusal; never part of a response body.
40    pub single_row_cursor: Option<String>,
41}
42
43/// Refusal for a page whose rows together pass a byte budget. Carries no cursor:
44/// a smaller `limit` reaches the same rows without skipping any.
45pub 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
54/// Refusal for a single event whose row cannot be served within a byte budget.
55/// `resume_after` continues strictly after that event, which stays readable by id.
56pub 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    /// Read a live ordered window. Actor aliases are authorized by the caller;
72    /// this boundary rechecks namespace visibility and cursor scope on every page.
73    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            // with_namespace transfers capability; intersection above is the policy check.
141            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
429/// A stopped row shares its key with a returned row or with another stopped row:
430/// a strict cursor cannot represent a position between them.
431fn 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;