use mongreldb_core::schema::{
AnnOptions, AnnQuantization, ColumnDef, ColumnFlags, IndexDef, IndexKind, IndexOptions, Schema,
TypeId,
};
use mongreldb_core::{Database, Table, Value};
use std::io::Write;
use std::process::exit;
use std::sync::Arc;
use std::time::{Duration, Instant};
fn schema() -> Schema {
Schema {
schema_id: 1,
columns: vec![ColumnDef {
id: 1,
name: "v".into(),
ty: TypeId::Int64,
flags: ColumnFlags::empty().with(ColumnFlags::PRIMARY_KEY),
default_value: None,
embedding_source: None,
}],
indexes: Vec::new(),
colocation: Vec::new(),
constraints: Default::default(),
clustered: false,
}
}
fn main() {
let mut args = std::env::args().skip(1);
let (dir, mode) = match (args.next(), args.next()) {
(Some(d), Some(m)) => (std::path::PathBuf::from(d), m),
_ => {
eprintln!("usage: crash_writer <db_dir> <mode>");
exit(2);
}
};
match mode.as_str() {
"committed" => {
let mut db = Table::create(&dir, schema(), 1).expect("create");
db.put(vec![(1, Value::Int64(4242))]).expect("put");
db.commit().expect("commit");
signal("COMMITTED_READY");
spin();
}
"committed-burst" => {
let mut db = Table::create(&dir, schema(), 1).expect("create");
for v in [10i64, 20, 30] {
db.put(vec![(1, Value::Int64(v))]).expect("put");
db.commit().expect("commit");
}
signal("BURST_READY");
spin();
}
"uncommitted" => {
let mut db = Table::create(&dir, schema(), 1).expect("create");
db.set_sync_byte_threshold(1);
db.put(vec![(1, Value::Int64(9999))]).expect("put");
signal("UNCOMMITTED_READY");
spin();
}
"flush-spill" => {
let mut db = Table::create(&dir, schema(), 1).expect("create");
db.set_mutable_run_spill_bytes(1);
db.put(vec![(1, Value::Int64(7777))]).expect("put");
db.flush().expect("flush");
signal("FLUSH_READY");
spin();
}
"ctas-building" => {
let db = Database::create(&dir).expect("create database");
let build = "__mongreldb_ctas_build_crash-query";
db.create_building_table(build, "target", "crash-query", schema())
.expect("create building table");
let mut txn = db.begin();
txn.put_building(build, vec![(1, Value::Int64(4242))])
.expect("put building row");
txn.commit().expect("commit building row");
signal("CTAS_BUILDING_READY");
spin();
}
"shared-commit" => {
let value: i64 = args.next().and_then(|v| v.parse().ok()).unwrap_or(1);
let db = match Database::open(&dir) {
Ok(db) => db,
Err(_) => Database::create(&dir).expect("create database"),
};
if db.table("t").is_err() {
db.create_table("t", schema()).expect("create table");
}
let t = db.table("t").expect("table");
{
let mut g = t.lock();
g.put(vec![(1, Value::Int64(value))]).expect("put");
g.commit().expect("commit");
}
signal("SHARED_READY");
spin();
}
"shared-commit-encrypted" => {
const PASSPHRASE: &str = "crash-harness-test-passphrase";
let value: i64 = args.next().and_then(|v| v.parse().ok()).unwrap_or(1);
let db = match Database::open_encrypted(&dir, PASSPHRASE) {
Ok(db) => db,
Err(_) => Database::create_encrypted(&dir, PASSPHRASE).expect("create database"),
};
if db.table("t").is_err() {
db.create_table("t", schema()).expect("create table");
}
let t = db.table("t").expect("table");
{
let mut g = t.lock();
g.put(vec![(1, Value::Int64(value))]).expect("put");
g.commit().expect("commit");
}
signal("SHARED_READY");
spin();
}
"durable-hook" => {
let hook = args.next().unwrap_or_else(|| {
eprintln!("durable-hook requires a hook name");
exit(2);
});
let hook = durable_hook(&hook);
match hook {
"wal.append.before"
| "wal.append.after"
| "wal.fsync.before"
| "wal.fsync.after"
| "commit.publish.before"
| "commit.publish.after" => {
let db = database_with_row(&dir);
arm_crash_hook(hook);
let mut transaction = db.begin();
transaction
.put("items", vec![(1, Value::Int64(2))])
.expect("stage hook row");
let _ = transaction.commit();
}
"catalog.publish.before" | "catalog.publish.after" => {
let db = database_with_row(&dir);
arm_crash_hook(hook);
let _ = db.create_table("published", schema());
}
"snapshot.install.before" | "snapshot.install.after" => {
let leader = database_with_row(&dir.join("leader"));
let snapshot = leader.replication_snapshot().expect("snapshot");
arm_crash_hook(hook);
let _ = snapshot.install(dir.join("follower"));
}
"index.publish.before" | "index.publish.after" => {
let db = database_with_row(&dir);
db.table("items")
.expect("items")
.lock()
.set_mutable_run_spill_bytes(1);
arm_crash_hook(hook);
let _ = db.table("items").expect("items").lock().flush();
}
_ => unreachable!(),
}
panic!("durable hook {hook} was not hit");
}
"index-ddl-hook" => {
let hook = args.next().unwrap_or_else(|| {
eprintln!("index-ddl-hook requires a hook name");
exit(2);
});
let hook = durable_hook(&hook);
if !matches!(hook, "index.publish.before" | "index.publish.after") {
eprintln!("index-ddl-hook requires an index publication hook");
exit(2);
}
let db = Database::create(&dir).expect("create database");
db.create_table("docs", embedding_schema())
.expect("create docs");
db.transaction(|transaction| {
transaction.put(
"docs",
vec![
(1, Value::Int64(1)),
(2, Value::Embedding(vec![1.0, 0.0, 0.0, 0.0])),
],
)?;
Ok(())
})
.expect("insert docs row");
arm_crash_hook(hook);
let _ = db.create_index("docs", dense_index());
panic!("index DDL hook {hook} was not hit");
}
other => {
eprintln!("unknown mode: {other}");
exit(2);
}
}
}
fn embedding_schema() -> Schema {
Schema {
columns: vec![
ColumnDef {
id: 1,
name: "id".into(),
ty: TypeId::Int64,
flags: ColumnFlags::empty().with(ColumnFlags::PRIMARY_KEY),
default_value: None,
embedding_source: None,
},
ColumnDef {
id: 2,
name: "embedding".into(),
ty: TypeId::Embedding { dim: 4 },
flags: ColumnFlags::empty(),
default_value: None,
embedding_source: None,
},
],
..Schema::default()
}
}
fn dense_index() -> IndexDef {
IndexDef {
name: "idx_embedding".into(),
column_id: 2,
kind: IndexKind::Ann,
predicate: None,
options: IndexOptions {
ann: Some(AnnOptions {
quantization: AnnQuantization::Dense,
..AnnOptions::default()
}),
..IndexOptions::default()
},
}
}
fn database_with_row(path: &std::path::Path) -> Database {
let db = Database::create(path).expect("create database");
db.create_table("items", schema()).expect("create table");
db.transaction(|transaction| {
transaction.put("items", vec![(1, Value::Int64(1))])?;
Ok(())
})
.expect("commit baseline row");
db
}
fn durable_hook(name: &str) -> &'static str {
match name {
"wal.append.before" => "wal.append.before",
"wal.append.after" => "wal.append.after",
"wal.fsync.before" => "wal.fsync.before",
"wal.fsync.after" => "wal.fsync.after",
"commit.publish.before" => "commit.publish.before",
"commit.publish.after" => "commit.publish.after",
"catalog.publish.before" => "catalog.publish.before",
"catalog.publish.after" => "catalog.publish.after",
"snapshot.install.before" => "snapshot.install.before",
"snapshot.install.after" => "snapshot.install.after",
"index.publish.before" => "index.publish.before",
"index.publish.after" => "index.publish.after",
other => {
eprintln!("unknown durable hook: {other}");
exit(2);
}
}
}
fn arm_crash_hook(hook: &'static str) {
mongreldb_fault::activate(
hook,
mongreldb_fault::Action::Callback(Arc::new(|_| {
signal("HOOK_READY");
spin();
})),
);
}
fn signal(msg: &str) {
println!("{msg}");
let _ = std::io::stdout().flush();
}
fn spin() {
let deadline = Instant::now() + Duration::from_secs(180);
while Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(20));
}
}