use std::path::PathBuf;
use clap::{Args, Parser, Subcommand};
use fathomdb::{
CheckIntegrityOpts, CorruptionLocator, DumpProfileReport, DumpRowCountsReport,
DumpSchemaReport, Engine, EngineError, EngineOpenError, ExciseReport, Finding, IntegrityReport,
MeanRecomputeReport, RebuildKind, RebuildReport, SafeExportArtifact, SchemaObject, Section,
TraceReport, TruncateWalReport, TruncateWalStatus, VerifyEmbedderReport, VerifyEmbedderStatus,
};
use serde_json::{json, Value};
pub mod exit_code {
pub const OK: i32 = 0;
pub const RECOVERY_ACCEPTED_LOSS: i32 = 64;
pub const DOCTOR_FOUND_ISSUES: i32 = 65;
pub const EXPORT_FAILURE: i32 = 66;
pub const UNRECOVERABLE: i32 = 70;
pub const LOCK_HELD: i32 = 71;
}
#[derive(Debug, Parser)]
#[command(name = "fathomdb", version, about = "FathomDB operator CLI", long_about = None)]
pub struct Cli {
#[command(subcommand)]
pub command: Command,
}
#[derive(Debug, Subcommand)]
pub enum Command {
Recover(RecoverArgs),
Doctor(DoctorArgs),
}
#[derive(Debug, Args)]
pub struct DoctorArgs {
#[command(subcommand)]
pub command: DoctorCommand,
}
#[derive(Debug, Args)]
pub struct RecoverArgs {
#[arg(long)]
pub accept_data_loss: bool,
#[arg(long)]
pub truncate_wal: bool,
#[arg(long)]
pub rebuild_vec0: bool,
#[arg(long)]
pub rebuild_projections: bool,
#[arg(long)]
pub excise_source: Option<String>,
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Subcommand)]
pub enum DoctorCommand {
CheckIntegrity(CheckIntegrityArgs),
SafeExport(SafeExportArgs),
VerifyEmbedder(VerifyEmbedderArgs),
Trace(TraceArgs),
DumpSchema(SimpleDoctorArgs),
DumpRowCounts(SimpleDoctorArgs),
DumpProfile(SimpleDoctorArgs),
WarmCache(WarmCacheArgs),
RecomputeMean(SimpleDoctorArgs),
}
#[derive(Debug, Args)]
pub struct WarmCacheArgs {
#[arg(long)]
pub json: bool,
}
#[derive(Debug, Args)]
pub struct SimpleDoctorArgs {
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Args)]
pub struct CheckIntegrityArgs {
#[arg(long)]
pub quick: bool,
#[arg(long)]
pub full: bool,
#[arg(long = "round-trip")]
pub round_trip: bool,
#[arg(long)]
pub pretty: bool,
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Args)]
pub struct SafeExportArgs {
pub out: PathBuf,
#[arg(long)]
pub manifest: Option<PathBuf>,
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Args)]
pub struct VerifyEmbedderArgs {
#[arg(long)]
pub identity: String,
#[arg(long)]
pub dimension: u32,
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Args)]
pub struct TraceArgs {
#[arg(long = "source-ref")]
pub source_ref: String,
#[arg(long)]
pub json: bool,
pub db_path: PathBuf,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CliOutcome {
Clean,
Findings,
ExportFailure,
RecoveryAcceptedLoss,
LockHeld,
Unrecoverable,
}
#[must_use]
pub fn outcome_to_exit_code(outcome: CliOutcome) -> i32 {
match outcome {
CliOutcome::Clean => exit_code::OK,
CliOutcome::Findings => exit_code::DOCTOR_FOUND_ISSUES,
CliOutcome::ExportFailure => exit_code::EXPORT_FAILURE,
CliOutcome::RecoveryAcceptedLoss => exit_code::RECOVERY_ACCEPTED_LOSS,
CliOutcome::LockHeld => exit_code::LOCK_HELD,
CliOutcome::Unrecoverable => exit_code::UNRECOVERABLE,
}
}
#[must_use]
pub fn engine_error_to_outcome(err: &EngineError) -> CliOutcome {
match err {
EngineError::Closing => CliOutcome::LockHeld,
_ => CliOutcome::Unrecoverable,
}
}
#[must_use]
pub fn engine_open_error_to_outcome(err: &EngineOpenError) -> CliOutcome {
match err {
EngineOpenError::DatabaseLocked { .. } => CliOutcome::LockHeld,
_ => CliOutcome::Unrecoverable,
}
}
#[must_use]
pub fn run(cli: Cli) -> i32 {
match cli.command {
Command::Recover(args) => run_recover(args),
Command::Doctor(d) => run_doctor(d.command),
}
}
fn run_recover(args: RecoverArgs) -> i32 {
if !args.accept_data_loss {
println!(
r#"{{"status":"refused","verb":"recover","code":"E_RECOVER_REQUIRES_ACCEPT_DATA_LOSS"}}"#
);
return exit_code::UNRECOVERABLE;
}
if args.rebuild_projections {
return wire_recover(&args.db_path, "rebuild-projections", |e| {
e.rebuild_projections().map(|r| rebuild_report_json("rebuild-projections", &r))
});
}
if args.rebuild_vec0 {
return wire_recover(&args.db_path, "rebuild-vec0", |e| {
e.rebuild_vec0().map(|r| rebuild_report_json("rebuild-vec0", &r))
});
}
if let Some(source_id) = args.excise_source.as_deref() {
return wire_recover(&args.db_path, "excise-source", |e| {
e.excise_source(source_id).map(|r| excise_report_json(&r))
});
}
if args.truncate_wal {
return wire_recover(&args.db_path, "truncate-wal", |e| {
e.truncate_wal().map(|r| truncate_wal_report_json(&r))
});
}
println!(r#"{{"status":"not_implemented","verb":"recover"}}"#);
exit_code::UNRECOVERABLE
}
fn run_doctor(cmd: DoctorCommand) -> i32 {
match cmd {
DoctorCommand::CheckIntegrity(args) => {
let opts = CheckIntegrityOpts {
quick: args.quick,
full: args.full,
round_trip: args.round_trip,
};
run_doctor_verb(&args.db_path, "check-integrity", |e| {
e.check_integrity(opts).map(|r| integrity_report_outcome(&r))
})
}
DoctorCommand::SafeExport(args) => {
let manifest = args.manifest.clone().unwrap_or_else(|| {
let mut p = args.out.clone();
let name = p
.file_name()
.map(|s| s.to_string_lossy().into_owned())
.unwrap_or_else(|| "export".to_string());
p.set_file_name(format!("{name}.manifest.json"));
p
});
run_doctor_verb_with_error_outcome(
&args.db_path,
"safe-export",
CliOutcome::ExportFailure,
|e| {
e.safe_export(&args.out, &manifest)
.map(|r| (safe_export_json(&r), CliOutcome::Clean))
},
)
}
DoctorCommand::Trace(args) => run_doctor_verb(&args.db_path, "trace", |e| {
e.trace_source_ref(&args.source_ref).map(|r| (trace_report_json(&r), CliOutcome::Clean))
}),
DoctorCommand::VerifyEmbedder(args) => {
let identity = args.identity.clone();
let dimension = args.dimension;
run_doctor_verb(&args.db_path, "verify-embedder", |e| {
e.verify_embedder(&identity, dimension)
.map(|r| (verify_embedder_report_json(&r), CliOutcome::Clean))
})
}
DoctorCommand::DumpSchema(args) => run_doctor_verb(&args.db_path, "dump-schema", |e| {
e.dump_schema().map(|r| (dump_schema_report_json(&r), CliOutcome::Clean))
}),
DoctorCommand::DumpRowCounts(args) => {
run_doctor_verb(&args.db_path, "dump-row-counts", |e| {
e.dump_row_counts().map(|r| (dump_row_counts_report_json(&r), CliOutcome::Clean))
})
}
DoctorCommand::DumpProfile(args) => run_doctor_verb(&args.db_path, "dump-profile", |e| {
e.dump_profile().map(|r| (dump_profile_report_json(&r), CliOutcome::Clean))
}),
DoctorCommand::WarmCache(args) => run_doctor_warm_cache(args),
DoctorCommand::RecomputeMean(args) => {
run_doctor_verb(&args.db_path, "recompute-mean", |e| {
e.recompute_mean().map(|r| (recompute_mean_report_json(&r), CliOutcome::Clean))
})
}
}
}
fn run_doctor_warm_cache(args: WarmCacheArgs) -> i32 {
#[cfg(feature = "default-embedder")]
{
match fathomdb_embedder::loader::load_pinned_default_embedder() {
Ok(weights) => {
if args.json {
let payload = json!({
"verb": "warm-cache",
"status": "ok",
"config_json": weights.config_json_path.to_string_lossy(),
"tokenizer_json": weights.tokenizer_json_path.to_string_lossy(),
"model_safetensors": weights.model_safetensors_path.to_string_lossy(),
"bytes_downloaded": weights.bytes_downloaded,
"events": weights
.events
.iter()
.map(warm_cache_event_json)
.collect::<Vec<_>>(),
});
println!("{payload}");
} else {
let kind = if weights.bytes_downloaded > 0 { "cold" } else { "warm" };
println!("warm-cache: ok ({kind})");
println!(" config.json: {}", weights.config_json_path.display());
println!(" tokenizer.json: {}", weights.tokenizer_json_path.display());
println!(" model.safetensors: {}", weights.model_safetensors_path.display());
println!(" bytes downloaded: {}", weights.bytes_downloaded);
println!(" events: {}", weights.events.len());
}
exit_code::OK
}
Err(err) => {
if args.json {
let payload = json!({
"verb": "warm-cache",
"status": "error",
"code": "EmbedderLoadError",
"detail": err.to_string(),
});
println!("{payload}");
} else {
eprintln!("warm-cache: error: {err}");
}
exit_code::UNRECOVERABLE
}
}
}
#[cfg(not(feature = "default-embedder"))]
{
let detail = "fathomdb CLI was built without the `default-embedder` feature; rebuild with --features default-embedder";
if args.json {
let payload = json!({
"verb": "warm-cache",
"status": "error",
"code": "DefaultEmbedderFeatureDisabled",
"detail": detail,
});
println!("{payload}");
} else {
eprintln!("warm-cache: error: {detail}");
}
exit_code::UNRECOVERABLE
}
}
#[cfg(feature = "default-embedder")]
fn warm_cache_event_json(ev: &fathomdb_embedder::EmbedderEvent) -> Value {
use fathomdb_embedder::EmbedderEvent;
match ev {
EmbedderEvent::DefaultEmbedderDownload {
file,
url,
bytes,
sha256,
cache_path,
duration_ms,
} => json!({
"kind": "download",
"file": file,
"url": url,
"bytes": bytes,
"sha256": sha256,
"cache_path": cache_path.to_string_lossy(),
"duration_ms": duration_ms,
}),
EmbedderEvent::DefaultEmbedderCacheHit { file, sha256, cache_path } => json!({
"kind": "cache_hit",
"file": file,
"sha256": sha256,
"cache_path": cache_path.to_string_lossy(),
}),
EmbedderEvent::MeanVecPinned { dim, doc_count } => json!({
"kind": "mean_vec_pinned",
"dim": dim,
"doc_count": doc_count,
}),
EmbedderEvent::MeanVecRecomputed { dim, doc_count, trigger } => json!({
"kind": "mean_vec_recomputed",
"dim": dim,
"doc_count": doc_count,
"trigger": trigger.as_str(),
}),
}
}
fn run_doctor_verb<F>(db_path: &std::path::Path, verb: &str, f: F) -> i32
where
F: FnOnce(&Engine) -> Result<(Value, CliOutcome), EngineError>,
{
run_doctor_verb_inner(db_path, verb, None, f)
}
fn run_doctor_verb_with_error_outcome<F>(
db_path: &std::path::Path,
verb: &str,
error_outcome: CliOutcome,
f: F,
) -> i32
where
F: FnOnce(&Engine) -> Result<(Value, CliOutcome), EngineError>,
{
run_doctor_verb_inner(db_path, verb, Some(error_outcome), f)
}
fn run_doctor_verb_inner<F>(
db_path: &std::path::Path,
verb: &str,
error_outcome: Option<CliOutcome>,
f: F,
) -> i32
where
F: FnOnce(&Engine) -> Result<(Value, CliOutcome), EngineError>,
{
let opened = match Engine::open(db_path.to_path_buf()) {
Ok(o) => o,
Err(err) => return emit_engine_open_error(verb, &err),
};
match f(&opened.engine) {
Ok((value, outcome)) => {
println!("{value}");
outcome_to_exit_code(outcome)
}
Err(err) => match error_outcome {
Some(outcome) => emit_engine_error_with_outcome(verb, &err, outcome),
None => emit_engine_error(verb, &err),
},
}
}
fn wire_recover<F>(db_path: &std::path::Path, sub_verb: &str, f: F) -> i32
where
F: FnOnce(&Engine) -> Result<Value, EngineError>,
{
let opened = match Engine::open(db_path.to_path_buf()) {
Ok(o) => o,
Err(err) => return emit_engine_open_error(sub_verb, &err),
};
match f(&opened.engine) {
Ok(value) => {
println!("{value}");
outcome_to_exit_code(CliOutcome::RecoveryAcceptedLoss)
}
Err(err) => emit_engine_error(sub_verb, &err),
}
}
fn emit_engine_error(verb: &str, err: &EngineError) -> i32 {
emit_engine_error_with_outcome(verb, err, engine_error_to_outcome(err))
}
fn emit_engine_error_with_outcome(verb: &str, err: &EngineError, outcome: CliOutcome) -> i32 {
let payload = json!({
"status": "error",
"verb": verb,
"code": engine_error_code(err),
"detail": err.to_string(),
});
println!("{payload}");
outcome_to_exit_code(outcome)
}
fn emit_engine_open_error(verb: &str, err: &EngineOpenError) -> i32 {
let outcome = engine_open_error_to_outcome(err);
let payload = json!({
"status": "error",
"verb": verb,
"code": engine_open_error_code(err),
"detail": err.to_string(),
});
println!("{payload}");
outcome_to_exit_code(outcome)
}
fn engine_error_code(err: &EngineError) -> &'static str {
match err {
EngineError::Storage => "StorageError",
EngineError::Projection => "ProjectionError",
EngineError::Vector => "VectorError",
EngineError::Embedder => "EmbedderError",
EngineError::EmbedderNotConfigured => "EmbedderNotConfiguredError",
EngineError::KindNotVectorIndexed => "KindNotVectorIndexedError",
EngineError::EmbedderDimensionMismatch { .. } => "EmbedderDimensionMismatchError",
EngineError::Scheduler => "SchedulerError",
EngineError::OpStore => "OpStoreError",
EngineError::WriteValidation => "WriteValidationError",
EngineError::SchemaValidation => "SchemaValidationError",
EngineError::Overloaded => "OverloadedError",
EngineError::Closing => "ClosingError",
}
}
fn engine_open_error_code(err: &EngineOpenError) -> &'static str {
match err {
EngineOpenError::DatabaseLocked { .. } => "DatabaseLockedError",
EngineOpenError::Corruption(_) => "CorruptionError",
EngineOpenError::IncompatibleSchemaVersion { .. } => "IncompatibleSchemaVersionError",
EngineOpenError::MigrationError { .. } => "MigrationError",
EngineOpenError::EmbedderIdentityMismatch { .. } => "EmbedderIdentityMismatchError",
EngineOpenError::EmbedderDimensionMismatch { .. } => "EmbedderDimensionMismatchError",
EngineOpenError::Embedder(_) => "EmbedderError",
EngineOpenError::Io { .. } => "IoError",
}
}
fn integrity_report_outcome(report: &IntegrityReport) -> (Value, CliOutcome) {
let any_findings = matches!(report.physical, Section::Findings(_))
|| matches!(report.logical, Section::Findings(_))
|| matches!(report.semantic, Section::Findings(_));
let body = json!({
"verb": "check-integrity",
"physical": section_json(&report.physical),
"logical": section_json(&report.logical),
"semantic": section_json(&report.semantic),
});
let outcome = if any_findings { CliOutcome::Findings } else { CliOutcome::Clean };
(body, outcome)
}
fn section_json(section: &Section) -> Value {
match section {
Section::Clean => json!({ "status": "clean", "findings": [] }),
Section::Findings(findings) => json!({
"status": "findings",
"findings": findings.iter().map(finding_json).collect::<Vec<_>>(),
}),
}
}
fn finding_json(f: &Finding) -> Value {
json!({
"code": f.code,
"stage": f.stage,
"locator": locator_json(&f.locator),
"doc_anchor": f.doc_anchor,
"detail": f.detail,
})
}
fn locator_json(loc: &CorruptionLocator) -> Value {
match loc {
CorruptionLocator::FileOffset { offset } => {
json!({ "kind": "file_offset", "offset": offset })
}
CorruptionLocator::PageId { page } => json!({ "kind": "page_id", "page": page }),
CorruptionLocator::TableRow { table, rowid } => {
json!({ "kind": "table_row", "table": table, "rowid": rowid })
}
CorruptionLocator::Vec0ShadowRow { partition, rowid } => {
json!({ "kind": "vec0_shadow_row", "partition": partition, "rowid": rowid })
}
CorruptionLocator::MigrationStep { from, to } => {
json!({ "kind": "migration_step", "from": from, "to": to })
}
CorruptionLocator::OpaqueSqliteError { sqlite_extended_code } => {
json!({
"kind": "opaque_sqlite_error",
"sqlite_extended_code": sqlite_extended_code,
})
}
}
}
fn safe_export_json(a: &SafeExportArtifact) -> Value {
json!({
"verb": "safe-export",
"export_path": a.export_path.to_string_lossy(),
"manifest_path": a.manifest_path.to_string_lossy(),
"manifest_sha256": a.manifest_sha256,
})
}
fn trace_report_json(t: &TraceReport) -> Value {
json!({
"verb": "trace",
"source_ref": t.source_ref,
"events": t.events.iter().map(|e| json!({
"write_cursor": e.write_cursor,
"kind": e.kind,
"table": e.table,
})).collect::<Vec<_>>(),
})
}
fn rebuild_report_json(verb: &'static str, r: &RebuildReport) -> Value {
let kind = match r.kind {
RebuildKind::Projections => "projections",
RebuildKind::Vec0 => "vec0",
};
json!({
"verb": verb,
"kind": kind,
"rows_invalidated": r.rows_invalidated,
"rows_rebuilt": r.rows_rebuilt,
"projection_cursor_after": r.projection_cursor_after,
})
}
fn excise_report_json(r: &ExciseReport) -> Value {
json!({
"verb": "excise-source",
"source_ref": r.source_ref,
"nodes_excised": r.nodes_excised,
"edges_excised": r.edges_excised,
"projections_invalidated": r.projections_invalidated,
})
}
fn verify_embedder_report_json(r: &VerifyEmbedderReport) -> Value {
let status = match r.status {
VerifyEmbedderStatus::Match => "match",
VerifyEmbedderStatus::IdentityMismatch => "identity_mismatch",
VerifyEmbedderStatus::DimensionMismatch => "dimension_mismatch",
VerifyEmbedderStatus::BothMismatch => "both_mismatch",
};
json!({
"verb": "verify-embedder",
"stored_identity": r.stored_identity,
"stored_dimension": r.stored_dimension,
"supplied_identity": r.supplied_identity,
"supplied_dimension": r.supplied_dimension,
"status": status,
})
}
fn schema_object_json(o: &SchemaObject) -> Value {
json!({ "name": o.name, "sql": o.sql })
}
fn dump_schema_report_json(r: &DumpSchemaReport) -> Value {
json!({
"verb": "dump-schema",
"user_version": r.user_version,
"tables": r.tables.iter().map(schema_object_json).collect::<Vec<_>>(),
"indexes": r.indexes.iter().map(schema_object_json).collect::<Vec<_>>(),
})
}
fn recompute_mean_report_json(r: &MeanRecomputeReport) -> Value {
json!({
"verb": "recompute-mean",
"status": "ok",
"dim": r.dim,
"old_doc_count": r.old_doc_count,
"doc_count_requantized": r.doc_count_requantized,
"drift_cos_before": r.drift_cos_before,
"mean_was_pinned": r.mean_was_pinned,
"elapsed_ms": r.elapsed_ms,
})
}
fn dump_row_counts_report_json(r: &DumpRowCountsReport) -> Value {
json!({
"verb": "dump-row-counts",
"counts": r.counts.iter().map(|c| json!({
"name": c.name,
"rows": c.rows,
})).collect::<Vec<_>>(),
})
}
fn dump_profile_report_json(r: &DumpProfileReport) -> Value {
json!({
"verb": "dump-profile",
"embedder_identity": r.embedder_identity,
"embedder_dimension": r.embedder_dimension,
"vectorized_kinds": r.vectorized_kinds,
})
}
fn truncate_wal_report_json(r: &TruncateWalReport) -> Value {
let status = match r.status {
TruncateWalStatus::Done => "done",
TruncateWalStatus::Busy => "busy",
};
json!({
"verb": "truncate-wal",
"status": status,
"busy": r.busy,
"log_frames": r.log_frames,
"checkpointed_frames": r.checkpointed_frames,
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn outcome_mapping_covers_cli_md_exit_classes() {
assert_eq!(outcome_to_exit_code(CliOutcome::Clean), 0);
assert_eq!(outcome_to_exit_code(CliOutcome::RecoveryAcceptedLoss), 64);
assert_eq!(outcome_to_exit_code(CliOutcome::Findings), 65);
assert_eq!(outcome_to_exit_code(CliOutcome::ExportFailure), 66);
assert_eq!(outcome_to_exit_code(CliOutcome::Unrecoverable), 70);
assert_eq!(outcome_to_exit_code(CliOutcome::LockHeld), 71);
}
#[test]
fn engine_error_storage_maps_to_unrecoverable() {
assert_eq!(engine_error_to_outcome(&EngineError::Storage), CliOutcome::Unrecoverable);
assert_eq!(engine_error_to_outcome(&EngineError::Closing), CliOutcome::LockHeld);
}
#[test]
fn engine_open_database_locked_maps_to_lock_held() {
let err = EngineOpenError::DatabaseLocked { holder_pid: Some(1234) };
assert_eq!(engine_open_error_to_outcome(&err), CliOutcome::LockHeld);
}
}