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        || !khive_types::is_lowercase_hex(fields[3])
256        || fields[4].len() != 64
257        || !khive_types::is_lowercase_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 hex(bytes: &[u8]) -> String {
288    const DIGITS: &[u8; 16] = b"0123456789abcdef";
289    let mut output = String::with_capacity(bytes.len() * 2);
290    for byte in bytes {
291        output.push(DIGITS[(byte >> 4) as usize] as char);
292        output.push(DIGITS[(byte & 15) as usize] as char);
293    }
294    output
295}
296
297fn encode_cursor(until_us: i64, key: &EventOrderKey, binding: &str) -> String {
298    format!(
299        "ep1:{until_us}:{}:{}:{binding}",
300        key.created_at_us,
301        hex(key.physical_id.as_bytes())
302    )
303}
304
305fn filter_binding(
306    token: &NamespaceToken,
307    since_us: i64,
308    until_us: i64,
309    kinds: &[EventKind],
310    actors: &[String],
311    namespaces: &[String],
312    exclusions: &[String],
313) -> RuntimeResult<String> {
314    let kinds = kinds.iter().map(|kind| kind.name()).collect::<Vec<_>>();
315    let bytes = serde_json::to_vec(&(
316        "event-page-v1",
317        &token.actor().kind,
318        &token.actor().id,
319        since_us,
320        until_us,
321        kinds,
322        actors,
323        namespaces,
324        exclusions,
325    ))
326    .map_err(|_| invalid("event page filter encoding failed"))?;
327    Ok(hex(&Sha256::digest(bytes)))
328}
329
330pub(crate) fn page_error(message: &str) -> StorageError {
331    StorageError::InvalidInput {
332        capability: StorageCapability::Events,
333        operation: "query_event_page".into(),
334        message: message.to_owned(),
335    }
336}
337
338pub(crate) fn compare_keys(a: &EventOrderKey, b: &EventOrderKey) -> Ordering {
339    a.created_at_us
340        .cmp(&b.created_at_us)
341        .then_with(|| a.physical_id.as_bytes().cmp(b.physical_id.as_bytes()))
342}
343
344fn physical_uuid(value: &str) -> Option<Uuid> {
345    matches!(value.len(), 32 | 36 | 38 | 45)
346        .then(|| Uuid::parse_str(value).ok())
347        .flatten()
348}
349
350pub(crate) fn validate_page_query(query: &EventPageQuery) -> StorageResult<()> {
351    if !(1..=crate::events_split::MAX_QUERY_EVENTS_PAGE_ROWS).contains(&query.max_rows)
352        || !valid_time(query.since_us)
353        || !valid_time(query.until_us)
354        || query.since_us >= query.until_us
355        || query.exclude_namespaces.len() > MAX_EXCLUSIONS
356        || query
357            .exclude_namespaces
358            .iter()
359            .any(|ns| Namespace::parse(ns).is_err())
360        || query.after.as_ref().is_some_and(|key| {
361            key.created_at_us < query.since_us
362                || key.created_at_us >= query.until_us
363                || physical_uuid(&key.physical_id).is_none()
364        })
365    {
366        return Err(page_error("invalid bounded event page query"));
367    }
368    Ok(())
369}
370
371pub(crate) fn validate_window(
372    query: &EventPageQuery,
373    namespace: Option<&str>,
374    window: &EventPageWindow,
375) -> StorageResult<()> {
376    if window.rows.len() > query.max_rows as usize {
377        return Err(page_error("event page row bound violated"));
378    }
379    let mut previous = query.after.as_ref();
380    for row in &window.rows {
381        let event = &row.event;
382        if row.order_key.created_at_us != event.created_at
383            || physical_uuid(&row.order_key.physical_id) != Some(event.id)
384            || namespace.is_some_and(|ns| ns != event.namespace)
385            || query.exclude_namespaces.contains(&event.namespace)
386            || (!query.kinds.is_empty() && !query.kinds.contains(&event.kind))
387            || (!query.actors.is_empty() && !query.actors.contains(&event.actor))
388            || event.created_at < query.since_us
389            || event.created_at >= query.until_us
390            || query
391                .after
392                .as_ref()
393                .is_some_and(|key| compare_keys(&row.order_key, key) != Ordering::Greater)
394            || previous.is_some_and(|key| compare_keys(&row.order_key, key) != Ordering::Greater)
395        {
396            return Err(page_error("event page row invariant violated"));
397        }
398        previous = Some(&row.order_key);
399    }
400    if let Some(stop) = &window.budget_stop {
401        if window.rows.len() >= query.max_rows as usize
402            || stop.created_at_us < query.since_us
403            || stop.created_at_us >= query.until_us
404            || physical_uuid(&stop.physical_id).is_none()
405            || previous.is_some_and(|key| compare_keys(stop, key) != Ordering::Greater)
406        {
407            return Err(page_error("event page row invariant violated"));
408        }
409    }
410    Ok(())
411}
412
413fn reject_duplicate_keys(rows: &[EventPageRow]) -> StorageResult<()> {
414    if rows
415        .windows(2)
416        .any(|pair| compare_keys(&pair[0].order_key, &pair[1].order_key) == Ordering::Equal)
417    {
418        return Err(page_error("event page duplicate ordering key"));
419    }
420    Ok(())
421}
422
423/// A stopped row shares its key with a returned row or with another stopped row:
424/// a strict cursor cannot represent a position between them.
425fn reject_stop_collisions(rows: &[EventPageRow], stops: &[EventOrderKey]) -> StorageResult<()> {
426    let collides = stops.iter().enumerate().any(|(index, stop)| {
427        stops[index + 1..]
428            .iter()
429            .any(|other| compare_keys(stop, other) == Ordering::Equal)
430            || rows
431                .iter()
432                .any(|row| compare_keys(&row.order_key, stop) == Ordering::Equal)
433    });
434    if collides {
435        return Err(page_error("event page duplicate ordering key"));
436    }
437    Ok(())
438}
439
440pub(crate) async fn split_page(
441    legacy: &dyn EventStore,
442    lane: &dyn EventStore,
443    query: EventPageQuery,
444) -> StorageResult<EventPageWindow> {
445    validate_page_query(&query)?;
446    let legacy = legacy.query_event_page(query.clone()).await?;
447    validate_window(&query, None, &legacy)?;
448    let lane = lane.query_event_page(query.clone()).await?;
449    validate_window(&query, None, &lane)?;
450    let stops = legacy
451        .budget_stop
452        .iter()
453        .chain(lane.budget_stop.iter())
454        .cloned()
455        .collect::<Vec<_>>();
456    let mut rows: Vec<EventPageRow> = legacy.rows;
457    rows.extend(lane.rows);
458    rows.sort_by(|a, b| compare_keys(&a.order_key, &b.order_key));
459    reject_stop_collisions(&rows, &stops)?;
460    let stop = stops.into_iter().min_by(compare_keys);
461    if let Some(stop) = &stop {
462        rows.truncate(
463            rows.partition_point(|row| compare_keys(&row.order_key, stop) == Ordering::Less),
464        );
465    }
466    reject_duplicate_keys(&rows)?;
467    let max_rows = query.max_rows as usize;
468    let budget_stop = if rows.len() >= max_rows { None } else { stop };
469    rows.truncate(max_rows);
470    Ok(EventPageWindow { rows, budget_stop })
471}
472
473struct ByteBudget(usize);
474
475impl Write for ByteBudget {
476    fn write(&mut self, bytes: &[u8]) -> std::io::Result<usize> {
477        if bytes.len() > self.0 {
478            return Err(std::io::Error::other("event page byte budget exceeded"));
479        }
480        self.0 -= bytes.len();
481        Ok(bytes.len())
482    }
483
484    fn flush(&mut self) -> std::io::Result<()> {
485        Ok(())
486    }
487}
488
489#[cfg(test)]
490#[path = "event_page_tests.rs"]
491mod tests;