use std::time::Instant;
use libsql::Value;
use macrame::branch::BranchId;
use macrame::graph::EdgeAssertion;
use macrame::{ConceptUpsert, Database};
const GENESIS: &str = "2020-01-01T00:00:00.000000Z";
const HEAD: &str = r#"WITH RECURSIVE lineage(branch_id, dist, cutoff) AS (VALUES (?6, ?7, ?8), (?9, ?10, ?11)),
"#;
const CHURNED: &str = r#"churned(entity_id, branch_id, cutoff) AS (
SELECT lc.source_id || '|' || lc.target_id || '|' || lc.edge_type || '|' || lc.valid_from,
lc.branch_id, g.cutoff
FROM links_current lc
JOIN lineage g ON g.branch_id = lc.branch_id
WHERE lc.source_id = ?1 AND lc.target_id = ?2 AND lc.edge_type = ?3 AND g.cutoff IS NOT NULL AND lc.recorded_at > g.cutoff
),
"#;
const CHURNED_MAT: &str = r#"churned(entity_id, branch_id, cutoff) AS MATERIALIZED (
SELECT lc.source_id || '|' || lc.target_id || '|' || lc.edge_type || '|' || lc.valid_from,
lc.branch_id, g.cutoff
FROM links_current lc
JOIN lineage g ON g.branch_id = lc.branch_id
WHERE lc.source_id = ?1 AND lc.target_id = ?2 AND lc.edge_type = ?3 AND g.cutoff IS NOT NULL AND lc.recorded_at > g.cutoff
),
"#;
const CUT_OPEN: &str = r#"links_cut(source_id, target_id, edge_type, valid_from, valid_to, weight, properties, branch_id) AS (
SELECT lc.source_id, lc.target_id, lc.edge_type, lc.valid_from, lc.valid_to,
lc.weight, lc.properties, lc.branch_id
FROM links_current lc
JOIN lineage g ON g.branch_id = lc.branch_id
WHERE lc.source_id = ?1 AND lc.target_id = ?2 AND lc.edge_type = ?3 AND (g.cutoff IS NULL OR lc.recorded_at <= g.cutoff)
UNION ALL
SELECT json_extract(payload, '$.source_id'),
json_extract(payload, '$.target_id'),
json_extract(payload, '$.edge_type'),
json_extract(payload, '$.valid_from'),
json_extract(payload, '$.valid_to'),
json_extract(payload, '$.weight'),
json_extract(payload, '$.properties'),
branch_id
FROM (
SELECT transaction_log.payload, transaction_log.branch_id,
ROW_NUMBER() OVER (
PARTITION BY transaction_log.entity_id, transaction_log.branch_id
ORDER BY transaction_log.seq_id DESC
) AS rn
"#;
const CUT_CLOSE: &str = r#" ) WHERE rn = 1
),
"#;
const ARM_SHIPPED: &str = r#" FROM transaction_log
JOIN churned k ON k.entity_id = transaction_log.entity_id
AND k.branch_id = transaction_log.branch_id
WHERE transaction_log.table_name = 'links'
AND transaction_log.recorded_at <= k.cutoff
"#;
const ARM_CROSS: &str = r#" FROM churned k
CROSS JOIN transaction_log ON transaction_log.entity_id = k.entity_id
AND transaction_log.branch_id = k.branch_id
WHERE transaction_log.table_name = 'links'
AND transaction_log.recorded_at <= k.cutoff
"#;
const ARM_INDEXED: &str = r#" FROM transaction_log INDEXED BY idx_txlog_entity
JOIN churned k ON k.entity_id = transaction_log.entity_id
AND k.branch_id = transaction_log.branch_id
WHERE transaction_log.table_name = 'links'
AND transaction_log.recorded_at <= k.cutoff
"#;
const VISIBLE: &str = r#"visible(source_id, target_id, edge_type, valid_from, valid_to, weight, properties, branch_id) AS (
SELECT source_id, target_id, edge_type, valid_from, valid_to, weight, properties, branch_id FROM (
SELECT l.source_id, l.target_id, l.edge_type, l.valid_from, l.valid_to, l.weight,
l.properties, l.branch_id,
ROW_NUMBER() OVER (
PARTITION BY l.source_id, l.target_id, l.edge_type, l.valid_from
ORDER BY g.dist
) AS rn
FROM links_cut l
JOIN lineage g ON g.branch_id = l.branch_id
) WHERE rn = 1
)
SELECT l.valid_from, l.valid_to FROM visible l WHERE l.valid_from <> ?4"#;
const TRUNK_GUARD: &str =
"SELECT l.valid_from, l.valid_to FROM links_current l \
WHERE l.valid_from <> ?4 AND l.source_id = ?1 AND l.target_id = ?2 AND l.edge_type = ?3";
fn resolved(churned: &str, arm: &str) -> String {
format!("{HEAD}{churned}{CUT_OPEN}{arm}{CUT_CLOSE}{VISIBLE}")
}
fn edge(i: usize) -> EdgeAssertion {
EdgeAssertion::new("hub", format!("t{i}"), "LINKS")
.valid_from(format!("2026-01-01T00:00:00.{:06}Z", i % 1_000_000))
.valid_to(format!("2026-01-01T00:00:00.{:06}Z", (i % 1_000_000) + 1))
}
fn params(i: usize, branch: &str, cutoff: &str) -> Vec<Value> {
vec![
Value::Text("hub".into()),
Value::Text(format!("t{i}")),
Value::Text("LINKS".into()),
Value::Text(format!("2026-01-01T00:00:00.{:06}Z", i % 1_000_000)),
Value::Text(branch.into()),
Value::Text(branch.into()),
Value::Integer(0),
Value::Null,
Value::Text("main".into()),
Value::Integer(1),
Value::Text(cutoff.into()),
]
}
async fn plan(conn: &libsql::Connection, sql: &str, p: Vec<Value>) -> Vec<String> {
let mut rows = conn
.query(&format!("EXPLAIN QUERY PLAN {sql}"), p)
.await
.unwrap();
let mut out = Vec::new();
while let Some(row) = rows.next().await.unwrap() {
out.push(row.get::<String>(3).unwrap());
}
out
}
async fn scalar(conn: &libsql::Connection, sql: &str) -> i64 {
conn.query(sql, ())
.await
.unwrap()
.next()
.await
.unwrap()
.unwrap()
.get::<i64>(0)
.unwrap()
}
async fn time_guard(
conn: &libsql::Connection,
sql: &str,
batch: usize,
branch: &str,
cutoff: &str,
) -> f64 {
let stmt_start = Instant::now();
let stmt = conn.prepare(sql).await.unwrap();
let compile = stmt_start.elapsed().as_secs_f64() * 1e3;
let start = Instant::now();
for i in 0..batch {
let mut rows = stmt.query(params(i, branch, cutoff)).await.unwrap();
while rows.next().await.unwrap().is_some() {}
stmt.reset();
}
let ms = start.elapsed().as_secs_f64() * 1e3;
println!(" (compile {compile:.3} ms)");
ms
}
async fn run_end_to_end(
db: &Database,
trunk: usize,
batch: usize,
) -> Result<(), Box<dyn std::error::Error>> {
println!("end to end, bulk_import of one {batch}-edge batch");
let start = Instant::now();
db.bulk_import((trunk..trunk + batch).map(edge).collect())
.await?;
println!(
" {:>12} {:>12.2} ms",
"main",
start.elapsed().as_secs_f64() * 1e3
);
let start = Instant::now();
db.bulk_import(
(trunk..trunk + batch)
.map(|i| edge(i).on_branch(BranchId::new("alt").unwrap()))
.collect(),
)
.await?;
println!(
" {:>12} {:>12.2} ms",
"alt",
start.elapsed().as_secs_f64() * 1e3
);
Ok(())
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let args: Vec<String> = std::env::args().collect();
let flag = |name: &str, default: usize| -> usize {
args.iter()
.position(|a| a == name)
.and_then(|i| args.get(i + 1))
.and_then(|v| v.parse().ok())
.unwrap_or(default)
};
let trunk = flag("--trunk", 2_000);
let batch = flag("--batch", 200);
let skip_end = args.iter().any(|a| a == "--skip-end-to-end");
let skip_guards = args.iter().any(|a| a == "--skip-guards");
let dir = tempfile::tempdir()?;
let db = Database::open(dir.path().join("probe.db")).await?;
println!("branch write guard, trunk {trunk} edges, batch {batch} edges\n");
db.upsert_concept(ConceptUpsert::new("hub", "hub").valid_from(GENESIS))
.await?;
for i in 0..(trunk + batch) {
db.upsert_concept(ConceptUpsert::new(format!("t{i}"), "t").valid_from(GENESIS))
.await?;
}
db.bulk_import((0..trunk).map(edge).collect()).await?;
db.fork(BranchId::new("alt")?, BranchId::new("main")?)
.await?;
let conn = db.read_conn();
let cutoff: String = conn
.query("SELECT forked_at FROM branches WHERE branch_id = 'alt'", ())
.await?
.next()
.await?
.unwrap()
.get(0)?;
println!("fixture");
println!(
" links_current {:>8} transaction_log {:>8} of those table_name='links' {:>8}",
scalar(conn, "SELECT count(*) FROM links_current").await,
scalar(conn, "SELECT count(*) FROM transaction_log").await,
scalar(
conn,
"SELECT count(*) FROM transaction_log WHERE table_name = 'links'"
)
.await,
);
println!(" fork cutoff {cutoff}\n");
let candidates: Vec<(&str, String)> = vec![
("shipped", resolved(CHURNED, ARM_SHIPPED)),
("cross", resolved(CHURNED, ARM_CROSS)),
("indexed", resolved(CHURNED, ARM_INDEXED)),
("materialized", resolved(CHURNED_MAT, ARM_SHIPPED)),
("mat+cross", resolved(CHURNED_MAT, ARM_CROSS)),
];
if skip_guards {
run_end_to_end(&db, trunk, batch).await?;
db.close().await?;
return Ok(());
}
println!("plans (branch shape, the parameters the write path binds)");
for (name, sql) in &candidates {
println!(" {name}");
for step in plan(conn, sql, params(0, "alt", &cutoff)).await {
println!(" {step}");
}
}
println!(" trunk");
for step in plan(conn, TRUNK_GUARD, params(0, "main", &cutoff)).await {
println!(" {step}");
}
println!();
println!("the guard alone, {batch} executions (one per row of the batch)");
for (name, sql) in &candidates {
let ms = time_guard(conn, sql, batch, "alt", &cutoff).await;
println!(
" {name:>12} {ms:>12.2} ms {:>9.4} ms/row",
ms / batch as f64
);
}
let ms = time_guard(conn, TRUNK_GUARD, batch, "main", &cutoff).await;
println!(
" {:>12} {ms:>12.2} ms {:>9.4} ms/row",
"trunk",
ms / batch as f64
);
println!();
if !skip_end {
run_end_to_end(&db, trunk, batch).await?;
}
db.close().await?;
Ok(())
}