use std::time::{Duration, Instant};
use macrame::prelude::*;
use macrame::{Branch, BranchId};
const T0: &str = "2026-01-01T00:00:00.000000Z";
const FOREVER: &str = "9999-12-31T23:59:59.999999Z";
async fn best_of<F, Fut>(rounds: usize, mut f: F) -> Duration
where
F: FnMut() -> Fut,
Fut: std::future::Future<Output = ()>,
{
let mut best = Duration::MAX;
for _ in 0..rounds {
let t = Instant::now();
f().await;
best = best.min(t.elapsed());
}
best
}
async fn one_row(conn: &libsql::Connection, sql: &str, params: Vec<libsql::Value>) {
let mut r = conn.query(sql, params).await.expect("query");
r.next().await.expect("step").expect("one row");
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Ancestor {
branch_id: String,
dist: i64,
cutoff: Option<String>,
}
fn resolve(branches: &[Branch], start: &str) -> Vec<Ancestor> {
let find = |name: &str| branches.iter().find(|b| b.id.as_str() == name);
let mut out = Vec::new();
let mut cur = start.to_string();
let mut cutoff: Option<String> = None;
for dist in 0..=(branches.len() as i64) {
out.push(Ancestor {
branch_id: cur.clone(),
dist,
cutoff: cutoff.clone(),
});
let Some(node) = find(&cur) else { break };
let (Some(parent), Some(forked)) = (node.parent.as_ref(), node.forked_at.as_ref()) else {
break;
};
cutoff = Some(match cutoff {
Some(c) if c.as_str() <= forked.as_str() => c,
_ => forked.clone(),
});
cur = parent.as_str().to_string();
}
out
}
fn ancestry_values(rows: &[Ancestor], first_slot: usize) -> String {
let tuples: Vec<String> = (0..rows.len())
.map(|i| {
let b = first_slot + i * 3;
format!("(?{}, ?{}, ?{})", b, b + 1, b + 2)
})
.collect();
format!(
"lineage(branch_id, dist, cutoff) AS (VALUES {})",
tuples.join(", ")
)
}
fn ancestry_values_lit_dist(rows: &[Ancestor], first_slot: usize) -> String {
let tuples: Vec<String> = (0..rows.len())
.map(|i| {
let b = first_slot + i * 2;
format!("(?{}, {}, ?{})", b, i, b + 1)
})
.collect();
format!(
"lineage(branch_id, dist, cutoff) AS (VALUES {})",
tuples.join(", ")
)
}
fn params_lit_dist(rows: &[Ancestor]) -> Vec<libsql::Value> {
let mut v = Vec::with_capacity(rows.len() * 2);
for r in rows {
v.push(libsql::Value::Text(r.branch_id.clone()));
v.push(match &r.cutoff {
Some(c) => libsql::Value::Text(c.clone()),
None => libsql::Value::Null,
});
}
v
}
fn ancestry_recursive(slot: usize) -> String {
format!(
r#"lineage(branch_id, dist, cutoff) AS (
SELECT ?{slot}, 0, NULL
UNION ALL
SELECT b.parent_id, g.dist + 1,
CASE WHEN g.cutoff IS NULL OR b.forked_at < g.cutoff
THEN b.forked_at ELSE g.cutoff END
FROM branches b JOIN lineage g ON b.branch_id = g.branch_id
WHERE b.parent_id IS NOT NULL
)"#
)
}
fn params_for(rows: &[Ancestor]) -> Vec<libsql::Value> {
let mut v = Vec::with_capacity(rows.len() * 3);
for r in rows {
v.push(libsql::Value::Text(r.branch_id.clone()));
v.push(libsql::Value::Integer(r.dist));
v.push(match &r.cutoff {
Some(c) => libsql::Value::Text(c.clone()),
None => libsql::Value::Null,
});
}
v
}
async fn read_ancestry(
conn: &libsql::Connection,
sql: &str,
params: Vec<libsql::Value>,
) -> Vec<Ancestor> {
let mut rows = conn.query(sql, params).await.expect("ancestry query");
let mut out = Vec::new();
while let Some(r) = rows.next().await.expect("row") {
out.push(Ancestor {
branch_id: r.get::<String>(0).expect("branch_id"),
dist: r.get::<i64>(1).expect("dist"),
cutoff: r.get::<Option<String>>(2).expect("cutoff"),
});
}
out.sort_by_key(|a| a.dist);
out
}
#[tokio::main(flavor = "multi_thread", worker_threads = 4)]
async fn main() {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.path().join("ancestry.db");
let db = Database::open(&path).await.expect("open");
let mut concepts = Vec::new();
for i in 0..400 {
concepts.push(ConceptUpsert::new(format!("c{i}"), "N").valid_from(T0));
}
db.write_concepts(concepts).await.expect("concepts");
let mut edges = Vec::new();
for i in 0..399 {
edges.push(
EdgeAssertion::new(format!("c{i}"), format!("c{}", i + 1), "LINKS")
.valid_from(T0)
.valid_to(FOREVER),
);
}
db.bulk_import(edges).await.expect("edges");
let mut chain = vec!["main".to_string()];
for i in 1..=16 {
let name = format!("b{i}");
let parent = chain.last().expect("parent").clone();
db.fork(
BranchId::new(&name).expect("name"),
BranchId::new(&parent).expect("parent"),
)
.await
.expect("fork");
chain.push(name);
}
let branches = db.branches().await.expect("branches");
let conn = db.diagnostic_conn().await.expect("diagnostic");
println!("=== 1. does libSQL accept a bound VALUES CTE? ===\n");
for depth in [1_usize, 4, 16] {
let start = &chain[depth];
let rust = resolve(&branches, start);
let sql = format!(
"WITH {} SELECT branch_id, dist, cutoff FROM lineage",
ancestry_values(&rust, 1)
);
let got = read_ancestry(&conn, &sql, params_for(&rust)).await;
println!(
" depth {depth:>2}: {} rows bound, read back {} -- {}",
rust.len(),
got.len(),
if got == rust { "identical" } else { "DIFFERS" }
);
}
println!("\n=== 2. the Rust walk against the CTE it replaces ===\n");
let mut disagreements = 0;
for start in &chain {
let cte = format!(
"WITH RECURSIVE {} SELECT branch_id, dist, cutoff FROM lineage",
ancestry_recursive(1)
);
let from_sql = read_ancestry(&conn, &cte, vec![libsql::Value::Text(start.clone())]).await;
let from_rust = resolve(&branches, start);
let same = from_sql == from_rust;
if !same {
disagreements += 1;
println!(" {start:>4}: DIFFERS\n sql {from_sql:?}\n rust {from_rust:?}");
}
}
println!(
" {} lineages compared, {} disagreements",
chain.len(),
disagreements
);
println!("\n=== 3. the ancestry alone: recursive against bound ===\n");
for depth in [1_usize, 4, 16] {
let start = chain[depth].clone();
let rust = resolve(&branches, &start);
let cte = format!(
"WITH RECURSIVE {} SELECT COUNT(*) FROM lineage",
ancestry_recursive(1)
);
let rec = best_of(200, || {
one_row(&conn, &cte, vec![libsql::Value::Text(start.clone())])
})
.await;
let vals = format!(
"WITH {} SELECT COUNT(*) FROM lineage",
ancestry_values(&rust, 1)
);
let bound = best_of(200, || one_row(&conn, &vals, params_for(&rust))).await;
println!(
" depth {depth:>2} ({} rows): recursive {:>7.2} µs bound {:>7.2} µs {:+.1}%",
rust.len(),
rec.as_secs_f64() * 1e6,
bound.as_secs_f64() * 1e6,
(bound.as_secs_f64() / rec.as_secs_f64() - 1.0) * 100.0,
);
}
println!("\n=== 4. joined the way the readers join it ===\n");
for depth in [1_usize, 2, 4, 6, 8, 12, 16] {
let start = chain[depth].clone();
let rust = resolve(&branches, &start);
let body = "SELECT COUNT(*) FROM links_current lc JOIN lineage g \
ON lc.branch_id = g.branch_id \
WHERE g.cutoff IS NOT NULL AND lc.recorded_at > g.cutoff";
let cte = format!("WITH RECURSIVE {} {body}", ancestry_recursive(1));
let rec = best_of(120, || {
one_row(&conn, &cte, vec![libsql::Value::Text(start.clone())])
})
.await;
let vals = format!("WITH {} {body}", ancestry_values(&rust, 1));
let bound = best_of(120, || one_row(&conn, &vals, params_for(&rust))).await;
println!(
" depth {depth:>2}: recursive {:>7.2} µs bound {:>7.2} µs {:+.1}%",
rec.as_secs_f64() * 1e6,
bound.as_secs_f64() * 1e6,
(bound.as_secs_f64() / rec.as_secs_f64() - 1.0) * 100.0,
);
}
println!("\n=== 7. is the slope the parameters? dist bound against dist literal ===\n");
for depth in [1_usize, 4, 8, 16] {
let start = chain[depth].clone();
let rust = resolve(&branches, &start);
let body = "SELECT COUNT(*) FROM links_current lc JOIN lineage g \
ON lc.branch_id = g.branch_id \
WHERE g.cutoff IS NOT NULL AND lc.recorded_at > g.cutoff";
let three = format!("WITH {} {body}", ancestry_values(&rust, 1));
let p3 = best_of(120, || one_row(&conn, &three, params_for(&rust))).await;
let two = format!("WITH {} {body}", ancestry_values_lit_dist(&rust, 1));
let p2 = best_of(120, || one_row(&conn, &two, params_lit_dist(&rust))).await;
println!(
" depth {depth:>2}: 3/row {:>7.2} \u{00b5}s ({:>2} params) 2/row {:>7.2} \u{00b5}s ({:>2} params) {:+.1}%",
p3.as_secs_f64() * 1e6,
rust.len() * 3,
p2.as_secs_f64() * 1e6,
rust.len() * 2,
(p2.as_secs_f64() / p3.as_secs_f64() - 1.0) * 100.0,
);
}
println!("\n=== 5. the read-side round trip: aggregates against rows ===\n");
let aggregates = "SELECT (SELECT COUNT(*) FROM branches), \
(SELECT COUNT(*) FROM branches WHERE branch_id = ?1), \
(SELECT COUNT(*) FROM branches \
WHERE branch_id = ?1 AND parent_id IS NULL)";
let agg = best_of(300, || {
one_row(&conn, aggregates, vec![libsql::Value::Text("b16".into())])
})
.await;
let rows_sql = "SELECT branch_id, parent_id, forked_at FROM branches";
let rows = best_of(300, || async {
let mut r = conn.query(rows_sql, ()).await.expect("rows");
let mut n = 0;
while r.next().await.expect("row").is_some() {
n += 1;
}
assert_eq!(n, 17, "17 lineages");
})
.await;
println!(
" three aggregates : {:>7.2} µs\n 17 rows loaded : {:>7.2} µs {:+.1}%",
agg.as_secs_f64() * 1e6,
rows.as_secs_f64() * 1e6,
(rows.as_secs_f64() / agg.as_secs_f64() - 1.0) * 100.0,
);
println!("\n=== 6. the running minimum, on the shape that needs it ===\n");
let db2dir = tempfile::tempdir().expect("tempdir2");
let db2 = Database::open(db2dir.path().join("clamp.db"))
.await
.expect("open2");
db2.write_concepts(vec![ConceptUpsert::new("x", "N").valid_from(T0)])
.await
.expect("w");
db2.fork(
BranchId::new("early").expect("n"),
BranchId::new("main").expect("p"),
)
.await
.expect("fork early");
db2.fork(
BranchId::new("late").expect("n"),
BranchId::new("early").expect("p"),
)
.await
.expect("fork late");
let b2 = db2.branches().await.expect("branches2");
let conn2 = db2.diagnostic_conn().await.expect("diag2");
let cte = format!(
"WITH RECURSIVE {} SELECT branch_id, dist, cutoff FROM lineage",
ancestry_recursive(1)
);
let from_sql = read_ancestry(&conn2, &cte, vec![libsql::Value::Text("late".into())]).await;
let from_rust = resolve(&b2, "late");
println!(" sql {from_sql:?}");
println!(" rust {from_rust:?}");
println!(
" {}",
if from_sql == from_rust {
"identical -- the clamp is reproduced"
} else {
"DIFFERS -- the clamp is not reproduced"
}
);
db.close().await.expect("close");
db2.close().await.expect("close2");
}