use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, PoisonError};
use std::time::Duration;
use async_trait::async_trait;
use rusqlite::types::{Value as SqlValue, ValueRef};
use rusqlite::{params, params_from_iter, Connection, OpenFlags, TransactionBehavior};
use rusqlite_isle::{AsyncIsle, AsyncIsleDriver, IsleError, RetryPolicy};
use serde_json::{Map, Value};
use tokio::sync::OnceCell;
use super::event::{
stamp, validate_event, FIELD_BEAT, FIELD_DATA, FIELD_EPOCH_MS, FIELD_KIND, FIELD_META,
FIELD_SEQ,
};
use super::event_store::{
stamp_schema_version, ChildScan, ChildrenDecision, Committed, Decision, EventStore, Split,
SplitDecision, CURRENT_SCHEMA_VERSION, SCHEMA_VERSION_FIELD,
};
use super::query::{session_slot, QueryParams, QueryPlan, QueryRows, STREAM_PARAM};
use super::{now_ms, KnlError, KnlResult};
const BUSY_TIMEOUT: Duration = Duration::from_secs(5);
#[derive(Clone, Default)]
pub struct IsleDrivers {
parked: Arc<Mutex<Vec<AsyncIsleDriver>>>,
}
impl std::fmt::Debug for IsleDrivers {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("IsleDrivers")
.field("parked", &self.len())
.finish()
}
}
impl IsleDrivers {
pub fn new() -> Self {
Self::default()
}
fn park(&self, driver: AsyncIsleDriver) {
self.parked
.lock()
.unwrap_or_else(PoisonError::into_inner)
.push(driver);
}
pub fn len(&self) -> usize {
self.parked
.lock()
.unwrap_or_else(PoisonError::into_inner)
.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub async fn shutdown(&self) -> Vec<IsleError> {
let drivers: Vec<AsyncIsleDriver> =
std::mem::take(&mut *self.parked.lock().unwrap_or_else(PoisonError::into_inner));
let mut failures = Vec::new();
for driver in drivers {
if let Err(e) = driver.shutdown().await {
failures.push(e);
}
}
failures
}
}
pub const EVENTS_TABLE: &str = "events";
const SCHEMA_DDL: &str = "CREATE TABLE IF NOT EXISTS events ( \
stream TEXT NOT NULL, \
seq INTEGER NOT NULL, \
epoch_ms INTEGER NOT NULL, \
kind TEXT NOT NULL, \
schema_version INTEGER NOT NULL, \
beat TEXT NULL, \
meta TEXT NOT NULL, \
data TEXT NOT NULL, \
PRIMARY KEY (stream, seq) \
); \
CREATE INDEX IF NOT EXISTS events_stream_kind_seq \
ON events (stream, kind, seq); \
CREATE INDEX IF NOT EXISTS events_stream_beat_seq \
ON events (stream, beat, seq); \
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 MEMORY_URI_PREFIX: &str = "file:knl-";
pub fn is_memory_database(database: &str) -> bool {
database.starts_with(MEMORY_URI_PREFIX)
}
#[derive(Debug, Clone)]
enum Db {
File(PathBuf),
Memory(String),
}
impl Db {
fn id(&self) -> String {
self.target().to_string_lossy().into_owned()
}
fn memory_uri(stream: &str) -> String {
format!("{MEMORY_URI_PREFIX}{stream}?mode=memory&cache=shared")
}
fn target(&self) -> PathBuf {
match self {
Self::File(path) => path.clone(),
Self::Memory(uri) => PathBuf::from(uri),
}
}
fn write_flags() -> OpenFlags {
OpenFlags::default() | OpenFlags::SQLITE_OPEN_URI
}
fn read_only_flags() -> OpenFlags {
OpenFlags::SQLITE_OPEN_READ_ONLY
| OpenFlags::SQLITE_OPEN_NO_MUTEX
| OpenFlags::SQLITE_OPEN_URI
}
async fn spawn_writer(&self, drivers: &IsleDrivers) -> KnlResult<AsyncIsle> {
let (isle, driver) = AsyncIsle::builder()
.thread_name("knl-events")
.open_flags(Self::write_flags())
.wal(BUSY_TIMEOUT)
.spawn(self.target(), |conn| conn.execute_batch(SCHEMA_DDL))
.await
.map_err(KnlError::from)?;
drivers.park(driver);
Ok(isle)
}
async fn spawn_reader(&self, drivers: &IsleDrivers) -> KnlResult<AsyncIsle> {
let (isle, driver) = AsyncIsle::builder()
.thread_name("knl-events-read")
.open_flags(Self::read_only_flags())
.busy_timeout(BUSY_TIMEOUT)
.spawn(self.target(), |conn| {
conn.execute_batch("PRAGMA query_only = 1;")
})
.await
.map_err(KnlError::from)?;
drivers.park(driver);
Ok(isle)
}
}
fn retry_policy() -> RetryPolicy {
RetryPolicy::default()
}
enum JobError {
Sqlite(rusqlite::Error),
Terminal(KnlError),
}
fn finish<T>(outcome: Result<T, JobError>) -> Result<KnlResult<T>, rusqlite::Error> {
match outcome {
Ok(value) => Ok(Ok(value)),
Err(JobError::Sqlite(error)) => Err(error),
Err(JobError::Terminal(error)) => Ok(Err(error)),
}
}
fn is_retryable(error: &rusqlite::Error) -> bool {
matches!(
error,
rusqlite::Error::SqliteFailure(inner, _)
if matches!(
inner.code,
rusqlite::ErrorCode::DatabaseBusy | rusqlite::ErrorCode::DatabaseLocked
)
)
}
pub struct SqliteEventStore {
writer: AsyncIsle,
db: Db,
reader: OnceCell<AsyncIsle>,
drivers: IsleDrivers,
stream: String,
db_id: String,
}
impl SqliteEventStore {
pub async fn open(
path: &Path,
stream: impl Into<String>,
drivers: &IsleDrivers,
) -> KnlResult<Self> {
Self::init(Db::File(path.to_path_buf()), stream.into(), drivers).await
}
pub async fn open_memory(stream: impl Into<String>, drivers: &IsleDrivers) -> KnlResult<Self> {
let stream = stream.into();
Self::init(Db::Memory(Db::memory_uri(&stream)), stream, drivers).await
}
async fn init(db: Db, stream: String, drivers: &IsleDrivers) -> KnlResult<Self> {
let writer = db.spawn_writer(drivers).await?;
let db_id = db.id();
Ok(Self {
writer,
db,
reader: OnceCell::new(),
drivers: drivers.clone(),
stream,
db_id,
})
}
async fn reader(&self) -> KnlResult<&AsyncIsle> {
self.reader
.get_or_try_init(|| self.db.spawn_reader(&self.drivers))
.await
}
pub async fn schema(&self) -> KnlResult<Vec<SchemaColumn>> {
self.reader()
.await?
.call(schema_of)
.await
.map_err(KnlError::from)
}
}
fn schema_of(conn: &mut Connection) -> rusqlite::Result<Vec<SchemaColumn>> {
let mut stmt = conn.prepare(&format!("PRAGMA table_info({EVENTS_TABLE})"))?;
let rows = stmt.query_map([], |row| {
Ok(SchemaColumn {
name: row.get::<_, String>("name")?,
declared_type: row.get::<_, String>("type")?,
pk: row.get::<_, i64>("pk")? > 0,
})
})?;
rows.collect::<rusqlite::Result<Vec<_>>>()
}
pub fn events_schema() -> KnlResult<Vec<SchemaColumn>> {
let mut conn = Connection::open_in_memory().map_err(KnlError::from)?;
conn.execute_batch(SCHEMA_DDL).map_err(KnlError::from)?;
schema_of(&mut conn).map_err(KnlError::from)
}
fn owned_kinds(kinds: Option<&[&str]>) -> Option<Vec<String>> {
kinds.map(|kinds| kinds.iter().map(|kind| (*kind).to_string()).collect())
}
#[async_trait]
impl EventStore for SqliteEventStore {
async fn append(&mut self, mut event: Map<String, Value>) -> KnlResult<Committed> {
validate_event(&event)?;
stamp_schema_version(&mut event);
let stream = self.stream.clone();
self.writer
.call_retry(retry_policy(), move |conn| {
finish(append_in(conn, &stream, &event))
})
.await
.map_err(KnlError::from)?
}
async fn append_many(&mut self, events: Vec<Map<String, Value>>) -> KnlResult<Vec<Committed>> {
if events.is_empty() {
return Ok(Vec::new());
}
for event in &events {
validate_event(event)?;
}
let stream = self.stream.clone();
self.writer
.call_retry(retry_policy(), move |conn| {
finish(append_many_in(conn, &stream, &events))
})
.await
.map_err(KnlError::from)?
}
async fn append_if(
&mut self,
kinds: Option<&[&str]>,
decide: Decision,
) -> KnlResult<Option<Committed>> {
let stream = self.stream.clone();
let kinds = owned_kinds(kinds);
self.writer
.call(move |conn| finish(append_if_in(conn, &stream, kinds.as_deref(), decide)))
.await
.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);
self.writer
.call(move |conn| {
finish(append_if_many_in(
conn,
&stream,
&other,
kinds.as_deref(),
decide,
))
})
.await
.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();
self.writer
.call(move |conn| finish(append_with_open_children_in(conn, &stream, &scan, decide)))
.await
.map_err(KnlError::from)?
}
fn database(&self) -> Option<&str> {
Some(&self.db_id)
}
async fn read_kinds(
&self,
kinds: Option<&[&str]>,
from_seq: u64,
limit: usize,
) -> KnlResult<Vec<Value>> {
if kinds.is_some_and(<[&str]>::is_empty) {
return Ok(Vec::new());
}
let capped = i64::try_from(limit).unwrap_or(i64::MAX);
let kinds = owned_kinds(kinds);
let stream = self.stream.clone();
self.writer
.call(move |conn| finish(read_in(conn, &stream, kinds.as_deref(), from_seq, capped)))
.await
.map_err(KnlError::from)?
}
async fn read_last(&self, n: usize) -> KnlResult<Vec<Value>> {
let capped = i64::try_from(n).unwrap_or(i64::MAX);
let stream = self.stream.clone();
self.writer
.call(move |conn| finish(read_last_in(conn, &stream, capped)))
.await
.map_err(KnlError::from)?
}
async fn head(&self) -> KnlResult<Option<u64>> {
let stream = self.stream.clone();
self.writer
.call(move |conn| head_in(conn, &stream))
.await
.map_err(KnlError::from)
}
async fn len(&self) -> KnlResult<usize> {
let stream = self.stream.clone();
self.writer
.call(move |conn| {
conn.query_row(
"SELECT COUNT(*) FROM events WHERE stream = ?1",
params![stream],
|row| row.get::<_, i64>(0),
)
})
.await
.map(|n| n as usize)
.map_err(KnlError::from)
}
async fn query(&self, plan: &QueryPlan) -> KnlResult<QueryRows> {
let timeout = plan.timeout;
let plan = plan.clone();
self.reader()
.await?
.call_timeout(timeout, move |conn| Ok(run_query(conn, &plan)))
.await
.map_err(KnlError::from)?
}
fn detach_append(&self, mut 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;
}
stamp_schema_version(&mut event);
let stream = self.stream.clone();
self.writer
.spawn_call(move |conn| finish(append_in(conn, &stream, &event)))
.detach();
}
}
fn query_error(error: rusqlite::Error) -> KnlError {
if let rusqlite::Error::SqliteFailure(inner, _) = &error {
if inner.code == rusqlite::ErrorCode::OperationInterrupted {
return KnlError::Timeout(format!("query interrupted: {error}"));
}
}
match error {
rusqlite::Error::Utf8Error(_)
| rusqlite::Error::FromSqlConversionFailure(..)
| rusqlite::Error::IntegralValueOutOfRange(..) => {
KnlError::Corruption(format!("sqlite: query: {error}"))
}
other => KnlError::from(other),
}
}
fn run_query(conn: &Connection, plan: &QueryPlan) -> KnlResult<QueryRows> {
let mut stmt = conn.prepare(&plan.sql).map_err(|error| {
if is_retryable(&error) {
KnlError::from(error)
} else {
KnlError::Validation(format!("sql: {error}"))
}
})?;
if !stmt.readonly() {
return Err(KnlError::Validation(
"a query may not write; only SELECT / WITH statements are run".to_string(),
));
}
bind(&mut stmt, plan)?;
let columns: Vec<String> = stmt
.column_names()
.into_iter()
.map(str::to_string)
.collect();
let mut rows = stmt.raw_query();
let mut out = Vec::new();
while out.len() < plan.limit {
let Some(row) = rows.next().map_err(query_error)? else {
return Ok(QueryRows {
rows: out,
truncated: false,
});
};
let mut record = Map::new();
for (index, column) in columns.iter().enumerate() {
if let Some(value) = read_value(row.get_ref(index).map_err(query_error)?)? {
record.insert(column.clone(), value);
}
}
out.push(record);
}
let truncated = rows.next().map_err(query_error)?.is_some();
Ok(QueryRows {
rows: out,
truncated,
})
}
fn bind(stmt: &mut rusqlite::Statement<'_>, plan: &QueryPlan) -> KnlResult<()> {
const NO_VALUES: &[Value] = &[];
let slots: Vec<String> = (0..plan.sessions.len()).map(session_slot).collect();
let given: &[Value] = match &plan.params {
QueryParams::Positional(values) => values,
_ => NO_VALUES,
};
let mut taken = 0;
for index in 1..=stmt.parameter_count() {
let name = stmt.parameter_name(index).map(str::to_string);
let Some(name) = name else {
let value = given.get(taken).ok_or_else(|| {
KnlError::Validation(format!(
"the query has more `?` parameters than the {} value(s) given",
given.len()
))
})?;
taken += 1;
stmt.raw_bind_parameter(index, SqlParam(value.clone()))
.map_err(query_error)?;
continue;
};
if name == STREAM_PARAM {
stmt.raw_bind_parameter(index, plan.stream.clone())
.map_err(query_error)?;
continue;
}
if let Some(slot) = slots.iter().position(|slot| *slot == name) {
stmt.raw_bind_parameter(index, plan.sessions[slot].clone())
.map_err(query_error)?;
continue;
}
let QueryParams::Named(named) = &plan.params else {
return Err(KnlError::Validation(format!(
"the query names the parameter {name:?}, so params must be a table of names \
to values"
)));
};
let value = named
.get(&name[1..])
.or_else(|| named.get(&name))
.ok_or_else(|| {
KnlError::Validation(format!("no value was given for the parameter {name:?}"))
})?;
stmt.raw_bind_parameter(index, SqlParam(value.clone()))
.map_err(query_error)?;
}
if given.len() > taken {
return Err(KnlError::Validation(format!(
"{} value(s) were given for {taken} `?` parameter(s)",
given.len()
)));
}
Ok(())
}
struct SqlParam(Value);
impl rusqlite::ToSql for SqlParam {
fn to_sql(&self) -> rusqlite::Result<rusqlite::types::ToSqlOutput<'_>> {
use rusqlite::types::ToSqlOutput;
let value = match &self.0 {
Value::Null => SqlValue::Null,
Value::Bool(b) => SqlValue::Integer(i64::from(*b)),
Value::Number(n) => {
if let Some(i) = n.as_i64() {
SqlValue::Integer(i)
} else if let Some(f) = n.as_f64() {
SqlValue::Real(f)
} else {
return Err(rusqlite::Error::ToSqlConversionFailure(
format!("{n} is not a SQLite number").into(),
));
}
}
Value::String(s) => SqlValue::Text(s.clone()),
other => {
return Err(rusqlite::Error::ToSqlConversionFailure(
format!("a {} is not a SQLite value", type_name_of(other)).into(),
));
}
};
Ok(ToSqlOutput::Owned(value))
}
}
fn type_name_of(value: &Value) -> &'static str {
match value {
Value::Null => "null",
Value::Bool(_) => "boolean",
Value::Number(_) => "number",
Value::String(_) => "string",
Value::Array(_) => "list",
Value::Object(_) => "table",
}
}
fn read_value(value: ValueRef<'_>) -> KnlResult<Option<Value>> {
Ok(match value {
ValueRef::Null => None,
ValueRef::Integer(i) => Some(Value::from(i)),
ValueRef::Real(f) => {
let number = serde_json::Number::from_f64(f).ok_or_else(|| {
KnlError::Storage(format!(
"sqlite: a REAL column is {f}, which has no value on the other side of the \
bridge"
))
})?;
Some(Value::Number(number))
}
ValueRef::Text(bytes) => {
let text = std::str::from_utf8(bytes).map_err(|e| {
KnlError::Corruption(format!("sqlite: a TEXT column is not valid UTF-8: {e}"))
})?;
Some(Value::from(text))
}
ValueRef::Blob(bytes) => Some(Value::from(String::from_utf8_lossy(bytes).into_owned())),
})
}
fn append_in(
conn: &mut Connection,
stream: &str,
event: &Map<String, Value>,
) -> Result<Committed, JobError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(JobError::Sqlite)?;
let seq = next_seq(&tx, stream).map_err(JobError::Sqlite)?;
let epoch_ms = now_ms();
let mut row = event.clone();
stamp(&mut row, seq, epoch_ms);
insert_row(&tx, stream, seq, epoch_ms, &row)?;
tx.commit().map_err(JobError::Sqlite)?;
Ok(Committed { seq, epoch_ms })
}
fn append_many_in(
conn: &mut Connection,
stream: &str,
events: &[Map<String, Value>],
) -> Result<Vec<Committed>, JobError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(JobError::Sqlite)?;
let committed = insert_batch(&tx, stream, events)?;
tx.commit().map_err(JobError::Sqlite)?;
Ok(committed)
}
fn insert_batch(
tx: &Connection,
stream: &str,
events: &[Map<String, Value>],
) -> Result<Vec<Committed>, JobError> {
let mut seq = next_seq(tx, stream).map_err(JobError::Sqlite)?;
let mut committed = Vec::with_capacity(events.len());
for event in events {
validate_event(event).map_err(JobError::Terminal)?;
let epoch_ms = now_ms();
let mut row = event.clone();
stamp_schema_version(&mut row);
stamp(&mut row, seq, epoch_ms);
insert_row(tx, stream, seq, epoch_ms, &row)?;
committed.push(Committed { seq, epoch_ms });
seq = seq.saturating_add(1);
}
Ok(committed)
}
fn append_if_many_in(
conn: &mut Connection,
stream: &str,
other: &str,
kinds: Option<&[String]>,
decide: SplitDecision,
) -> Result<Option<Split<Committed>>, JobError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(JobError::Sqlite)?;
let seen = Split {
own: read_in(&tx, stream, kinds, 0, i64::MAX)?,
other: read_in(&tx, other, None, 0, 1)?,
};
let Some(split) = decide(seen) else {
return Ok(None);
};
let own = insert_batch(&tx, stream, &split.own)?;
let other = insert_batch(&tx, other, &split.other)?;
tx.commit().map_err(JobError::Sqlite)?;
Ok(Some(Split { own, other }))
}
fn append_with_open_children_in(
conn: &mut Connection,
stream: &str,
scan: &ChildScan,
decide: ChildrenDecision,
) -> Result<Committed, JobError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(JobError::Sqlite)?;
let children = open_children_in(&tx, stream, scan).map_err(JobError::Sqlite)?;
let committed = insert_batch(&tx, stream, &[decide(children)])?;
tx.commit().map_err(JobError::Sqlite)?;
committed
.into_iter()
.next()
.ok_or_else(|| JobError::Terminal(KnlError::Storage("the close wrote nothing".to_string())))
}
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: &Connection,
stream: &str,
scan: &ChildScan,
) -> rusqlite::Result<Vec<String>> {
let mut stmt = conn.prepare(&child_scan_sql(scan))?;
let rows = stmt.query_map(params![stream], |row| row.get::<_, String>(0))?;
rows.collect()
}
fn append_if_in(
conn: &mut Connection,
stream: &str,
kinds: Option<&[String]>,
decide: Decision,
) -> Result<Option<Committed>, JobError> {
let tx = conn
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(JobError::Sqlite)?;
let events = read_in(&tx, stream, kinds, 0, i64::MAX)?;
let Some(event) = decide(events) else {
return Ok(None);
};
validate_event(&event).map_err(JobError::Terminal)?;
let seq = next_seq(&tx, stream).map_err(JobError::Sqlite)?;
let epoch_ms = now_ms();
let mut row = event.clone();
stamp_schema_version(&mut row);
stamp(&mut row, seq, epoch_ms);
insert_row(&tx, stream, seq, epoch_ms, &row)?;
tx.commit().map_err(JobError::Sqlite)?;
Ok(Some(Committed { seq, epoch_ms }))
}
const READ_COLUMNS: &str = "seq, epoch_ms, kind, schema_version, beat, meta, data";
struct StoredRow {
seq: i64,
epoch_ms: i64,
kind: String,
schema_version: i64,
beat: Option<String>,
meta: String,
data: String,
}
fn read_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<StoredRow> {
Ok(StoredRow {
seq: row.get(0)?,
epoch_ms: row.get(1)?,
kind: row.get(2)?,
schema_version: row.get(3)?,
beat: row.get(4)?,
meta: row.get(5)?,
data: row.get(6)?,
})
}
fn event_of(row: StoredRow) -> KnlResult<Value> {
let meta = decode_object(&row.meta, FIELD_META)?;
let data = decode_object(&row.data, FIELD_DATA)?;
let mut event = Map::new();
event.insert(FIELD_KIND.to_string(), Value::from(row.kind));
if let Some(beat) = row.beat {
event.insert(FIELD_BEAT.to_string(), Value::from(beat));
}
event.insert(FIELD_META.to_string(), meta);
event.insert(FIELD_DATA.to_string(), data);
event.insert(FIELD_SEQ.to_string(), Value::from(row.seq as u64));
event.insert(FIELD_EPOCH_MS.to_string(), Value::from(row.epoch_ms as u64));
event.insert(
SCHEMA_VERSION_FIELD.to_string(),
Value::from(row.schema_version as u64),
);
Ok(Value::Object(event))
}
fn decode_object(text: &str, column: &str) -> KnlResult<Value> {
let value = serde_json::from_str::<Value>(text)
.map_err(|e| KnlError::Corruption(format!("sqlite: corrupt event {column}: {e}")))?;
if !value.is_object() {
return Err(KnlError::Corruption(format!(
"sqlite: corrupt event {column}: stored as {}, not a table",
super::event::json_type_name(&value)
)));
}
Ok(value)
}
fn read_query(
stream: &str,
kinds: Option<&[String]>,
from_seq: u64,
limit: i64,
) -> (String, Vec<SqlValue>) {
let mut sql = format!("SELECT {READ_COLUMNS} FROM events WHERE stream = ? AND seq >= ?");
let mut args = vec![
SqlValue::Text(stream.to_string()),
SqlValue::Integer(from_seq as i64),
];
if let Some(kinds) = kinds {
let placeholders = vec!["?"; kinds.len()].join(", ");
sql.push_str(&format!(" AND kind IN ({placeholders})"));
args.extend(kinds.iter().map(|kind| SqlValue::Text(kind.clone())));
}
sql.push_str(" ORDER BY seq ASC LIMIT ?");
args.push(SqlValue::Integer(limit));
(sql, args)
}
fn read_last_query(stream: &str, n: i64) -> (String, Vec<SqlValue>) {
(
format!("SELECT {READ_COLUMNS} FROM events WHERE stream = ? ORDER BY seq DESC LIMIT ?"),
vec![SqlValue::Text(stream.to_string()), SqlValue::Integer(n)],
)
}
fn read_last_in(conn: &Connection, stream: &str, n: i64) -> Result<Vec<Value>, JobError> {
let (sql, args) = read_last_query(stream, n);
let mut stmt = conn.prepare(&sql).map_err(JobError::Sqlite)?;
let rows = stmt
.query_map(params_from_iter(args.iter()), read_row)
.map_err(JobError::Sqlite)?;
let mut events = Vec::new();
for row in rows {
let row = row.map_err(JobError::Sqlite)?;
events.push(event_of(row).map_err(JobError::Terminal)?);
}
events.reverse();
Ok(events)
}
fn read_in(
conn: &Connection,
stream: &str,
kinds: Option<&[String]>,
from_seq: u64,
limit: i64,
) -> Result<Vec<Value>, JobError> {
if kinds.is_some_and(<[String]>::is_empty) {
return Ok(Vec::new());
}
let (sql, args) = read_query(stream, kinds, from_seq, limit);
let mut stmt = conn.prepare(&sql).map_err(JobError::Sqlite)?;
let rows = stmt
.query_map(params_from_iter(args.iter()), read_row)
.map_err(JobError::Sqlite)?;
let mut events = Vec::new();
for row in rows {
let row = row.map_err(JobError::Sqlite)?;
events.push(event_of(row).map_err(JobError::Terminal)?);
}
Ok(events)
}
fn next_seq(conn: &Connection, stream: &str) -> Result<u64, rusqlite::Error> {
conn.query_row(
"SELECT COALESCE(MAX(seq), 0) + 1 FROM events WHERE stream = ?1",
params![stream],
|row| row.get::<_, i64>(0),
)
.map(|n| n as u64)
}
fn head_in(conn: &Connection, stream: &str) -> Result<Option<u64>, rusqlite::Error> {
let max: Option<i64> = conn.query_row(
"SELECT MAX(seq) FROM events WHERE stream = ?1",
params![stream],
|row| row.get::<_, Option<i64>>(0),
)?;
Ok(max.map(|n| n as u64))
}
fn insert_row(
conn: &Connection,
stream: &str,
seq: u64,
epoch_ms: u64,
event: &Map<String, Value>,
) -> Result<(), JobError> {
let kind = event.get(FIELD_KIND).and_then(Value::as_str).unwrap_or("");
let schema_version = event
.get(SCHEMA_VERSION_FIELD)
.and_then(Value::as_u64)
.unwrap_or(CURRENT_SCHEMA_VERSION);
let beat = event.get(FIELD_BEAT).and_then(Value::as_str);
let meta = encode_object(event.get(FIELD_META), FIELD_META)?;
let data = encode_object(event.get(FIELD_DATA), FIELD_DATA)?;
conn.execute(
"INSERT INTO events (stream, seq, epoch_ms, kind, schema_version, beat, meta, data) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
stream,
seq as i64,
epoch_ms as i64,
kind,
schema_version as i64,
beat,
meta,
data
],
)
.map_err(JobError::Sqlite)?;
Ok(())
}
fn encode_object(value: Option<&Value>, column: &str) -> Result<String, JobError> {
let Some(value) = value else {
return Ok("{}".to_string());
};
serde_json::to_string(value).map_err(|e| {
JobError::Terminal(KnlError::Storage(format!(
"sqlite: encode event {column}: {e}"
)))
})
}
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<IsleError> for KnlError {
fn from(error: IsleError) -> Self {
match error {
IsleError::Sqlite(error) => KnlError::from(error),
IsleError::QueueFull => {
KnlError::Busy("sqlite: the connection thread is at capacity".to_string())
}
IsleError::Timeout => KnlError::Timeout("sqlite: the deadline elapsed".to_string()),
IsleError::Cancelled => KnlError::Timeout("sqlite: the job was cancelled".to_string()),
other => KnlError::Storage(format!("sqlite: {other}")),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::knl::event::{kind_of, seq_of};
use crate::knl::query::QueryOpts;
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, IsleDrivers) {
let drivers = IsleDrivers::new();
let store = SqliteEventStore::open_memory(uuid::Uuid::new_v4().to_string(), &drivers)
.await
.expect("open");
(store, drivers)
}
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, _drivers) = 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, _drivers) = 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 append_if_decides_inside_the_transaction_and_writes_only_a_some() {
let (mut store, _drivers) = 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, _drivers) = mem_store().await;
store
.append_if(None, decide(|_| Some(obj(json!({ "text": "no kind" })))))
.await
.expect_err("kind is required");
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, _drivers) = 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 failed batch wrote nothing"
);
assert_eq!(
store.append(ev(5)).await.expect("append").seq,
4,
"no seq burnt"
);
assert!(store
.append_many(Vec::new())
.await
.expect("empty")
.is_empty());
assert_eq!(store.len().await.expect("len"), 4);
}
#[tokio::test]
async fn append_if_many_writes_both_streams_or_neither() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut parent = SqliteEventStore::open(&path, "p", &drivers)
.await
.expect("open the parent");
let child = SqliteEventStore::open(&path, "c", &drivers)
.await
.expect("open the child");
parent.append(ev(1)).await.expect("seed");
let committed = parent
.append_if_many(
"c",
None,
Box::new(|events| {
assert_eq!(events.own.len(), 1, "the decision reads its own stream");
assert!(
events.other.is_empty(),
"and is shown that the other one is empty"
);
Some(Split {
own: vec![ev(2)],
other: vec![ev(3), ev(4)],
})
}),
)
.await
.expect("both sides")
.expect("the decision wrote");
assert_eq!(
committed.own.iter().map(|c| c.seq).collect::<Vec<_>>(),
[2],
"this stream numbers on from its own head"
);
assert_eq!(
committed.other.iter().map(|c| c.seq).collect::<Vec<_>>(),
[1, 2],
"and the other from its own, which was empty"
);
assert_eq!(child.len().await.expect("len"), 2, "the other side landed");
assert_eq!(
parent
.append_if_many("c", None, Box::new(|_| None))
.await
.expect("append_if_many"),
None
);
assert_eq!(parent.len().await.expect("len"), 2);
assert_eq!(child.len().await.expect("len"), 2);
parent
.append_if_many("c", None, Box::new(|_| Some(Split::own(vec![ev(5)]))))
.await
.expect("append_if_many")
.expect("the decision wrote");
assert_eq!(parent.len().await.expect("len"), 3);
assert_eq!(child.len().await.expect("len"), 2, "and nothing else");
parent
.append_if_many(
"c",
None,
Box::new(|_| {
Some(Split {
own: vec![ev(6)],
other: vec![obj(json!({ "text": "no kind" }))],
})
}),
)
.await
.expect_err("kind is required");
assert_eq!(parent.len().await.expect("len"), 3, "nothing was written");
assert_eq!(child.len().await.expect("len"), 2);
}
#[test]
fn the_child_scan_reads_by_the_parent_index() {
let conn = Connection::open_in_memory().expect("an in-memory database");
conn.execute_batch(SCHEMA_DDL).expect("the schema");
let insert = |stream: &str, seq: i64, kind: &str, data: &str| {
conn.execute(
"INSERT INTO events \
(stream, seq, epoch_ms, kind, schema_version, beat, meta, data) \
VALUES (?1, ?2, ?3, ?4, 1, NULL, '{}', ?5)",
params![stream, seq, seq * 10, kind, data],
)
.expect("insert");
};
insert("c1", 1, "session_opened", r#"{"parent":"p"}"#);
insert("c2", 1, "session_opened", r#"{"parent":"p"}"#);
insert("c2", 2, "session_closed", "{}");
insert("c3", 1, "session_opened", r#"{"parent":"elsewhere"}"#);
insert("p", 1, "session_opened", "{}");
let scan = ChildScan {
opened: "session_opened".to_string(),
closed: "session_closed".to_string(),
parent_field: "parent".to_string(),
};
assert_eq!(
open_children_in(&conn, "p", &scan).expect("scan"),
vec!["c1".to_string()],
"the ended child and the other parent's are not this stream's open children"
);
let sql = format!("EXPLAIN QUERY PLAN {}", child_scan_sql(&scan));
let mut stmt = conn.prepare(&sql).expect("prepare the plan");
let plan: Vec<String> = stmt
.query_map(params!["p"], |row| row.get::<_, String>(3))
.expect("the plan's rows")
.collect::<rusqlite::Result<Vec<_>>>()
.expect("the plan's rows");
assert!(
plan.iter()
.any(|step| step.contains("events_session_opened_parent")),
"the openings must be looked up by the index: {plan:?}"
);
assert!(
!plan
.iter()
.any(|step| step.starts_with("SCAN events AS opened")),
"and not found by walking the table: {plan:?}"
);
}
#[test]
fn a_word_written_into_the_scan_stays_one_word() {
assert_eq!(sql_literal("parent"), "'parent'");
assert_eq!(sql_literal("a'b"), "'a''b'");
let conn = Connection::open_in_memory().expect("an in-memory database");
conn.execute_batch(SCHEMA_DDL).expect("the schema");
conn.execute(
"INSERT INTO events \
(stream, seq, epoch_ms, kind, schema_version, beat, meta, data) \
VALUES ('c', 1, 10, 'it''s open', 1, NULL, '{}', ?1)",
params![r#"{"pa'rent":"p"}"#],
)
.expect("insert");
let scan = ChildScan {
opened: "it's open".to_string(),
closed: "it's over".to_string(),
parent_field: "pa'rent".to_string(),
};
assert_eq!(
open_children_in(&conn, "p", &scan).expect("the statement parses and runs"),
vec!["c".to_string()],
"the quoted words are still the words being matched"
);
}
#[tokio::test]
async fn database_is_the_same_for_two_streams_of_one_database() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let elsewhere = dir.path().join("other.db");
let drivers = IsleDrivers::new();
let a = SqliteEventStore::open(&path, "a", &drivers)
.await
.expect("open a");
let b = SqliteEventStore::open(&path, "b", &drivers)
.await
.expect("open b");
let far = SqliteEventStore::open(&elsewhere, "a", &drivers)
.await
.expect("open far");
assert_eq!(a.database(), b.database(), "two streams, one database");
assert_ne!(a.database(), far.database(), "two databases");
assert_eq!(
a.database(),
Some(path.to_string_lossy().as_ref()),
"the target it was opened by"
);
let (mem, _mem_drivers) = mem_store().await;
let uri = mem.database().expect("a database").to_string();
assert!(uri.contains("mode=memory"), "{uri}");
let beside = SqliteEventStore::open(std::path::Path::new(&uri), "beside", &drivers)
.await
.expect("open beside");
assert_eq!(
beside.database(),
Some(uri.as_str()),
"opening that target reaches the same database"
);
}
#[tokio::test]
async fn open_children_are_the_unended_streams_that_name_this_one() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
fn opened(parent: &str) -> Map<String, Value> {
obj(json!({
"kind": "session_opened",
"data": { "scope_id": "sc", "owner": "anon", "parent": parent }
}))
}
let ended = obj(json!({ "kind": "session_closed", "data": { "reason": "done" } }));
let mut parent = SqliteEventStore::open(&path, "p", &drivers)
.await
.expect("open p");
let mut running = SqliteEventStore::open(&path, "kid-a", &drivers)
.await
.expect("open kid-a");
running.append(opened("p")).await.expect("opened");
let mut over = SqliteEventStore::open(&path, "kid-b", &drivers)
.await
.expect("open kid-b");
over.append(opened("p")).await.expect("opened");
over.append(ended.clone()).await.expect("closed");
let mut theirs = SqliteEventStore::open(&path, "kid-c", &drivers)
.await
.expect("open kid-c");
theirs.append(opened("q")).await.expect("opened");
let mut root = SqliteEventStore::open(&path, "r", &drivers)
.await
.expect("open r");
root.append(obj(
json!({ "kind": "session_opened", "data": { "scope_id": "sc", "owner": "anon" } }),
))
.await
.expect("opened");
let scan = ChildScan {
opened: "session_opened".to_string(),
closed: "session_closed".to_string(),
parent_field: "parent".to_string(),
};
let seen: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen);
let committed = parent
.append_with_open_children(
&scan,
Box::new(move |children| {
*recorded.lock().expect("not poisoned") = children;
ended.clone()
}),
)
.await
.expect("the close lands");
assert_eq!(
*seen.lock().expect("not poisoned"),
["kid-a"],
"only the unended streams that named this one"
);
assert_eq!(committed.seq, 1, "and the event it built was appended");
assert_eq!(parent.len().await.expect("len"), 1);
}
#[tokio::test]
async fn read_kinds_selects_by_kind_and_keeps_the_streams_order() {
let (mut store, _drivers) = 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");
store.append(ev(2)).await.expect("more noise");
let ledger = store
.read_kinds(Some(&["budget_granted", "budget_spent"]), 0, usize::MAX)
.await
.expect("read_kinds");
let kinds: Vec<&str> = ledger.iter().map(kind_of).collect();
assert_eq!(kinds, ["budget_granted", "budget_spent"]);
assert_eq!(seq_of(&ledger[0]), 1);
assert_eq!(seq_of(&ledger[1]), 3, "the seq is the stream's");
assert_eq!(
store
.read_kinds(Some(&["budget_granted"]), 2, usize::MAX)
.await
.expect("read_kinds")
.len(),
0
);
assert_eq!(
store
.read_kinds(Some(&["budget_granted", "budget_spent"]), 0, 1)
.await
.expect("read_kinds")
.len(),
1
);
assert!(store
.read_kinds(Some(&[]), 0, usize::MAX)
.await
.expect("read_kinds")
.is_empty());
assert_eq!(
store
.read_kinds(None, 0, usize::MAX)
.await
.expect("read_kinds")
.len(),
4
);
}
#[tokio::test]
async fn append_if_filters_the_decisions_input_and_numbers_against_the_stream() {
let (mut store, _drivers) = mem_store().await;
store
.append(budget("budget_granted", 100))
.await
.expect("grant");
store.append(ev(1)).await.expect("noise");
store.append(ev(2)).await.expect("more noise");
let seen: Arc<Mutex<Vec<String>>> = Arc::default();
let recorded = Arc::clone(&seen);
let committed = store
.append_if(
Some(&["budget_granted"]),
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(4),
"the write lands after everything, not after the filtered read"
);
assert_eq!(store.len().await.expect("len"), 4);
}
#[tokio::test]
async fn append_if_across_two_handles_decides_on_the_other_handles_write() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open b");
let only_once = || {
decide(|events: Vec<Value>| {
(!events.iter().any(|e| kind_of(e) == "marker"))
.then(|| obj(json!({ "kind": "marker" })))
})
};
let first = a.append_if(None, only_once()).await.expect("a decides");
assert_eq!(first.map(|c| c.seq), Some(1), "a wrote the marker");
let second = b.append_if(None, only_once()).await.expect("b decides");
assert_eq!(second, None, "b saw a's marker and wrote nothing");
assert_eq!(b.len().await.expect("len"), 1, "exactly one marker");
}
#[tokio::test]
async fn read_pages_by_from_seq_and_limit() {
let (mut store, _drivers) = mem_store().await;
for i in 1..=5 {
store.append(ev(i)).await.expect("append");
}
assert_eq!(store.read(0, usize::MAX).await.expect("read").len(), 5);
assert_eq!(store.read(1, usize::MAX).await.expect("read").len(), 5);
assert_eq!(store.read(3, usize::MAX).await.expect("read").len(), 3);
assert_eq!(store.read(6, usize::MAX).await.expect("read").len(), 0);
let page = store.read(2, 2).await.expect("read");
assert_eq!(page.len(), 2);
assert_eq!(kind_of(&page[0]), "e2");
assert_eq!(kind_of(&page[1]), "e3");
assert!(store.read(0, 0).await.expect("read").is_empty());
}
#[tokio::test]
async fn read_last_takes_the_end_of_the_stream_in_seq_order() {
let (sql, args) = read_last_query("s-1", 5);
assert!(
sql.contains("ORDER BY seq DESC LIMIT ?"),
"the read must stop at n rows: {sql}"
);
assert_eq!(args.len(), 2, "the stream and the cap are bound: {sql}");
let (mut store, _drivers) = mem_store().await;
for i in 1..=200 {
store.append(ev(i)).await.expect("append");
}
let tail = store.read_last(5).await.expect("read_last");
assert_eq!(tail.len(), 5);
assert_eq!(kind_of(&tail[0]), "e196", "oldest of the five first");
assert_eq!(kind_of(&tail[4]), "e200", "the head last");
assert_eq!(store.read_last(usize::MAX).await.expect("all").len(), 200);
assert!(store.read_last(0).await.expect("none").is_empty());
}
#[tokio::test]
async fn head_is_none_when_empty_then_tracks_the_max() {
let (mut store, _drivers) = 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));
store
.append(obj(json!({ "text": "no kind" })))
.await
.expect_err("kind is required");
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, _drivers) = mem_store().await;
store
.append(obj(json!({
"kind": "note",
"beat": "b1",
"meta": { "label": "a", "attempt": 2, "retried": true },
"data": { "text": "hi", "nested": { "deep": [1, 2] } }
})))
.await
.expect("append");
store
.append(obj(json!({ "kind": "note" })))
.await
.expect("append a bare one");
let stored = store.read(0, usize::MAX).await.expect("read");
assert_eq!(kind_of(&stored[0]), "note");
assert_eq!(stored[0]["beat"], json!("b1"));
assert_eq!(
stored[0]["meta"],
json!({ "label": "a", "attempt": 2, "retried": true })
);
assert_eq!(
stored[0]["data"],
json!({ "text": "hi", "nested": { "deep": [1, 2] } }),
"data comes back at any depth"
);
assert_eq!(seq_of(&stored[0]), 1);
assert!(stored[0].get("epoch_ms").is_some(), "{}", stored[0]);
assert_eq!(
stored[0].get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION)
);
assert_eq!(stored[1].get("beat"), None, "{}", stored[1]);
assert_eq!(stored[1]["meta"], json!({}));
assert_eq!(stored[1]["data"], json!({}));
}
#[tokio::test]
async fn the_beat_is_a_column_of_its_own_with_an_index() {
let (mut store, _drivers) = mem_store().await;
store
.append(obj(json!({ "kind": "e1", "beat": "b1" })))
.await
.expect("append");
store.append(ev(2)).await.expect("append with no beat");
let rows = ask(
&store,
"SELECT seq, beat FROM events WHERE stream = $stream ORDER BY seq",
)
.await
.expect("query");
assert_eq!(rows.rows[0]["beat"], Value::from("b1"));
assert!(
!rows.rows[1].contains_key("beat"),
"an undeclared beat is NULL: {:?}",
rows.rows[1]
);
let indexes = store
.writer
.call(|conn| {
let mut stmt = conn.prepare("PRAGMA index_list(events)")?;
let names = stmt
.query_map([], |row| row.get::<_, String>("name"))?
.collect::<rusqlite::Result<Vec<_>>>();
names
})
.await
.expect("index_list");
assert!(
indexes.iter().any(|name| name == "events_stream_beat_seq"),
"the (stream, beat, seq) index must exist: {indexes:?}"
);
assert!(
indexes.iter().any(|name| name == "events_stream_kind_seq"),
"…beside the by-kind one: {indexes:?}"
);
}
#[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("events.db");
{
let drivers = IsleDrivers::new();
let mut store = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open");
store
.append(obj(
json!({ "kind": "note", "data": { "text": "durable" } }),
))
.await
.expect("append note");
store.append(ev(2)).await.expect("append e2");
drop(store);
assert!(drivers.shutdown().await.is_empty(), "the writer joined");
}
let drivers = IsleDrivers::new();
let store = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("reopen");
let events = store.read(0, usize::MAX).await.expect("read");
assert_eq!(events.len(), 2);
assert_eq!(kind_of(&events[0]), "note");
assert_eq!(events[0]["data"], json!({ "text": "durable" }));
assert_eq!(seq_of(&events[0]), 1);
assert_eq!(
events[0].get(SCHEMA_VERSION_FIELD).and_then(Value::as_u64),
Some(CURRENT_SCHEMA_VERSION),
"the schema version survives the round-trip too"
);
assert_eq!(store.head().await.expect("head"), Some(2));
}
#[tokio::test]
async fn two_streams_in_one_db_file_do_not_see_each_others_events() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "stream-a", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "stream-b", &drivers)
.await
.expect("open b");
a.append(obj(json!({ "kind": "only_a" })))
.await
.expect("append a");
b.append(obj(json!({ "kind": "only_b1" })))
.await
.expect("append b1");
b.append(obj(json!({ "kind": "only_b2" })))
.await
.expect("append b2");
assert_eq!(a.len().await.expect("len"), 1);
assert_eq!(b.len().await.expect("len"), 2);
assert_eq!(a.head().await.expect("head"), Some(1));
assert_eq!(b.head().await.expect("head"), Some(2));
assert_eq!(
kind_of(&a.read(0, usize::MAX).await.expect("read")[0]),
"only_a"
);
let b_events = b.read(0, usize::MAX).await.expect("read");
let b_kinds: Vec<&str> = b_events.iter().map(kind_of).collect();
assert_eq!(b_kinds, ["only_b1", "only_b2"]);
}
#[tokio::test]
async fn read_errors_on_a_corrupt_row_instead_of_dropping_it() {
for (seq, meta, data, column) in [
(2_i64, "{}", "{not valid json", "data"),
(3_i64, "not valid json either", "{}", "meta"),
(4_i64, "{}", "7", "data"),
] {
let (mut store, _drivers) = mem_store().await;
store.append(ev(1)).await.expect("append");
let stream = store.stream.clone();
let (meta, data) = (meta.to_string(), data.to_string());
store
.writer
.call(move |conn| {
conn.execute(
"INSERT INTO events \
(stream, seq, epoch_ms, kind, schema_version, beat, meta, data) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
stream,
seq,
0_i64,
"note",
1_i64,
None::<String>,
meta,
data
],
)
})
.await
.expect("insert corrupt row");
let err = store
.read(0, usize::MAX)
.await
.expect_err("a corrupt row must surface, not be dropped");
assert!(
err.reason().contains(&format!("corrupt event {column}")),
"{}",
err.reason()
);
assert_eq!(err.kind(), KnlError::CORRUPTION);
assert!(!err.is_retryable());
}
}
#[test]
fn a_contended_lock_is_busy_and_every_other_fault_is_storage() {
fn failure(code: rusqlite::ErrorCode) -> rusqlite::Error {
rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error {
code,
extended_code: 0,
},
Some("under test".to_string()),
)
}
for code in [
rusqlite::ErrorCode::DatabaseBusy,
rusqlite::ErrorCode::DatabaseLocked,
] {
let error = KnlError::from(failure(code));
assert_eq!(error.kind(), KnlError::BUSY, "{code:?}: {error}");
assert!(error.is_retryable(), "{code:?}: {error}");
}
for code in [
rusqlite::ErrorCode::DatabaseCorrupt,
rusqlite::ErrorCode::ReadOnly,
rusqlite::ErrorCode::DiskFull,
] {
let error = KnlError::from(failure(code));
assert_eq!(error.kind(), KnlError::STORAGE, "{code:?}: {error}");
assert!(
!error.is_retryable(),
"the kernel does not promise a retry it cannot back: {error}"
);
}
let error = KnlError::from(rusqlite::Error::QueryReturnedNoRows);
assert_eq!(error.kind(), KnlError::STORAGE, "{error}");
}
#[tokio::test]
async fn a_write_that_stays_contended_surfaces_as_busy() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut store = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open");
store.append(ev(1)).await.expect("seed");
let blocker = Connection::open(&path).expect("open blocker");
blocker
.execute_batch("BEGIN EXCLUSIVE")
.expect("take the write lock");
store
.writer
.call(|conn| conn.busy_timeout(Duration::from_millis(0)))
.await
.expect("no waiting");
let err = store
.append(ev(2))
.await
.expect_err("a write against a held lock must not succeed");
assert_eq!(err.kind(), KnlError::BUSY, "{err}");
assert!(err.is_retryable(), "{err}");
}
#[tokio::test]
async fn two_handles_on_one_stream_both_append_in_arrival_order() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open b");
a.append(ev(1)).await.expect("seed");
assert_eq!(a.append(ev(2)).await.expect("a appends").seq, 2);
assert_eq!(b.append(ev(3)).await.expect("b appends").seq, 3);
let events = b.read(0, usize::MAX).await.expect("read");
let kinds: Vec<&str> = events.iter().map(kind_of).collect();
assert_eq!(kinds, ["e1", "e2", "e3"]);
assert_eq!(b.head().await.expect("head"), Some(3));
}
#[tokio::test]
async fn immediate_tx_appends_round_trip_across_two_handles() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "s", &drivers)
.await
.expect("open b");
assert_eq!(a.append(ev(1)).await.expect("a1").seq, 1);
assert_eq!(b.append(ev(2)).await.expect("b2").seq, 2);
assert_eq!(a.append(ev(3)).await.expect("a3").seq, 3);
assert_eq!(b.head().await.expect("head"), Some(3));
assert_eq!(b.read(0, usize::MAX).await.expect("read").len(), 3);
}
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 = crate::knl::query::plan(sql, params, opts, &store.stream)?;
store.query(&plan).await
}
fn kinds_of(rows: &QueryRows) -> Vec<&str> {
rows.rows
.iter()
.map(|row| row["kind"].as_str().expect("kind is a string"))
.collect()
}
#[tokio::test]
async fn the_reader_sees_the_writers_rows_in_memory() {
let (mut store, _drivers) = mem_store().await;
store.append(ev(1)).await.expect("append e1");
store.append(ev(2)).await.expect("append e2");
let rows = ask(
&store,
"SELECT seq, kind FROM events WHERE stream = $stream ORDER BY seq",
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1", "e2"]);
assert_eq!(rows.rows[0]["seq"], Value::from(1));
assert!(!rows.truncated);
store.append(ev(3)).await.expect("append e3");
let again = ask(&store, "SELECT kind FROM events ORDER BY seq")
.await
.expect("query");
assert_eq!(kinds_of(&again), ["e1", "e2", "e3"]);
}
#[tokio::test]
async fn stream_binds_to_this_stores_own_stream() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "stream-a", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "stream-b", &drivers)
.await
.expect("open b");
a.append(obj(json!({ "kind": "only_a" }))).await.expect("a");
b.append(obj(json!({ "kind": "only_b" }))).await.expect("b");
let rows = ask(&a, "SELECT kind FROM events WHERE stream = $stream")
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["only_a"]);
let rows = ask(&b, "SELECT kind FROM events WHERE stream = $stream")
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["only_b"]);
}
#[tokio::test]
async fn sessions_reads_across_the_set_it_was_given() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("events.db");
let drivers = IsleDrivers::new();
let mut a = SqliteEventStore::open(&path, "stream-a", &drivers)
.await
.expect("open a");
let mut b = SqliteEventStore::open(&path, "stream-b", &drivers)
.await
.expect("open b");
a.append(obj(json!({ "kind": "from_a" }))).await.expect("a");
b.append(obj(json!({ "kind": "from_b1" })))
.await
.expect("b1");
b.append(obj(json!({ "kind": "from_b2" })))
.await
.expect("b2");
let opts = QueryOpts {
sessions: Some(vec!["stream-a".to_string(), "stream-b".to_string()]),
..QueryOpts::default()
};
let rows = ask_with(
&a,
"SELECT stream, kind FROM events WHERE stream IN $sessions ORDER BY stream, seq",
QueryParams::None,
&opts,
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["from_a", "from_b1", "from_b2"]);
let rows = ask(&a, "SELECT kind FROM events WHERE stream IN $sessions")
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["from_a"]);
}
#[tokio::test]
async fn a_bound_value_with_a_quote_in_it_is_a_value() {
let (mut store, _drivers) = mem_store().await;
store
.append(obj(json!({ "kind": "it's a kind" })))
.await
.expect("append");
store.append(ev(1)).await.expect("append e1");
let rows = ask_with(
&store,
"SELECT kind FROM events WHERE kind = ?",
QueryParams::Positional(vec![json!("it's a kind")]),
&QueryOpts::default(),
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["it's a kind"]);
let named = QueryParams::Named(
json!({ "kind": "x' OR 1=1 --" })
.as_object()
.expect("an object")
.clone(),
);
let rows = ask_with(
&store,
"SELECT kind FROM events WHERE kind = :kind",
named,
&QueryOpts::default(),
)
.await
.expect("query");
assert!(rows.rows.is_empty(), "{:?}", rows.rows);
}
#[tokio::test]
async fn the_row_cap_is_reported_when_it_cuts() {
let (mut store, _drivers) = 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 ORDER BY seq",
QueryParams::None,
&capped,
)
.await
.expect("query");
assert_eq!(kinds_of(&rows), ["e1", "e2"]);
assert!(rows.truncated, "the cap cut three rows off");
let exact = QueryOpts {
limit: 5,
..QueryOpts::default()
};
let rows = ask_with(
&store,
"SELECT kind FROM events 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, _drivers) = 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 the_reader_connection_refuses_a_write() {
let (mut store, _drivers) = mem_store().await;
store.append(ev(1)).await.expect("append");
let reader = store.reader().await.expect("open the reader");
let err = reader
.call(|conn| {
conn.execute(
"INSERT INTO events \
(stream, seq, epoch_ms, kind, schema_version, beat, meta, data) \
VALUES ('x', 1, 0, 'note', 1, NULL, '{}', '{}')",
[],
)
})
.await
.expect_err("the reader must not be able to write");
assert!(
matches!(KnlError::from(err), KnlError::Storage(_)),
"a write through the reader is refused by SQLite itself"
);
assert_eq!(store.len().await.expect("len"), 1);
}
#[tokio::test]
async fn a_write_or_a_second_statement_is_refused_before_the_connection() {
let (store, _drivers) = 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 every_parameter_is_answered_and_every_value_is_used() {
let (store, _drivers) = mem_store().await;
let err = ask(&store, "SELECT * FROM events WHERE kind = :kind")
.await
.expect_err("an unanswered parameter must be refused");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
assert!(err.reason().contains(":kind"), "{}", err.reason());
let err = ask_with(
&store,
"SELECT * FROM events WHERE kind = ?",
QueryParams::Positional(vec![json!("a"), json!("b")]),
&QueryOpts::default(),
)
.await
.expect_err("a value with no parameter must be refused");
assert_eq!(err.kind(), KnlError::VALIDATION, "{err}");
}
#[tokio::test]
async fn the_sqlite_types_map_onto_values_and_null_is_absence() {
let (store, _drivers) = mem_store().await;
let rows = ask(
&store,
"SELECT 1 AS whole, 1.5 AS fraction, 'text' AS words, NULL AS absent, \
CAST('bytes' AS BLOB) AS raw",
)
.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_eq!(row["raw"], Value::from("bytes"));
assert!(
!row.contains_key("absent"),
"a NULL column is absent, so it reads as nil: {row:?}"
);
}
#[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,
[
"stream",
"seq",
"epoch_ms",
"kind",
"schema_version",
"beat",
"meta",
"data"
]
);
let pk: Vec<&str> = columns
.iter()
.filter(|c| c.pk)
.map(|c| c.name.as_str())
.collect();
assert_eq!(pk, ["stream", "seq"], "the log is keyed by (stream, seq)");
let declared: Vec<&str> = columns.iter().map(|c| c.declared_type.as_str()).collect();
assert_eq!(
declared,
["TEXT", "INTEGER", "INTEGER", "TEXT", "INTEGER", "TEXT", "TEXT", "TEXT"]
);
let (store, _drivers) = 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 a_live_store_reports_the_published_schema() {
let (store, _drivers) = mem_store().await;
assert_eq!(
store.schema().await.expect("schema"),
events_schema().expect("published schema")
);
}
}