khive-db 0.10.0

SQLite storage backend: entities, edges, notes, events, FTS5, sqlite-vec vectors.
Documentation
use std::sync::Arc;

use khive_storage::event::{EventOrderKey, EventPageQuery, EventPageRow, EventPageWindow};
use khive_storage::{StorageCapability, StorageError};
use rusqlite::types::ValueRef;
use uuid::Uuid;

use super::{map_err, read_event, SqlEventStore};

const QUERY: &str = include_str!("../../sql/events-cursor-page.sql");
const MAX_ROWS: u32 = 4096;
const MAX_RAW_TEXT_BYTES: usize = 1024 * 1024;
const OPERATION: &str = "query_event_page";

fn invalid(message: &'static str) -> StorageError {
    StorageError::InvalidInput {
        capability: StorageCapability::Events,
        operation: OPERATION.into(),
        message: message.into(),
    }
}

fn physical_uuid(id: &str) -> Option<Uuid> {
    if !matches!(id.len(), 32 | 36 | 38 | 45) {
        return None;
    }
    Uuid::parse_str(id).ok()
}

fn validate(query: &EventPageQuery) -> Result<(), StorageError> {
    if !(1..=MAX_ROWS).contains(&query.max_rows) {
        return Err(invalid("event page max_rows must be in 1..=4096"));
    }
    if query.since_us >= query.until_us {
        return Err(invalid("event page requires since_us < until_us"));
    }
    if query.exclude_namespaces.len() > 32 {
        return Err(invalid(
            "event page permits at most 32 namespace exclusions",
        ));
    }
    if query.after.as_ref().is_some_and(|after| {
        after.created_at_us < query.since_us
            || after.created_at_us >= query.until_us
            || physical_uuid(&after.physical_id).is_none()
    }) {
        return Err(invalid("invalid event page ordering key"));
    }
    Ok(())
}

pub(super) async fn query_event_page(
    store: &SqlEventStore,
    query: EventPageQuery,
) -> Result<EventPageWindow, StorageError> {
    validate(&query)?;
    let namespace = store.namespace.clone();
    super::super::run_pooled_store_read(
        Arc::clone(&store.pool),
        StorageCapability::Events,
        OPERATION,
        move |conn| read_window(conn, &namespace, &query),
    )
    .await
}

fn read_window(
    conn: &rusqlite::Connection,
    namespace: &str,
    query: &EventPageQuery,
) -> Result<EventPageWindow, StorageError> {
    let kinds = serde_json::to_string(
        &query
            .kinds
            .iter()
            .map(|kind| kind.name())
            .collect::<Vec<_>>(),
    )
    .map_err(|error| StorageError::driver(StorageCapability::Events, OPERATION, error))?;
    let actors = serde_json::to_string(&query.actors)
        .map_err(|error| StorageError::driver(StorageCapability::Events, OPERATION, error))?;
    let exclusions = serde_json::to_string(&query.exclude_namespaces)
        .map_err(|error| StorageError::driver(StorageCapability::Events, OPERATION, error))?;
    let mut statement = conn
        .prepare(QUERY)
        .map_err(|error| map_err(error, OPERATION))?;
    let mut rows = statement
        .query(rusqlite::params![
            namespace,
            query.since_us,
            query.until_us,
            kinds,
            actors,
            exclusions,
            query.after.as_ref().map(|key| key.created_at_us),
            query.after.as_ref().map(|key| key.physical_id.as_str()),
            i64::from(query.max_rows),
        ])
        .map_err(|error| map_err(error, OPERATION))?;
    let mut result = Vec::new();
    let mut raw_bytes = 0usize;
    let mut budget_stop = None;
    while let Some(row) = rows.next().map_err(|error| map_err(error, OPERATION))? {
        let mut row_bytes = 0usize;
        for column in 0..18 {
            let bytes = match row
                .get_ref(column)
                .map_err(|error| map_err(error, OPERATION))?
            {
                ValueRef::Text(value) | ValueRef::Blob(value) => value.len(),
                _ => 0,
            };
            row_bytes = row_bytes.saturating_add(bytes);
        }
        if raw_bytes.saturating_add(row_bytes) > MAX_RAW_TEXT_BYTES {
            let physical_id: String = row.get(0).map_err(|error| map_err(error, OPERATION))?;
            physical_uuid(&physical_id)
                .ok_or_else(|| invalid("invalid stored event page ordering key"))?;
            let created_at_us: i64 = row.get(15).map_err(|error| map_err(error, OPERATION))?;
            budget_stop = Some(EventOrderKey {
                created_at_us,
                physical_id,
            });
            break;
        }
        raw_bytes += row_bytes;
        let physical_id: String = row.get(0).map_err(|error| map_err(error, OPERATION))?;
        let id = physical_uuid(&physical_id)
            .ok_or_else(|| invalid("invalid stored event page ordering key"))?;
        let created_at_us: i64 = row.get(15).map_err(|error| map_err(error, OPERATION))?;
        let event = read_event(row).map_err(|error| map_err(error, OPERATION))?;
        if event.id != id || event.created_at != created_at_us {
            return Err(invalid(
                "event page row does not match its stored ordering key",
            ));
        }
        result.push(EventPageRow {
            order_key: EventOrderKey {
                created_at_us,
                physical_id,
            },
            event,
        });
    }
    Ok(EventPageWindow {
        rows: result,
        budget_stop,
    })
}

#[cfg(test)]
#[path = "event_cursor_tests.rs"]
mod tests;