use std::path::Path;
use std::sync::{Arc, Mutex, PoisonError};
use async_trait::async_trait;
use eventsdb_core::store::EventStore as EventsdbStore;
use eventsdb_core::upcast::Current as Upcasted;
use eventsdb_sqlite::SqliteEventLog;
use serde_json::{Map, Value};
use super::event::{validate_event, FIELD_EPOCH_MS, FIELD_SEQ};
use super::event_store::{
stamp_schema_version, ChildScan, ChildrenDecision, Committed, Decision, EventStore, Split,
SplitDecision,
};
use super::logs::Logs;
use super::query::{QueryPlan, QueryRows};
use super::{KnlError, KnlResult};
pub const EVENTS_TABLE: &str = "events";
pub(super) const CHILD_INDEX_DDL: &str = "CREATE INDEX IF NOT EXISTS \
events_session_opened_parent \
ON events (json_extract(data, '$.parent')) \
WHERE kind = 'session_opened';";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SchemaColumn {
pub name: String,
pub declared_type: String,
pub pk: bool,
}
const EVENTS_COLUMNS: [(&str, &str, bool); 8] = [
("position", "INTEGER", true),
("stream", "TEXT", false),
("seq", "INTEGER", false),
("epoch_ms", "INTEGER", false),
("kind", "TEXT", false),
("schema_version", "INTEGER", false),
("meta", "TEXT", false),
("data", "TEXT", false),
];
pub fn events_schema() -> KnlResult<Vec<SchemaColumn>> {
Ok(EVENTS_COLUMNS
.iter()
.map(|(name, declared_type, pk)| SchemaColumn {
name: (*name).to_string(),
declared_type: (*declared_type).to_string(),
pk: *pk,
})
.collect())
}
pub struct SqliteEventStore {
log: Arc<SqliteEventLog>,
handle: eventsdb_sqlite::SqliteEventStore,
stream: String,
}
impl SqliteEventStore {
pub async fn open(path: &Path, stream: impl Into<String>, logs: &Logs) -> KnlResult<Self> {
Ok(Self::on(logs.file(path).await?, stream))
}
pub async fn open_memory(stream: impl Into<String>, logs: &Logs) -> KnlResult<Self> {
Ok(Self::on(logs.memory().await?, stream))
}
pub fn on(log: Arc<SqliteEventLog>, stream: impl Into<String>) -> Self {
let stream = stream.into();
let handle = log.stream_handle(&stream);
Self {
log,
handle,
stream,
}
}
}
fn owned_kinds(kinds: Option<&[&str]>) -> Option<Vec<String>> {
kinds.map(|kinds| kinds.iter().map(|kind| (*kind).to_string()).collect())
}
fn borrowed_kinds(kinds: &Option<Vec<String>>) -> Option<Vec<&str>> {
kinds
.as_ref()
.map(|kinds| kinds.iter().map(String::as_str).collect())
}
fn prepared(mut event: Map<String, Value>) -> Map<String, Value> {
event.remove(FIELD_SEQ);
event.remove(FIELD_EPOCH_MS);
stamp_schema_version(&mut event);
event
}
fn committed_of(committed: eventsdb_core::position::Committed) -> Committed {
Committed {
seq: committed.seq,
epoch_ms: committed.epoch_ms,
}
}
fn values_of(events: Vec<Upcasted>) -> Vec<Value> {
events
.into_iter()
.map(|event| Value::Object(event.into_inner()))
.collect()
}
fn values_of_ref(events: &[Upcasted]) -> Vec<Value> {
events
.iter()
.map(|event| Value::Object((**event).clone()))
.collect()
}
type Parked = Arc<Mutex<Option<KnlError>>>;
fn taken(parked: &Parked) -> Option<KnlError> {
parked.lock().unwrap_or_else(PoisonError::into_inner).take()
}
fn park(parked: &Parked, error: KnlError) {
*parked.lock().unwrap_or_else(PoisonError::into_inner) = Some(error);
}
fn rolled_back() -> eventsdb_core::Error {
eventsdb_core::Error::validation("the kernel refused an event this write was to record")
}
#[async_trait]
impl EventStore for SqliteEventStore {
async fn append(&mut self, event: Map<String, Value>) -> KnlResult<Committed> {
validate_event(&event)?;
self.handle
.append(prepared(event))
.await
.map(committed_of)
.map_err(KnlError::from)
}
async fn append_many(&mut self, events: Vec<Map<String, Value>>) -> KnlResult<Vec<Committed>> {
for event in &events {
validate_event(event)?;
}
let events: Vec<Map<String, Value>> = events.into_iter().map(prepared).collect();
self.handle
.append_many(events)
.await
.map(|committed| committed.into_iter().map(committed_of).collect())
.map_err(KnlError::from)
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: Decision,
) -> KnlResult<Option<Committed>> {
let parked: Parked = Arc::default();
let sink = Arc::clone(&parked);
let answer: eventsdb_core::store::Decision = Box::new(move |seen: &[Upcasted]| {
let event = decide(values_of_ref(seen))?;
match validate_event(&event) {
Ok(()) => Some(prepared(event)),
Err(refusal) => {
park(&sink, refusal);
None
}
}
});
let committed = self.handle.append_if(kinds, answer).await;
match taken(&parked) {
Some(refusal) => Err(refusal),
None => committed
.map(|committed| committed.map(committed_of))
.map_err(KnlError::from),
}
}
async fn append_if_many(
&mut self,
other: &str,
kinds: Option<&[&str]>,
decide: SplitDecision,
) -> KnlResult<Option<Split<Committed>>> {
let stream = self.stream.clone();
let other = other.to_string();
let kinds = owned_kinds(kinds);
let parked: Parked = Arc::default();
let sink = Arc::clone(&parked);
let committed = self
.log
.with_transaction(move |tx| {
let selection = borrowed_kinds(&kinds);
let seen = Split {
own: values_of(tx.read(&stream, selection.as_deref(), 0, usize::MAX)?),
other: values_of(tx.read(&other, None, 0, 1)?),
};
let Some(split) = decide(seen) else {
return Ok(None);
};
for event in split.own.iter().chain(split.other.iter()) {
if let Err(refusal) = validate_event(event) {
park(&sink, refusal);
return Err(rolled_back());
}
}
let own = tx.append_many(
&stream,
split.own.into_iter().map(prepared).collect::<Vec<_>>(),
)?;
let elsewhere = tx.append_many(
&other,
split.other.into_iter().map(prepared).collect::<Vec<_>>(),
)?;
Ok(Some(Split {
own: own.into_iter().map(committed_of).collect(),
other: elsewhere.into_iter().map(committed_of).collect(),
}))
})
.await;
match taken(&parked) {
Some(refusal) => Err(refusal),
None => committed.map_err(KnlError::from),
}
}
async fn append_with_open_children(
&mut self,
scan: &ChildScan,
decide: ChildrenDecision,
) -> KnlResult<Committed> {
let stream = self.stream.clone();
let scan = scan.clone();
let parked: Parked = Arc::default();
let sink = Arc::clone(&parked);
let committed = self
.log
.with_transaction(move |tx| {
let conn: &rusqlite::Connection = tx;
let children = open_children_in(conn, &stream, &scan)
.map_err(|e| eventsdb_core::Error::storage(e.to_string()))?;
let event = decide(children);
if let Err(refusal) = validate_event(&event) {
park(&sink, refusal);
return Err(rolled_back());
}
tx.append(&stream, prepared(event)).map(committed_of)
})
.await;
match taken(&parked) {
Some(refusal) => Err(refusal),
None => committed.map_err(KnlError::from),
}
}
fn database(&self) -> Option<&str> {
Some(self.log.database())
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
self.handle
.read_kinds(kinds, from_seq, limit)
.await
.map(values_of)
.map_err(KnlError::from)
}
async fn read_last(&self, n: usize) -> KnlResult<Vec<Value>> {
self.handle
.read_last(n)
.await
.map(values_of)
.map_err(KnlError::from)
}
async fn head(&self) -> KnlResult<Option<u64>> {
self.handle.head().await.map_err(KnlError::from)
}
async fn len(&self) -> KnlResult<usize> {
self.handle.len().await.map_err(KnlError::from)
}
async fn query(&self, plan: &QueryPlan) -> KnlResult<QueryRows> {
let cap = i64::try_from(plan.limit)
.unwrap_or(i64::MAX)
.saturating_add(1);
let sql = format!("SELECT * FROM ({}) LIMIT {cap}", plan.sql);
let rows = self
.log
.query_timeout(&sql, plan.values.clone(), plan.timeout)
.await
.map_err(query_error)?;
let truncated = rows.len() > plan.limit;
Ok(QueryRows {
rows: rows
.into_iter()
.take(plan.limit)
.map(|row| {
row.into_iter()
.filter(|(_, value)| !value.is_null())
.collect()
})
.collect(),
truncated,
})
}
fn detach_append(&self, event: Map<String, Value>) {
if let Err(e) = validate_event(&event) {
tracing::warn!(error = %e, "knl: a detached append was refused before it was submitted");
return;
}
if let Err(e) = self.log.detach_append(&self.stream, prepared(event)) {
tracing::warn!(error = %e, "knl: a detached append was not accepted by the log");
}
}
}
fn sql_literal(text: &str) -> String {
format!("'{}'", text.replace('\'', "''"))
}
fn child_scan_sql(scan: &ChildScan) -> String {
let opened = sql_literal(&scan.opened);
let closed = sql_literal(&scan.closed);
let path = sql_literal(&format!("$.{}", scan.parent_field));
format!(
"SELECT opened.stream \
FROM events AS opened \
WHERE opened.kind = {opened} \
AND json_extract(opened.data, {path}) = ?1 \
AND NOT EXISTS ( \
SELECT 1 FROM events AS ending \
WHERE ending.stream = opened.stream AND ending.kind = {closed} \
) \
ORDER BY opened.epoch_ms, opened.stream"
)
}
fn open_children_in(
conn: &rusqlite::Connection,
stream: &str,
scan: &ChildScan,
) -> rusqlite::Result<Vec<String>> {
let mut stmt = conn.prepare(&child_scan_sql(scan))?;
let rows = stmt.query_map(rusqlite::params![stream], |row| row.get::<_, String>(0))?;
rows.collect()
}
fn is_retryable(error: &rusqlite::Error) -> bool {
matches!(
error,
rusqlite::Error::SqliteFailure(inner, _)
if matches!(
inner.code,
rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked
)
)
}
impl From<rusqlite::Error> for KnlError {
fn from(error: rusqlite::Error) -> Self {
if is_retryable(&error) {
return KnlError::Busy(format!("sqlite: busy/locked: {error}"));
}
KnlError::Storage(format!("sqlite: {error}"))
}
}
impl From<eventsdb_core::Error> for KnlError {
fn from(error: eventsdb_core::Error) -> Self {
use eventsdb_core::Error as Failure;
match error {
Failure::Validation(reason) => KnlError::Validation(reason),
Failure::Busy(reason) => KnlError::Busy(reason),
Failure::Timeout(reason) => KnlError::Timeout(reason),
Failure::Storage(reason) => KnlError::Storage(reason),
Failure::Corruption(reason) => KnlError::Corruption(reason),
Failure::Unsupported(reason) => KnlError::Unsupported(reason),
other => KnlError::Storage(other.to_string()),
}
}
}
const SQLITE_STATEMENT_ERRORS: [&str; 7] = [
"no such column",
"no such table",
"no such function",
"syntax error",
"near \"",
"wrong number of arguments",
"ambiguous column name",
];
fn is_statement_error(message: &str) -> bool {
SQLITE_STATEMENT_ERRORS
.iter()
.any(|phrase| message.contains(phrase))
}
fn query_error(error: eventsdb_core::Error) -> KnlError {
match KnlError::from(error) {
KnlError::Storage(reason) if is_statement_error(&reason) => {
KnlError::Validation(format!("sql: {reason}"))
}
other => other,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::knl::event::{kind_of, seq_of, FIELD_DATA, FIELD_META};
use crate::knl::query::{self, QueryOpts, QueryParams};
use crate::knl::CURRENT_SCHEMA_VERSION;
use serde_json::json;
fn obj(value: Value) -> Map<String, Value> {
match value {
Value::Object(map) => map,
other => panic!("test fixture must be an object, got {other}"),
}
}
fn ev(i: usize) -> Map<String, Value> {
obj(json!({ "kind": format!("e{i}") }))
}
fn budget(kind: &str, amount: i64) -> Map<String, Value> {
obj(json!({ "kind": kind, "data": { "amount": amount } }))
}
async fn mem_store() -> (SqliteEventStore, Logs) {
let logs = Logs::new();
let store = SqliteEventStore::open_memory(uuid::Uuid::new_v4().to_string(), &logs)
.await
.expect("open");
(store, logs)
}
fn decide(
f: impl FnOnce(Vec<Value>) -> Option<Map<String, Value>> + Send + 'static,
) -> Decision {
Box::new(f)
}
#[tokio::test]
async fn append_assigns_gap_free_monotonic_seq_from_one() {
let (mut store, _logs) = mem_store().await;
assert!(store.is_empty().await.expect("is_empty"));
assert_eq!(store.len().await.expect("len"), 0);
let a = store.append(ev(1)).await.expect("append e1");
let b = store.append(ev(2)).await.expect("append e2");
let c = store.append(ev(3)).await.expect("append e3");
assert_eq!((a.seq, b.seq, c.seq), (1, 2, 3));
assert_eq!(store.len().await.expect("len"), 3);
assert!(!store.is_empty().await.expect("is_empty"));
let stored = store.read(0, usize::MAX).await.expect("read");
let stored_epoch = stored[0]
.get("epoch_ms")
.and_then(Value::as_u64)
.expect("epoch is on the stored event");
assert_eq!(stored_epoch, a.epoch_ms);
}
#[tokio::test]
async fn a_rejected_append_records_nothing_and_burns_no_seq() {
let (mut store, _logs) = mem_store().await;
store
.append(obj(json!({ "text": "no kind" })))
.await
.expect_err("kind is required");
assert_eq!(store.len().await.expect("len"), 0);
assert_eq!(store.append(ev(1)).await.expect("append").seq, 1);
}
#[tokio::test]
async fn a_caller_supplied_coordinate_is_overwritten() {
let (mut store, _logs) = mem_store().await;
store.append(ev(1)).await.expect("seed");
let committed = store
.append(obj(json!({ "kind": "e2", "seq": 99, "epoch_ms": 7 })))
.await
.expect("append");
assert_eq!(committed.seq, 2, "the store numbers the stream");
assert_ne!(committed.epoch_ms, 7, "the store reads the clock");
let stored = store.read(0, usize::MAX).await.expect("read");
assert_eq!(seq_of(&stored[1]), 2);
assert_eq!(
stored[1].get("_schema_version").and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"every append is stamped with the kernel's version"
);
}
#[tokio::test]
async fn append_if_decides_inside_the_transaction_and_writes_only_a_some() {
let (mut store, _logs) = mem_store().await;
store.append(ev(1)).await.expect("seed");
let seen_kinds: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen_kinds);
let committed = store
.append_if(
None,
decide(move |events| {
*recorded.lock().expect("not poisoned") =
events.iter().map(|e| kind_of(e).to_string()).collect();
Some(ev(2))
}),
)
.await
.expect("append_if");
assert_eq!(
*seen_kinds.lock().expect("not poisoned"),
["e1"],
"decide saw the durable stream"
);
assert_eq!(committed.map(|c| c.seq), Some(2));
let nothing = store
.append_if(None, decide(|_| None))
.await
.expect("append_if");
assert_eq!(nothing, None);
assert_eq!(store.len().await.expect("len"), 2, "a None commits nothing");
assert_eq!(store.append(ev(3)).await.expect("append").seq, 3);
}
#[tokio::test]
async fn append_if_validates_the_event_the_decision_returns() {
let (mut store, _logs) = mem_store().await;
let err = store
.append_if(None, decide(|_| Some(obj(json!({ "text": "no kind" })))))
.await
.expect_err("kind is required");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert_eq!(store.len().await.expect("len"), 0);
}
#[tokio::test]
async fn append_if_holds_a_kernel_kind_to_its_data() {
let (mut store, _logs) = mem_store().await;
let err = store
.append_if(
None,
decide(|_| Some(obj(json!({ "kind": "budget_spent", "data": {} })))),
)
.await
.expect_err("a kernel kind needs its data");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert!(err.reason().contains("amount"), "{}", err.reason());
assert_eq!(store.len().await.expect("len"), 0);
}
#[tokio::test]
async fn append_many_is_one_transaction_that_lands_whole_or_not_at_all() {
let (mut store, _logs) = mem_store().await;
store.append(ev(1)).await.expect("seed");
let committed = store
.append_many(vec![ev(2), ev(3)])
.await
.expect("the batch");
assert_eq!(
committed.iter().map(|c| c.seq).collect::<Vec<_>>(),
[2, 3],
"numbered on from the head that was there"
);
let stored = store.read(0, usize::MAX).await.expect("read");
let kinds: Vec<&str> = stored.iter().map(kind_of).collect();
assert_eq!(kinds, ["e1", "e2", "e3"]);
store
.append_many(vec![ev(4), obj(json!({ "text": "no kind" }))])
.await
.expect_err("kind is required");
assert_eq!(
store.len().await.expect("len"),
3,
"a batch that fails lands nothing"
);
}
#[tokio::test]
async fn append_if_many_writes_both_streams_or_neither() {
let logs = Logs::new();
let parent = uuid::Uuid::new_v4().to_string();
let child = uuid::Uuid::new_v4().to_string();
let log = logs.memory().await.expect("the log");
let mut ledger = SqliteEventStore::on(Arc::clone(&log), parent.clone());
let opened = SqliteEventStore::on(Arc::clone(&log), child.clone());
assert_eq!(
ledger.database(),
opened.database(),
"both streams are in one database"
);
ledger
.append(budget("budget_granted", 100))
.await
.expect("the grant");
let seen: Arc<Mutex<(usize, usize)>> = Arc::default();
let recorded = Arc::clone(&seen);
let child_stream = child.clone();
let committed = ledger
.append_if_many(
&child,
Some(&["budget_granted"]),
Box::new(move |split: Split<Value>| {
*recorded.lock().expect("not poisoned") = (split.own.len(), split.other.len());
Some(Split {
own: vec![obj(json!({
"kind": "budget_reserved",
"data": { "amount": 10, "child": child_stream },
}))],
other: vec![
obj(json!({
"kind": "session_opened",
"data": { "scope_id": "sc-1", "owner": "o", "parent": "p" },
})),
budget("budget_granted", 10),
],
})
}),
)
.await
.expect("append_if_many")
.expect("a Some writes");
assert_eq!(
*seen.lock().expect("not poisoned"),
(1, 0),
"its own kinds, and an empty other stream"
);
assert_eq!(committed.own.iter().map(|c| c.seq).collect::<Vec<_>>(), [2]);
assert_eq!(
committed.other.iter().map(|c| c.seq).collect::<Vec<_>>(),
[1, 2],
"the other stream is numbered from its own head"
);
let nothing = ledger
.append_if_many(&child, None, Box::new(|_| None))
.await
.expect("append_if_many");
assert_eq!(nothing, None);
assert_eq!(ledger.len().await.expect("len"), 2);
assert_eq!(opened.len().await.expect("len"), 2);
let err = ledger
.append_if_many(
&child,
None,
Box::new(|_| {
Some(Split {
own: vec![budget("budget_spent", 1)],
other: vec![obj(json!({ "text": "no kind" }))],
})
}),
)
.await
.expect_err("kind is required");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert_eq!(ledger.len().await.expect("len"), 2);
assert_eq!(opened.len().await.expect("len"), 2);
}
#[tokio::test]
async fn the_child_scan_reads_by_the_parent_index() {
let logs = Logs::new();
let log = logs.memory().await.expect("the log");
let scan = ChildScan {
opened: "session_opened".to_string(),
closed: "session_closed".to_string(),
parent_field: "parent".to_string(),
};
let rows = log
.query(
&format!("EXPLAIN QUERY PLAN {}", child_scan_sql(&scan)),
vec![Value::from("p-1")],
)
.await
.expect("the plan");
let plan: String = rows
.iter()
.filter_map(|row| row.get("detail").and_then(Value::as_str))
.collect::<Vec<_>>()
.join(" | ");
assert!(
plan.contains("events_session_opened_parent"),
"the scan must read by the parent index: {plan}"
);
}
#[test]
fn a_word_written_into_the_scan_stays_one_word() {
let sql = child_scan_sql(&ChildScan {
opened: "it's opened".to_string(),
closed: "it's closed".to_string(),
parent_field: "it's parent".to_string(),
});
assert!(sql.contains("'it''s opened'"), "{sql}");
assert!(sql.contains("'it''s closed'"), "{sql}");
assert!(sql.contains("'$.it''s parent'"), "{sql}");
assert_eq!(
sql.matches('\'').count() % 2,
0,
"every literal is closed: {sql}"
);
}
#[tokio::test]
async fn database_is_the_same_for_two_streams_of_one_database() {
let logs = Logs::new();
let dir = tempfile::tempdir().expect("tempdir");
let here = dir.path().join("knl.db");
let there = dir.path().join("other.db");
let a = SqliteEventStore::open(&here, "s-1", &logs)
.await
.expect("open a");
let b = SqliteEventStore::open(&here, "s-2", &logs)
.await
.expect("open b");
let elsewhere = SqliteEventStore::open(&there, "s-1", &logs)
.await
.expect("open elsewhere");
assert_eq!(a.database(), b.database());
assert_ne!(a.database(), elsewhere.database());
let (mem, _mem_logs) = mem_store().await;
assert_ne!(mem.database(), a.database());
assert!(mem.database().is_some());
}
#[tokio::test]
async fn open_children_are_the_unended_streams_that_name_this_one() {
let logs = Logs::new();
let log = logs.memory().await.expect("the log");
let scan = ChildScan {
opened: "session_opened".to_string(),
closed: "session_closed".to_string(),
parent_field: "parent".to_string(),
};
async fn opened(log: &Arc<SqliteEventLog>, id: &str, parent: &str) -> SqliteEventStore {
let mut store = SqliteEventStore::on(Arc::clone(log), id);
store
.append(obj(json!({
"kind": "session_opened",
"data": { "scope_id": "sc", "owner": "o", "parent": parent },
})))
.await
.expect("the opening");
store
}
let mut parent = SqliteEventStore::on(Arc::clone(&log), "parent");
let _open_child = opened(&log, "child-open", "parent").await;
let mut ended = opened(&log, "child-ended", "parent").await;
let _elsewhere = opened(&log, "child-of-other", "another").await;
ended
.append(obj(json!({
"kind": "session_closed",
"data": { "reason": "done" },
})))
.await
.expect("the ending");
let recorded: Arc<Mutex<Vec<String>>> = Arc::default();
let seen = Arc::clone(&recorded);
let committed = parent
.append_with_open_children(
&scan,
Box::new(move |children| {
*seen.lock().expect("not poisoned") = children.clone();
obj(json!({
"kind": "session_closed",
"data": { "reason": "done", "open_children": children },
}))
}),
)
.await
.expect("the close");
assert_eq!(committed.seq, 1, "the close is the parent's first event");
assert_eq!(
*recorded.lock().expect("not poisoned"),
["child-open"],
"only this stream's children, and only the open ones"
);
}
#[tokio::test]
async fn read_kinds_selects_by_kind_and_keeps_the_streams_order() {
let (mut store, _logs) = mem_store().await;
store
.append(budget("budget_granted", 100))
.await
.expect("grant");
store.append(ev(1)).await.expect("noise");
store
.append(budget("budget_spent", 10))
.await
.expect("spend");
let all = store.read(0, usize::MAX).await.expect("read");
assert_eq!(
all.iter().map(kind_of).collect::<Vec<_>>(),
["budget_granted", "e1", "budget_spent"]
);
let ledger = store
.read_kinds(Some(&["budget_granted", "budget_spent"]), 0, usize::MAX)
.await
.expect("read_kinds");
assert_eq!(
ledger.iter().map(seq_of).collect::<Vec<_>>(),
[1, 3],
"the seq the stream gave them, not a fresh numbering"
);
let nothing = store
.read_kinds(Some(&[]), 0, usize::MAX)
.await
.expect("read_kinds");
assert!(nothing.is_empty(), "an empty selection selects nothing");
}
#[tokio::test]
async fn append_if_filters_the_decisions_input_and_numbers_against_the_stream() {
let (mut store, _logs) = mem_store().await;
store
.append(budget("budget_granted", 100))
.await
.expect("grant");
store.append(ev(1)).await.expect("noise");
let seen: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen);
let committed = store
.append_if(
Some(&["budget_granted", "budget_spent"]),
decide(move |events| {
*recorded.lock().expect("not poisoned") =
events.iter().map(|e| kind_of(e).to_string()).collect();
Some(budget("budget_spent", 10))
}),
)
.await
.expect("append_if");
assert_eq!(
*seen.lock().expect("not poisoned"),
["budget_granted"],
"only the kinds asked for"
);
assert_eq!(
committed.map(|c| c.seq),
Some(3),
"numbered against the whole stream"
);
}
#[tokio::test]
async fn append_if_across_two_handles_decides_on_the_other_handles_write() {
let logs = Logs::new();
let log = logs.memory().await.expect("the log");
let stream = uuid::Uuid::new_v4().to_string();
let mut a = SqliteEventStore::on(Arc::clone(&log), stream.clone());
let mut b = SqliteEventStore::on(Arc::clone(&log), stream);
fn claim() -> Decision {
decide(|events: Vec<Value>| {
if events.is_empty() {
Some(obj(json!({ "kind": "claim" })))
} else {
None
}
})
}
assert!(a
.append_if(Some(&["claim"]), claim())
.await
.expect("a")
.is_some());
assert!(
b.append_if(Some(&["claim"]), claim())
.await
.expect("b")
.is_none(),
"the second handle decided against what the first wrote"
);
assert_eq!(a.len().await.expect("len"), 1);
}
#[tokio::test]
async fn read_pages_by_from_seq_and_limit() {
let (mut store, _logs) = mem_store().await;
for i in 1..=5 {
store.append(ev(i)).await.expect("append");
}
let page = store.read(2, 2).await.expect("read");
assert_eq!(page.iter().map(seq_of).collect::<Vec<_>>(), [2, 3]);
let rest = store.read(4, usize::MAX).await.expect("read");
assert_eq!(rest.iter().map(seq_of).collect::<Vec<_>>(), [4, 5]);
assert!(store.read(6, 10).await.expect("read").is_empty());
}
#[tokio::test]
async fn read_last_takes_the_end_of_the_stream_in_seq_order() {
let (mut store, _logs) = mem_store().await;
for i in 1..=5 {
store.append(ev(i)).await.expect("append");
}
let tail = store.read_last(2).await.expect("read_last");
assert_eq!(tail.iter().map(seq_of).collect::<Vec<_>>(), [4, 5]);
assert!(store.read_last(0).await.expect("read_last").is_empty());
assert_eq!(store.read_last(50).await.expect("read_last").len(), 5);
}
#[tokio::test]
async fn head_is_none_when_empty_then_tracks_the_max() {
let (mut store, _logs) = mem_store().await;
assert_eq!(store.head().await.expect("head"), None);
store.append(ev(1)).await.expect("append");
assert_eq!(store.head().await.expect("head"), Some(1));
store.append(ev(2)).await.expect("append");
assert_eq!(store.head().await.expect("head"), Some(2));
}
#[tokio::test]
async fn read_reconstructs_the_written_event_out_of_its_columns() {
let (mut store, _logs) = mem_store().await;
let committed = store
.append(obj(json!({
"kind": "llm_response",
"meta": { "beat": "b-1", "attempt": 2, "final": true },
"data": { "content": { "text": "hi" }, "usage": { "input_tokens": 3 } },
})))
.await
.expect("append");
let stored = store.read(0, usize::MAX).await.expect("read");
let event = &stored[0];
assert_eq!(kind_of(event), "llm_response");
assert_eq!(seq_of(event), committed.seq);
assert_eq!(
event.get(FIELD_META),
Some(&json!({ "beat": "b-1", "attempt": 2, "final": true }))
);
assert_eq!(
event.get(FIELD_DATA),
Some(&json!({ "content": { "text": "hi" }, "usage": { "input_tokens": 3 } }))
);
assert_eq!(
event.get("_schema_version").and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION)
);
}
#[tokio::test]
async fn the_beat_is_a_meta_label_with_an_index() {
let (mut store, _logs) = mem_store().await;
store
.append(obj(json!({ "kind": "e1", "meta": { "beat": "b-1" } })))
.await
.expect("append");
store.append(ev(2)).await.expect("append");
let rows = ask(
&store,
"SELECT json_extract(meta, '$.beat') AS beat FROM events \
WHERE stream = $stream ORDER BY seq",
)
.await
.expect("query");
assert_eq!(rows.rows[0].get("beat"), Some(&json!("b-1")));
assert!(
!rows.rows[1].contains_key("beat"),
"an undeclared beat reads as nil: {:?}",
rows.rows[1]
);
let indexes = ask(
&store,
"SELECT name FROM sqlite_master WHERE type = 'index' AND tbl_name = 'events'",
)
.await
.expect("query");
let names: Vec<&str> = indexes
.rows
.iter()
.filter_map(|row| row.get("name").and_then(Value::as_str))
.collect();
assert!(
names.contains(&"events_meta_beat"),
"the beat label is indexed: {names:?}"
);
assert!(
!names.contains(&"events_stream_beat_seq"),
"and the column's old index is gone: {names:?}"
);
}
#[tokio::test]
async fn events_persist_across_a_reopen_of_the_same_path_and_stream() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
let stream = "s-1";
{
let logs = Logs::new();
let mut store = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("open");
store.append(ev(1)).await.expect("append");
store.append(ev(2)).await.expect("append");
assert!(logs.shutdown().await.is_empty(), "the log closed cleanly");
}
let logs = Logs::new();
let reopened = SqliteEventStore::open(&path, stream, &logs)
.await
.expect("reopen");
let stored = reopened.read(0, usize::MAX).await.expect("read");
assert_eq!(stored.iter().map(kind_of).collect::<Vec<_>>(), ["e1", "e2"]);
assert_eq!(reopened.head().await.expect("head"), Some(2));
}
#[tokio::test]
async fn two_streams_in_one_db_file_do_not_see_each_others_events() {
let logs = Logs::new();
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
let mut a = SqliteEventStore::open(&path, "s-a", &logs)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "s-b", &logs)
.await
.expect("open b");
a.append(ev(1)).await.expect("a");
b.append(ev(2)).await.expect("b");
assert_eq!(
a.read(0, usize::MAX)
.await
.expect("read a")
.iter()
.map(kind_of)
.collect::<Vec<_>>(),
["e1"]
);
assert_eq!(
b.read(0, usize::MAX)
.await
.expect("read b")
.iter()
.map(kind_of)
.collect::<Vec<_>>(),
["e2"]
);
assert_eq!(a.head().await.expect("head"), Some(1));
assert_eq!(b.head().await.expect("head"), Some(1));
}
#[tokio::test]
async fn read_errors_on_a_corrupt_row_instead_of_dropping_it() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
{
let logs = Logs::new();
let mut store = SqliteEventStore::open(&path, "s-1", &logs)
.await
.expect("open");
store.append(ev(1)).await.expect("append");
assert!(logs.shutdown().await.is_empty(), "the log closed cleanly");
}
let conn = rusqlite::Connection::open(&path).expect("open the file");
conn.execute(
"INSERT INTO events (stream, seq, epoch_ms, kind, schema_version, meta, data) \
VALUES ('s-1', 2, 0, 'e2', 2, '{}', 'not json')",
[],
)
.expect("write a bad row");
drop(conn);
let logs = Logs::new();
let store = SqliteEventStore::open(&path, "s-1", &logs)
.await
.expect("reopen");
let err = store
.read(0, usize::MAX)
.await
.expect_err("a row that will not decode must surface");
assert_eq!(err.kind(), KnlError::CORRUPTION, "{err}");
}
#[test]
fn every_store_error_has_a_kernel_class() {
use eventsdb_core::Error as Failure;
let cases = [
(Failure::Validation("v".into()), KnlError::VALIDATION),
(Failure::Busy("b".into()), KnlError::BUSY),
(Failure::Timeout("t".into()), KnlError::TIMEOUT),
(Failure::Storage("s".into()), KnlError::STORAGE),
(Failure::Corruption("c".into()), KnlError::CORRUPTION),
(Failure::Unsupported("u".into()), KnlError::UNSUPPORTED),
(
Failure::Truncated {
requested: 1,
removed_up_to: 2,
},
KnlError::STORAGE,
),
];
for (failure, expected) in cases {
let translated = KnlError::from(failure);
assert_eq!(translated.kind(), expected, "{translated}");
}
assert!(
!KnlError::from(Failure::Timeout("t".into())).is_retryable(),
"a deadline is not contention"
);
assert!(KnlError::from(Failure::Busy("b".into())).is_retryable());
}
#[tokio::test]
async fn a_write_that_stays_contended_surfaces_as_busy() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
let log = SqliteEventLog::open_with(
&path,
eventsdb_sqlite::OpenOptions::default()
.busy_timeout(std::time::Duration::from_millis(50))
.upcasters(crate::knl::kernel_upcasters()),
)
.await
.expect("open");
let mut store = SqliteEventStore::on(Arc::new(log), "s-1");
store.append(ev(1)).await.expect("the first append");
let blocker = rusqlite::Connection::open(&path).expect("open the file");
blocker
.execute_batch("BEGIN IMMEDIATE; CREATE TABLE IF NOT EXISTS held (x)")
.expect("hold the lock");
let err = store
.append(ev(2))
.await
.expect_err("a write that stays contended must surface");
assert_eq!(err.kind(), KnlError::BUSY, "{err}");
assert!(err.is_retryable(), "busy is the class that says ask again");
blocker.execute_batch("ROLLBACK").expect("release");
store.append(ev(3)).await.expect("the lock is free again");
}
#[tokio::test]
async fn two_handles_on_one_stream_both_append_in_arrival_order() {
let logs = Logs::new();
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
let mut a = SqliteEventStore::open(&path, "s-1", &logs)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "s-1", &logs)
.await
.expect("open b");
assert_eq!(a.append(ev(1)).await.expect("a").seq, 1);
assert_eq!(b.append(ev(2)).await.expect("b").seq, 2);
assert_eq!(a.append(ev(3)).await.expect("a").seq, 3);
let stored = b.read(0, usize::MAX).await.expect("read");
assert_eq!(
stored.iter().map(kind_of).collect::<Vec<_>>(),
["e1", "e2", "e3"]
);
}
async fn ask(store: &SqliteEventStore, sql: &str) -> KnlResult<QueryRows> {
ask_with(store, sql, QueryParams::None, &QueryOpts::default()).await
}
async fn ask_with(
store: &SqliteEventStore,
sql: &str,
params: QueryParams,
opts: &QueryOpts,
) -> KnlResult<QueryRows> {
let plan = query::plan(sql, params, opts, &store.stream)?;
store.query(&plan).await
}
fn kinds_of(rows: &QueryRows) -> Vec<&str> {
rows.rows
.iter()
.filter_map(|row| row.get("kind").and_then(Value::as_str))
.collect()
}
#[tokio::test]
async fn a_query_reads_what_the_writer_wrote() {
let (mut store, _logs) = mem_store().await;
store.append(ev(1)).await.expect("append");
store.append(ev(2)).await.expect("append");
let rows = ask(
&store,
"SELECT kind, seq FROM events WHERE stream = $stream ORDER BY seq",
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1", "e2"]);
assert!(!rows.truncated);
assert_eq!(rows.rows[0].get("seq"), Some(&json!(1)));
}
#[tokio::test]
async fn stream_binds_to_this_stores_own_stream() {
let logs = Logs::new();
let log = logs.memory().await.expect("the log");
let mut mine = SqliteEventStore::on(Arc::clone(&log), "s-mine");
let mut theirs = SqliteEventStore::on(Arc::clone(&log), "s-theirs");
mine.append(ev(1)).await.expect("mine");
theirs.append(ev(2)).await.expect("theirs");
let rows = ask(&mine, "SELECT kind FROM events WHERE stream = $stream")
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1"]);
}
#[tokio::test]
async fn sessions_reads_across_the_set_it_was_given() {
let logs = Logs::new();
let log = logs.memory().await.expect("the log");
let mut one = SqliteEventStore::on(Arc::clone(&log), "s-one");
let mut two = SqliteEventStore::on(Arc::clone(&log), "s-two");
let mut three = SqliteEventStore::on(Arc::clone(&log), "s-three");
one.append(ev(1)).await.expect("one");
two.append(ev(2)).await.expect("two");
three.append(ev(3)).await.expect("three");
let opts = QueryOpts {
sessions: Some(vec!["s-one".to_string(), "s-two".to_string()]),
..QueryOpts::default()
};
let rows = ask_with(
&one,
"SELECT kind FROM events WHERE stream IN $sessions ORDER BY position",
QueryParams::None,
&opts,
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1", "e2"]);
}
#[tokio::test]
async fn a_bound_value_with_a_quote_in_it_is_a_value() {
let (mut store, _logs) = mem_store().await;
store
.append(obj(json!({ "kind": "it's fine" })))
.await
.expect("append");
let rows = ask_with(
&store,
"SELECT kind FROM events WHERE stream = $stream AND kind = :kind",
QueryParams::Named(obj(json!({ "kind": "it's fine" }))),
&QueryOpts::default(),
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["it's fine"]);
}
#[tokio::test]
async fn the_row_cap_is_reported_when_it_cuts() {
let (mut store, _logs) = mem_store().await;
for i in 1..=5 {
store.append(ev(i)).await.expect("append");
}
let capped = QueryOpts {
limit: 2,
..QueryOpts::default()
};
let rows = ask_with(
&store,
"SELECT kind FROM events WHERE stream = $stream ORDER BY seq",
QueryParams::None,
&capped,
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1", "e2"]);
assert!(rows.truncated, "the answer was cut");
let exact = QueryOpts {
limit: 5,
..QueryOpts::default()
};
let rows = ask_with(
&store,
"SELECT kind FROM events WHERE stream = $stream ORDER BY seq",
QueryParams::None,
&exact,
)
.await
.expect("query");
assert_eq!(rows.rows.len(), 5);
assert!(!rows.truncated, "nothing was cut off");
}
#[tokio::test]
async fn a_query_that_runs_too_long_is_a_timeout() {
let (store, _logs) = mem_store().await;
let hurried = QueryOpts {
timeout_ms: 50,
..QueryOpts::default()
};
let err = ask_with(
&store,
"WITH RECURSIVE forever(x) AS (SELECT 1 UNION ALL SELECT x + 1 FROM forever) \
SELECT COUNT(*) FROM forever",
QueryParams::None,
&hurried,
)
.await
.expect_err("an endless query must be cut short");
assert_eq!(err.kind(), KnlError::TIMEOUT, "{err}");
assert!(!err.is_retryable(), "a slow query is not a retry: {err}");
assert!(ask(&store, "SELECT 1 AS one").await.is_ok());
}
#[tokio::test]
async fn a_write_or_a_second_statement_is_refused_before_the_connection() {
let (store, _logs) = mem_store().await;
for sql in [
"INSERT INTO events (stream) VALUES ('x')",
"UPDATE events SET kind = 'x'",
"PRAGMA table_info(events)",
"ATTACH DATABASE '/tmp/other.db' AS other",
"SELECT 1; DROP TABLE events",
] {
let err = ask(&store, sql).await.expect_err("must be refused");
assert_eq!(err.kind(), KnlError::VALIDATION, "{sql:?}: {err}");
}
}
#[tokio::test]
async fn a_statement_that_does_not_compile_is_the_callers() {
let (store, _logs) = mem_store().await;
for (sql, phrase) in [
("SELECT beat FROM events", "no such column"),
("SELECT * FROM chronicle", "no such table"),
("SELECT nonesuch(1) AS x", "no such function"),
("SELECT 1 + FROM events", "syntax error"),
("SELECT * FROM events ORDER seq", "near \""),
("SELECT abs(1, 2) AS x", "wrong number of arguments"),
(
"SELECT seq FROM events AS a, events AS b",
"ambiguous column name",
),
] {
let err = ask(&store, sql).await.expect_err("must not compile");
assert_eq!(err.kind(), KnlError::VALIDATION, "{sql:?}: {err}");
assert!(
err.reason().contains(phrase),
"{sql:?}: expected SQLite to say {phrase:?}, got {err}"
);
assert!(
err.reason().starts_with("sql: "),
"{sql:?}: the reason names which half was wrong: {err}"
);
}
}
#[tokio::test]
async fn a_real_store_fault_is_still_storage() {
let (store, _logs) = mem_store().await;
let err = ask(&store, "SELECT zeroblob(1000000001) AS huge")
.await
.expect_err("over SQLITE_MAX_LENGTH");
assert_eq!(err.kind(), KnlError::STORAGE, "{err}");
}
#[tokio::test]
async fn the_sqlite_types_map_onto_values_and_null_is_absence() {
let (store, _logs) = mem_store().await;
let rows = ask(
&store,
"SELECT 1 AS whole, 1.5 AS fraction, 'text' AS words, NULL AS absent",
)
.await
.expect("query");
let row = &rows.rows[0];
assert_eq!(row["whole"], Value::from(1));
assert_eq!(row["fraction"], Value::from(1.5));
assert_eq!(row["words"], Value::from("text"));
assert!(
!row.contains_key("absent"),
"a NULL column is absent, so it reads as nil: {row:?}"
);
}
#[tokio::test]
async fn a_cell_with_no_json_value_is_refused_and_the_refusal_names_it() {
let (store, _logs) = mem_store().await;
for (sql, column, advice) in [
("SELECT CAST('bytes' AS BLOB) AS raw", "raw", "hex(raw)"),
("SELECT 9e999 AS boundless", "boundless", "CAST(boundless"),
("SELECT CAST(x'ff' AS TEXT) AS garbled", "garbled", "UTF-8"),
] {
let err = ask(&store, sql).await.expect_err("must be refused");
assert_eq!(err.kind(), KnlError::UNSUPPORTED, "{sql:?}: {err}");
assert!(
err.reason().contains(column),
"{sql:?}: the refusal names the column: {err}"
);
assert!(
err.reason().contains(advice),
"{sql:?}: expected {advice:?} in the refusal: {err}"
);
}
}
#[tokio::test]
async fn the_published_schema_is_the_events_table() {
let columns = events_schema().expect("schema");
let names: Vec<&str> = columns.iter().map(|c| c.name.as_str()).collect();
assert_eq!(
names,
[
"position",
"stream",
"seq",
"epoch_ms",
"kind",
"schema_version",
"meta",
"data"
]
);
let pk: Vec<&str> = columns
.iter()
.filter(|c| c.pk)
.map(|c| c.name.as_str())
.collect();
assert_eq!(pk, ["position"], "the log is keyed by its global order");
let (store, _logs) = mem_store().await;
let sql = format!("SELECT {} FROM {EVENTS_TABLE}", names.join(", "));
ask(&store, &sql)
.await
.expect("the published columns are the real ones");
}
#[tokio::test]
async fn the_published_schema_is_the_live_one() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("knl.db");
let logs = Logs::new();
let log = logs.file(&path).await.expect("open");
let rows = log
.query(&format!("PRAGMA table_info({EVENTS_TABLE})"), vec![])
.await
.expect("table_info");
let live: Vec<SchemaColumn> = rows
.iter()
.map(|row| SchemaColumn {
name: row["name"].as_str().expect("a column name").to_string(),
declared_type: row["type"].as_str().expect("a declared type").to_string(),
pk: row["pk"].as_i64().expect("a pk flag") > 0,
})
.collect();
assert_eq!(live, events_schema().expect("published schema"));
drop(log);
assert!(logs.shutdown().await.is_empty(), "the log closed cleanly");
}
}