use core_api::{BatchOp, GraphDb, OpenOptions, SharedDb};
use core_storage::GraphError;
use std::path::Path;
use std::time::Duration;
const CHUNK: usize = 50;
const WRITE_RETRIES: usize = 8;
fn chunk_ops(prefix: &str, lo: usize, hi: usize) -> Vec<BatchOp> {
(lo..hi)
.map(|k| BatchOp::InsertNode {
label: "Person".into(),
key: format!("{prefix}-{k}"),
props: vec![("team".into(), core_api::Value::Str(prefix.into()))],
})
.collect()
}
fn main() {
let args: Vec<String> = std::env::args().collect();
if args.len() < 3 {
eprintln!("usage: mp_worker <dir> <command> [args...]");
std::process::exit(2);
}
let dir = Path::new(&args[1]);
let cmd = args[2].as_str();
let code = match cmd {
"write" => {
let prefix = args.get(3).map(String::as_str).unwrap_or("n");
let n: usize = args.get(4).and_then(|s| s.parse().ok()).unwrap_or(0);
run_write(dir, prefix, n)
}
"read" => run_read(dir),
"ro-read" => run_ro_read(dir),
"busy" => {
let ms: u64 = args.get(3).and_then(|s| s.parse().ok()).unwrap_or(500);
run_busy(dir, Duration::from_millis(ms))
}
"snapshot" => run_snapshot(dir),
"snapshot-shared" => run_snapshot_shared(dir),
other => {
eprintln!("unknown command: {other}");
std::process::exit(2);
}
};
std::process::exit(code);
}
fn exit_code(e: &GraphError) -> i32 {
match e {
GraphError::Busy { .. } => 3,
_ => 1,
}
}
fn run_write(dir: &Path, prefix: &str, n: usize) -> i32 {
let db = match SharedDb::open(dir) {
Ok(db) => db,
Err(e) => {
eprintln!("open: {e}");
return exit_code(&e);
}
};
let mut i = 0usize;
while i < n {
let hi = (i + CHUNK).min(n);
let ops = chunk_ops(prefix, i, hi);
let mut ops = Some(ops);
let mut attempt = 0;
loop {
match db.submit_batch(ops.take().expect("ops present on each attempt")) {
Ok(_) => break,
Err(GraphError::Busy { .. }) if attempt < WRITE_RETRIES => {
attempt += 1;
std::thread::sleep(Duration::from_millis(20 * attempt as u64));
ops = Some(chunk_ops(prefix, i, hi));
}
Err(e) => {
eprintln!("submit_batch: {e}");
return exit_code(&e);
}
}
}
i = hi;
}
0
}
fn run_read(dir: &Path) -> i32 {
match SharedDb::open(dir) {
Ok(db) => {
println!("{}", db.read().node_count());
0
}
Err(e) => {
eprintln!("open: {e}");
exit_code(&e)
}
}
}
fn run_ro_read(dir: &Path) -> i32 {
let opts = OpenOptions {
read_only: true,
..OpenOptions::default()
};
match GraphDb::open_with_options(dir, opts) {
Ok(db) => {
println!("{}", db.node_count());
0
}
Err(e) => {
eprintln!("open: {e}");
exit_code(&e)
}
}
}
fn run_busy(dir: &Path, wait: Duration) -> i32 {
let db = match SharedDb::open(dir) {
Ok(db) => db,
Err(e) => {
eprintln!("open: {e}");
return exit_code(&e);
}
};
let code = match db.write_with_wait(wait) {
Ok(mut guard) => match guard.insert_node("Person", "busy-probe", vec![]) {
Ok(()) => 0,
Err(e) => {
eprintln!("insert: {e}");
exit_code(&e)
}
},
Err(e) => {
eprintln!("write_with_wait: {e}");
exit_code(&e)
}
};
code
}
fn run_snapshot(dir: &Path) -> i32 {
match GraphDb::open(dir) {
Ok(mut db) => match db.snapshot() {
Ok(()) => 0,
Err(e) => {
eprintln!("snapshot: {e}");
exit_code(&e)
}
},
Err(e) => {
eprintln!("open: {e}");
exit_code(&e)
}
}
}
fn run_snapshot_shared(dir: &Path) -> i32 {
let db = match SharedDb::open(dir) {
Ok(db) => db,
Err(e) => {
eprintln!("open: {e}");
return exit_code(&e);
}
};
let code = match db.write().snapshot() {
Ok(()) => 0,
Err(e) => {
eprintln!("snapshot: {e}");
exit_code(&e)
}
};
code
}