use eventsdb::sqlite::Projection;
use eventsdb::sqlite::rusqlite::{self, OptionalExtension, Transaction};
use eventsdb::{Error, Result};
use serde_json::Value as Json;
pub const STREAM_PREFIX: &str = "card-";
pub const ALIAS_PREFIX: &str = "alias-";
pub const PRUNE_STREAM: &str = "prune";
pub struct CardsProjection {
name: String,
}
impl Default for CardsProjection {
fn default() -> Self {
CardsProjection::new()
}
}
impl CardsProjection {
pub const NAME: &'static str = "cards_v4";
pub fn new() -> CardsProjection {
CardsProjection {
name: CardsProjection::NAME.to_string(),
}
}
#[cfg(test)]
pub fn under(name: &str) -> CardsProjection {
CardsProjection {
name: name.to_string(),
}
}
pub const KINDS: [&'static str; 10] = [
"card_opened",
"samples_appended",
"eval_recorded",
"checkpoint_saved",
"card_closed",
"tag_set",
"tag_unset",
"alias_bound",
"alias_released",
"cards_pruned",
];
}
const CREATE: &str = "\
CREATE TABLE IF NOT EXISTS cb_cards (
id TEXT PRIMARY KEY,
pkg TEXT,
scenario TEXT,
source TEXT,
created_by TEXT,
note TEXT,
model TEXT,
trace_id TEXT,
work_url TEXT,
fingerprint TEXT,
params_json TEXT,
state TEXT NOT NULL,
opened_ms INTEGER,
closed_ms INTEGER,
opened_position INTEGER,
error TEXT,
stats_json TEXT,
cost_json TEXT,
mean_score REAL,
n INTEGER,
pass_rate REAL,
passed INTEGER,
elapsed_ms INTEGER,
llm_calls INTEGER,
sample_batches INTEGER NOT NULL DEFAULT 0,
sample_rows INTEGER NOT NULL DEFAULT 0,
eval_count INTEGER NOT NULL DEFAULT 0,
checkpoint_count INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS cb_cards_pkg ON cb_cards (pkg);
CREATE INDEX IF NOT EXISTS cb_cards_state ON cb_cards (state);
CREATE INDEX IF NOT EXISTS cb_cards_opened_ms ON cb_cards (opened_ms);
CREATE INDEX IF NOT EXISTS cb_cards_model ON cb_cards (model);
CREATE INDEX IF NOT EXISTS cb_cards_trace ON cb_cards (trace_id);
CREATE INDEX IF NOT EXISTS cb_cards_print ON cb_cards (fingerprint);
CREATE TABLE IF NOT EXISTS cb_samples (
card_id TEXT NOT NULL,
seq INTEGER NOT NULL,
n INTEGER,
rows_json TEXT,
blob TEXT,
size INTEGER,
epoch_ms INTEGER,
PRIMARY KEY (card_id, seq)
);
CREATE TABLE IF NOT EXISTS cb_evals (
card_id TEXT NOT NULL,
seq INTEGER NOT NULL,
source TEXT,
data_json TEXT,
epoch_ms INTEGER,
PRIMARY KEY (card_id, seq)
);
CREATE TABLE IF NOT EXISTS cb_tags (
card_id TEXT NOT NULL,
key TEXT NOT NULL,
value TEXT NOT NULL,
set_ms INTEGER,
PRIMARY KEY (card_id, key)
);
CREATE INDEX IF NOT EXISTS cb_tags_key ON cb_tags (key, value);
CREATE TABLE IF NOT EXISTS cb_checkpoints (
card_id TEXT NOT NULL,
seq INTEGER NOT NULL,
blob TEXT,
size INTEGER,
format TEXT,
note TEXT,
epoch_ms INTEGER,
PRIMARY KEY (card_id, seq)
);
CREATE TABLE IF NOT EXISTS cb_lineage (
child TEXT NOT NULL,
parent TEXT NOT NULL,
PRIMARY KEY (child, parent)
);
CREATE INDEX IF NOT EXISTS cb_lineage_parent ON cb_lineage (parent);
CREATE TABLE IF NOT EXISTS cb_blobs (
hash TEXT PRIMARY KEY,
size INTEGER,
refs INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS cb_aliases (
name TEXT PRIMARY KEY,
card_id TEXT NOT NULL,
pkg TEXT,
bound_ms INTEGER,
note TEXT
);
CREATE INDEX IF NOT EXISTS cb_aliases_card ON cb_aliases (card_id);
CREATE TABLE IF NOT EXISTS cb_alias_log (
name TEXT NOT NULL,
seq INTEGER NOT NULL,
kind TEXT NOT NULL,
card_id TEXT,
epoch_ms INTEGER,
note TEXT,
PRIMARY KEY (name, seq)
);
";
const DROP: &str = "\
DROP TABLE IF EXISTS cb_cards;
DROP TABLE IF EXISTS cb_samples;
DROP TABLE IF EXISTS cb_evals;
DROP TABLE IF EXISTS cb_tags;
DROP TABLE IF EXISTS cb_checkpoints;
DROP TABLE IF EXISTS cb_lineage;
DROP TABLE IF EXISTS cb_blobs;
DROP TABLE IF EXISTS cb_aliases;
DROP TABLE IF EXISTS cb_alias_log;
";
impl Projection for CardsProjection {
fn name(&self) -> &str {
&self.name
}
fn kinds(&self) -> Option<Vec<String>> {
Some(
CardsProjection::KINDS
.iter()
.map(|k| k.to_string())
.collect(),
)
}
fn init(&mut self, tx: &Transaction<'_>) -> Result<()> {
if shape_outdated(tx)? {
tx.execute_batch(DROP).map_err(storage)?;
}
tx.execute_batch(CREATE).map_err(storage)
}
fn reset(&mut self, tx: &Transaction<'_>) -> Result<()> {
tx.execute_batch(DROP).map_err(storage)
}
fn tolerates_truncation(&self) -> bool {
true
}
fn apply(&mut self, tx: &Transaction<'_>, event: &eventsdb::Recorded) -> Result<()> {
let kind = event.kind();
let seq = event.seq() as i64;
let position = event.position.get() as i64;
let epoch_ms = num(event.event.get("epoch_ms")).unwrap_or(0);
let meta = event.event.get("meta");
let data = event.event.get("data");
match kind {
"card_opened" => opened(
tx,
card_id(&event.stream, kind)?,
epoch_ms,
position,
meta,
data,
),
"samples_appended" => samples(tx, card_id(&event.stream, kind)?, seq, epoch_ms, data),
"eval_recorded" => eval(tx, card_id(&event.stream, kind)?, seq, epoch_ms, meta, data),
"tag_set" => tag_set(tx, card_id(&event.stream, kind)?, epoch_ms, meta),
"tag_unset" => tag_unset(tx, card_id(&event.stream, kind)?, meta),
"checkpoint_saved" => {
checkpoint(tx, card_id(&event.stream, kind)?, seq, epoch_ms, data)
}
"card_closed" => closed(tx, card_id(&event.stream, kind)?, epoch_ms, meta, data),
"alias_bound" => bound(
tx,
alias_name(&event.stream, kind)?,
seq,
epoch_ms,
meta,
data,
),
"alias_released" => released(
tx,
alias_name(&event.stream, kind)?,
seq,
epoch_ms,
meta,
data,
),
"cards_pruned" => {
on_prune_stream(&event.stream, kind)?;
pruned(tx, data)
}
other => Err(Error::storage(format!(
"the {} projection was handed a {other:?} event, which is not one of the \
kinds it asked for ({})",
self.name,
CardsProjection::KINDS.join(", ")
))),
}
}
}
pub const SHAPE_MARKS: [(&str, &str); 2] = [("cb_cards", "params_json"), ("cb_evals", "source")];
pub fn shape_probe(table: &str) -> String {
format!("PRAGMA table_info({table})")
}
pub fn shape_outdated(tx: &Transaction<'_>) -> Result<bool> {
for (table, column) in SHAPE_MARKS {
let mut stmt = tx.prepare(&shape_probe(table)).map_err(storage)?;
let names = stmt
.query_map([], |row| row.get::<_, String>(1))
.map_err(storage)?
.collect::<rusqlite::Result<Vec<String>>>()
.map_err(storage)?;
if !names.is_empty() && !names.iter().any(|n| n == column) {
return Ok(true);
}
}
Ok(false)
}
fn card_id<'a>(stream: &'a str, kind: &str) -> Result<&'a str> {
stream.strip_prefix(STREAM_PREFIX).ok_or_else(|| {
Error::storage(format!(
"a {kind:?} event is on stream {stream:?}, which is not a card's: \
a card's stream is {STREAM_PREFIX}<id>"
))
})
}
fn alias_name<'a>(stream: &'a str, kind: &str) -> Result<&'a str> {
stream.strip_prefix(ALIAS_PREFIX).ok_or_else(|| {
Error::storage(format!(
"a {kind:?} event is on stream {stream:?}, which is not an alias's: \
an alias's stream is {ALIAS_PREFIX}<name>"
))
})
}
fn opened(
tx: &Transaction<'_>,
id: &str,
epoch_ms: i64,
position: i64,
meta: Option<&Json>,
data: Option<&Json>,
) -> Result<()> {
let params = data.and_then(|d| d.get("params"));
tx.execute(
"INSERT INTO cb_cards (id, pkg, scenario, source, created_by, note,
model, trace_id, work_url, fingerprint, params_json,
state, opened_ms, opened_position)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, 'open', ?12, ?13)
ON CONFLICT (id) DO NOTHING",
rusqlite::params![
id,
text(meta, "pkg"),
text(meta, "scenario"),
text(meta, "source"),
text(meta, "created_by"),
text(data, "note"),
text(meta, "model"),
text(meta, "trace_id"),
text(meta, "work_url"),
text(meta, "fingerprint"),
params.map(Json::to_string),
epoch_ms,
position,
],
)
.map_err(storage)?;
let parents = data.and_then(|d| d.get("parents")).and_then(Json::as_array);
for parent in parents.into_iter().flatten() {
let Some(parent) = parent.as_str() else {
continue;
};
tx.execute(
"INSERT INTO cb_lineage (child, parent) VALUES (?1, ?2)
ON CONFLICT (child, parent) DO NOTHING",
rusqlite::params![id, parent],
)
.map_err(storage)?;
}
Ok(())
}
fn samples(
tx: &Transaction<'_>,
id: &str,
seq: i64,
epoch_ms: i64,
data: Option<&Json>,
) -> Result<()> {
let n = num(data.and_then(|d| d.get("n")));
let rows = data
.and_then(|d| d.get("rows"))
.map(|rows| rows.to_string());
let blob = text(data, "blob");
let size = num(data.and_then(|d| d.get("size")));
tx.execute(
"INSERT INTO cb_samples (card_id, seq, n, rows_json, blob, size, epoch_ms)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT (card_id, seq) DO NOTHING",
rusqlite::params![id, seq, n, rows, blob.as_deref(), size, epoch_ms],
)
.map_err(storage)?;
tx.execute(
"UPDATE cb_cards SET sample_batches = sample_batches + 1,
sample_rows = sample_rows + ?2
WHERE id = ?1",
rusqlite::params![id, n.unwrap_or(0)],
)
.map_err(storage)?;
if let Some(hash) = blob {
reference_blob(tx, &hash, size)?;
}
Ok(())
}
fn eval(
tx: &Transaction<'_>,
id: &str,
seq: i64,
epoch_ms: i64,
meta: Option<&Json>,
data: Option<&Json>,
) -> Result<()> {
tx.execute(
"INSERT INTO cb_evals (card_id, seq, source, data_json, epoch_ms)
VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT (card_id, seq) DO NOTHING",
rusqlite::params![
id,
seq,
text(meta, "source"),
data.map(Json::to_string),
epoch_ms
],
)
.map_err(storage)?;
tx.execute(
"UPDATE cb_cards SET eval_count = eval_count + 1 WHERE id = ?1",
rusqlite::params![id],
)
.map_err(storage)?;
Ok(())
}
fn checkpoint(
tx: &Transaction<'_>,
id: &str,
seq: i64,
epoch_ms: i64,
data: Option<&Json>,
) -> Result<()> {
let blob = text(data, "blob");
let size = num(data.and_then(|d| d.get("size")));
tx.execute(
"INSERT INTO cb_checkpoints (card_id, seq, blob, size, format, note, epoch_ms)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT (card_id, seq) DO NOTHING",
rusqlite::params![
id,
seq,
blob.as_deref(),
size,
text(data, "format"),
text(data, "note"),
epoch_ms,
],
)
.map_err(storage)?;
tx.execute(
"UPDATE cb_cards SET checkpoint_count = checkpoint_count + 1 WHERE id = ?1",
rusqlite::params![id],
)
.map_err(storage)?;
if let Some(hash) = blob {
reference_blob(tx, &hash, size)?;
}
Ok(())
}
fn tag_set(tx: &Transaction<'_>, id: &str, epoch_ms: i64, meta: Option<&Json>) -> Result<()> {
let (Some(key), Some(value)) = (text(meta, "key"), text(meta, "value")) else {
return Err(Error::storage(format!(
"a tag_set on {STREAM_PREFIX}{id} carries no meta.key and meta.value"
)));
};
tx.execute(
"INSERT INTO cb_tags (card_id, key, value, set_ms) VALUES (?1, ?2, ?3, ?4)
ON CONFLICT (card_id, key) DO UPDATE SET value = excluded.value,
set_ms = excluded.set_ms",
rusqlite::params![id, key, value, epoch_ms],
)
.map_err(storage)?;
Ok(())
}
fn tag_unset(tx: &Transaction<'_>, id: &str, meta: Option<&Json>) -> Result<()> {
let Some(key) = text(meta, "key") else {
return Err(Error::storage(format!(
"a tag_unset on {STREAM_PREFIX}{id} carries no meta.key"
)));
};
tx.execute(
"DELETE FROM cb_tags WHERE card_id = ?1 AND key = ?2",
rusqlite::params![id, key],
)
.map_err(storage)?;
Ok(())
}
fn closed(
tx: &Transaction<'_>,
id: &str,
epoch_ms: i64,
meta: Option<&Json>,
data: Option<&Json>,
) -> Result<()> {
let state = match text(meta, "outcome").as_deref() {
Some("ok") => "closed_ok",
_ => "closed_failed",
};
let stats = data.and_then(|d| d.get("stats"));
let cost = data.and_then(|d| d.get("cost"));
tx.execute(
"UPDATE cb_cards SET state = ?2, closed_ms = ?3, error = ?4,
stats_json = ?5, cost_json = ?6,
mean_score = ?7, n = ?8, pass_rate = ?9, passed = ?10,
elapsed_ms = ?11, llm_calls = ?12
WHERE id = ?1",
rusqlite::params![
id,
state,
epoch_ms,
text(data, "error"),
stats.map(Json::to_string),
cost.map(Json::to_string),
real(stats.and_then(|s| s.get("mean_score"))),
num(stats.and_then(|s| s.get("n"))),
real(stats.and_then(|s| s.get("pass_rate"))),
num(stats.and_then(|s| s.get("passed"))),
num(cost.and_then(|c| c.get("elapsed_ms"))),
num(cost.and_then(|c| c.get("llm_calls"))),
],
)
.map_err(storage)?;
Ok(())
}
fn bound(
tx: &Transaction<'_>,
name: &str,
seq: i64,
epoch_ms: i64,
meta: Option<&Json>,
data: Option<&Json>,
) -> Result<()> {
let card_id = text(meta, "card_id");
let note = text(data, "note");
alias_logged(
tx,
name,
seq,
"alias_bound",
card_id.as_deref(),
epoch_ms,
note.as_deref(),
)?;
let Some(card_id) = card_id else {
return Err(Error::storage(format!(
"an alias_bound on {ALIAS_PREFIX}{name} carries no meta.card_id, \
so it binds the name to nothing"
)));
};
tx.execute(
"INSERT INTO cb_aliases (name, card_id, pkg, bound_ms, note) VALUES (?1, ?2, ?3, ?4, ?5)
ON CONFLICT (name) DO UPDATE SET card_id = excluded.card_id, pkg = excluded.pkg,
bound_ms = excluded.bound_ms, note = excluded.note",
rusqlite::params![name, card_id, text(meta, "pkg"), epoch_ms, note],
)
.map_err(storage)?;
Ok(())
}
fn released(
tx: &Transaction<'_>,
name: &str,
seq: i64,
epoch_ms: i64,
meta: Option<&Json>,
data: Option<&Json>,
) -> Result<()> {
alias_logged(
tx,
name,
seq,
"alias_released",
text(meta, "card_id").as_deref(),
epoch_ms,
text(data, "note").as_deref(),
)?;
tx.execute(
"DELETE FROM cb_aliases WHERE name = ?1",
rusqlite::params![name],
)
.map_err(storage)?;
Ok(())
}
fn alias_logged(
tx: &Transaction<'_>,
name: &str,
seq: i64,
kind: &str,
card_id: Option<&str>,
epoch_ms: i64,
note: Option<&str>,
) -> Result<()> {
tx.execute(
"INSERT INTO cb_alias_log (name, seq, kind, card_id, epoch_ms, note)
VALUES (?1, ?2, ?3, ?4, ?5, ?6)
ON CONFLICT (name, seq) DO NOTHING",
rusqlite::params![name, seq, kind, card_id, epoch_ms, note],
)
.map_err(storage)?;
Ok(())
}
fn on_prune_stream(stream: &str, kind: &str) -> Result<()> {
if stream == PRUNE_STREAM {
return Ok(());
}
Err(Error::storage(format!(
"a {kind:?} event is on stream {stream:?}: the prune journal is the one stream \
{PRUNE_STREAM:?}"
)))
}
fn pruned(tx: &Transaction<'_>, data: Option<&Json>) -> Result<()> {
let cards = data.and_then(|d| d.get("cards")).and_then(Json::as_array);
for card in cards.into_iter().flatten() {
if let Some(id) = card.as_str() {
purge(tx, id)?;
}
}
Ok(())
}
fn purge(tx: &Transaction<'_>, id: &str) -> Result<()> {
let present: Option<i64> = tx
.query_row(
"SELECT 1 FROM cb_cards WHERE id = ?1",
rusqlite::params![id],
|row| row.get(0),
)
.optional()
.map_err(storage)?;
if present.is_none() {
return Ok(());
}
let named: Option<String> = tx
.query_row(
"SELECT name FROM cb_aliases WHERE card_id = ?1 ORDER BY name LIMIT 1",
rusqlite::params![id],
|row| row.get(0),
)
.optional()
.map_err(storage)?;
if let Some(name) = named {
return Err(Error::storage(format!(
"card {id} was pruned while the alias {name:?} still points at it: a card with \
a name is not prunable, and the name would be left pointing at nothing"
)));
}
let hashes: Vec<String> = {
let mut stmt = tx
.prepare(
"SELECT blob FROM cb_samples WHERE card_id = ?1 AND blob IS NOT NULL \
UNION ALL \
SELECT blob FROM cb_checkpoints WHERE card_id = ?1 AND blob IS NOT NULL",
)
.map_err(storage)?;
let rows = stmt
.query_map(rusqlite::params![id], |row| row.get::<_, String>(0))
.map_err(storage)?;
rows.collect::<rusqlite::Result<Vec<String>>>()
.map_err(storage)?
};
for hash in hashes {
tx.execute(
"UPDATE cb_blobs SET refs = refs - 1 WHERE hash = ?1",
rusqlite::params![hash],
)
.map_err(storage)?;
}
for sql in [
"DELETE FROM cb_samples WHERE card_id = ?1",
"DELETE FROM cb_evals WHERE card_id = ?1",
"DELETE FROM cb_tags WHERE card_id = ?1",
"DELETE FROM cb_checkpoints WHERE card_id = ?1",
"DELETE FROM cb_lineage WHERE child = ?1 OR parent = ?1",
"DELETE FROM cb_cards WHERE id = ?1",
] {
tx.execute(sql, rusqlite::params![id]).map_err(storage)?;
}
Ok(())
}
fn reference_blob(tx: &Transaction<'_>, hash: &str, size: Option<i64>) -> Result<()> {
tx.execute(
"INSERT INTO cb_blobs (hash, size, refs) VALUES (?1, ?2, 1)
ON CONFLICT (hash) DO UPDATE SET refs = refs + 1, size = COALESCE(excluded.size, size)",
rusqlite::params![hash, size],
)
.map_err(storage)?;
Ok(())
}
fn text(object: Option<&Json>, key: &str) -> Option<String> {
object
.and_then(|o| o.get(key))
.and_then(Json::as_str)
.map(str::to_string)
}
fn num(value: Option<&Json>) -> Option<i64> {
value.and_then(|v| v.as_i64().or_else(|| v.as_f64().map(|f| f as i64)))
}
fn real(value: Option<&Json>) -> Option<f64> {
value.and_then(Json::as_f64)
}
fn storage(e: rusqlite::Error) -> Error {
Error::storage(e.to_string())
}