pub mod apply;
pub mod bridge;
pub mod discovery;
pub mod ledger;
#[cfg(test)]
mod live_tests;
pub mod mapping;
pub mod palace_io;
#[cfg(test)]
mod r5_critic_tests;
pub mod report;
pub mod retract;
#[cfg(test)]
mod safety_tests;
pub mod screen;
#[cfg(test)]
mod tests;
use std::path::{Path, PathBuf};
use std::time::Duration;
use colored::Colorize;
use apply::{execute, plan_store, HandleSink, PalaceSink, PalaceView, StoreCounts};
use bridge::{export_store, resolve_python, CommandRunner, KuzuExport, SystemRunner};
use discovery::{discover, resolve_from, DiscoveredStore};
use ledger::Ledger;
use palace_io::{
check_env_override, open_palace_for_write, open_snapshot_view, store_is_live, target_palace,
DaemonProbe, SystemDaemonProbe,
};
use report::{print_report, print_totals};
pub const DEFAULT_MAX_DEPTH: usize = 5;
pub const DEFAULT_BRIDGE_TIMEOUT_SECS: u64 = 1800;
#[derive(Debug, thiserror::Error)]
pub enum KuzuImportError {
#[error("kuzu-memory interpreter not found: {0}")]
InterpreterNotFound(String),
#[error("export bridge exited with status {code:?}: {stderr}")]
BridgeFailed { code: Option<i32>, stderr: String },
#[error("export bridge did not finish and was {0}; raise --bridge-timeout-secs")]
BridgeTimedOut(String),
#[error("export output is malformed: {0}")]
MalformedExport(String),
#[error("store lacks required Memory columns: {}", .0.join(", "))]
SchemaColumnsMissing(Vec<String>),
#[error("{} is not a kuzu-memory store (expected a .kuzu-memory dir or memories.db)", .0.display())]
NotAStore(PathBuf),
#[error("could not resolve a palace: {0}")]
PalaceResolve(String),
#[error(
"TRUSTY_MEMORY_PALACE is set ({0:?}), and a --discover/--root walk would send every \
store into that one palace; unset it for this run \
(`env -u TRUSTY_MEMORY_PALACE trusty-memory import kuzu ...`) or import one store with --from"
)]
EnvPalaceWithWalk(String),
#[error(
"the trusty-memory daemon is running ({0}); stop the trusty-memory daemon first \
(`trusty-memory stop`), then re-run the import (a --dry-run needs no stop)"
)]
DaemonRunning(String),
#[error(
"palace '{0}' is locked by another process โ a trusty-memory daemon (`trusty-memory \
stop`) or another trusty-memory command such as a second import; let it finish or stop \
it, then re-run"
)]
PalaceLocked(String),
#[error(
"palace '{0}' has unreadable drawer rows, so what it already holds is unknown; \
refusing to import into it until it is repaired"
)]
DrawersUnreadable(String),
#[error("palace write failed: {0}")]
Palace(String),
#[error("I/O: {0}")]
Io(String),
}
impl KuzuImportError {
pub fn kind(&self) -> &'static str {
match self {
Self::InterpreterNotFound(_) => "interpreter_not_found",
Self::BridgeFailed { .. } => "bridge_failed",
Self::BridgeTimedOut(_) => "bridge_timed_out",
Self::MalformedExport(_) => "malformed_export",
Self::SchemaColumnsMissing(_) => "schema_columns_missing",
Self::NotAStore(_) => "not_a_store",
Self::PalaceResolve(_) => "palace_resolve",
Self::EnvPalaceWithWalk(_) => "env_palace_with_walk",
Self::DaemonRunning(_) => "daemon_running",
Self::PalaceLocked(_) => "palace_locked",
Self::DrawersUnreadable(_) => "drawers_unreadable",
Self::Palace(_) => "palace",
Self::Io(_) => "io",
}
}
}
#[derive(Debug, Clone, clap::Subcommand)]
pub enum ImportSource {
Kuzu(KuzuImportArgs),
}
#[derive(Debug, Clone, clap::Args)]
pub struct KuzuImportArgs {
#[arg(long)]
pub discover: bool,
#[arg(long, value_name = "PATH")]
pub root: Vec<PathBuf>,
#[arg(long, value_name = "PATH", conflicts_with_all = ["discover", "root"])]
pub from: Option<PathBuf>,
#[arg(long, value_name = "NAME", requires = "from")]
pub palace: Option<String>,
#[arg(long)]
pub dry_run: bool,
#[arg(long)]
pub update: bool,
#[arg(long, value_name = "N", default_value_t = DEFAULT_MAX_DEPTH)]
pub max_depth: usize,
#[arg(long, value_name = "PATH")]
pub python: Option<PathBuf>,
#[arg(long, value_name = "SECS", default_value_t = DEFAULT_BRIDGE_TIMEOUT_SECS)]
pub bridge_timeout_secs: u64,
}
impl Default for KuzuImportArgs {
fn default() -> Self {
Self {
discover: false,
root: Vec::new(),
from: None,
palace: None,
dry_run: false,
update: false,
max_depth: DEFAULT_MAX_DEPTH,
python: None,
bridge_timeout_secs: DEFAULT_BRIDGE_TIMEOUT_SECS,
}
}
}
#[derive(Debug)]
pub enum StoreStatus {
Imported,
WouldImport,
UpToDate,
Empty,
Partial,
Failed(KuzuImportError),
}
#[derive(Debug)]
pub struct StoreReport {
pub store: PathBuf,
pub palace: Option<String>,
pub source: Option<&'static str>,
pub counts: StoreCounts,
pub status: StoreStatus,
pub flush_error: Option<String>,
}
impl StoreReport {
pub fn is_bad(&self) -> bool {
matches!(self.status, StoreStatus::Failed(_) | StoreStatus::Partial)
|| self.flush_error.is_some()
}
}
pub enum Target<'a> {
DryRun(&'a dyn PalaceView),
Write(&'a dyn PalaceSink),
}
pub struct ImportEnv<'a> {
pub runner: &'a dyn CommandRunner,
pub daemon: &'a dyn DaemonProbe,
pub python: &'a Path,
pub data_root: &'a Path,
pub env_palace: Option<String>,
}
pub async fn handle_import(source: ImportSource) -> anyhow::Result<()> {
match source {
ImportSource::Kuzu(args) => handle_import_kuzu(args).await,
}
}
pub async fn handle_import_kuzu(args: KuzuImportArgs) -> anyhow::Result<()> {
let stores = select_stores(&args)?;
if stores.is_empty() {
println!("No kuzu-memory stores found.");
return Ok(());
}
let python = resolve_python(args.python.as_deref(), std::env::var_os("PATH").as_deref())?;
let data_dir = trusty_common::resolve_data_dir("trusty-memory")?;
let data_root = crate::resolve_palace_registry_dir(data_dir);
let runner = SystemRunner::new(Duration::from_secs(args.bridge_timeout_secs));
let env = ImportEnv {
runner: &runner,
daemon: &SystemDaemonProbe,
python: &python,
data_root: &data_root,
env_palace: trusty_common::palace_id::palace_override_from_env(),
};
let reports = run_import(&args, &stores, &env).await?;
print_totals(&reports, args.update);
let bad = reports.iter().filter(|r| r.is_bad()).count();
if bad > 0 {
anyhow::bail!(
"{bad} of {} store(s) did not import completely",
reports.len()
);
}
Ok(())
}
pub async fn run_import(
args: &KuzuImportArgs,
stores: &[DiscoveredStore],
env: &ImportEnv<'_>,
) -> Result<Vec<StoreReport>, KuzuImportError> {
check_env_override(args.from.is_none(), env.env_palace.as_deref())?;
if args.dry_run {
println!("{} Dry run โ nothing will be written.", "ยท".dimmed());
} else if let Some(daemon) = env.daemon.live_daemon().await {
return Err(KuzuImportError::DaemonRunning(daemon));
}
let mut reports = Vec::new();
for store in stores {
let report = import_one(env, store, args).await;
print_report(&report, args.update);
reports.push(report);
}
Ok(reports)
}
fn select_stores(args: &KuzuImportArgs) -> anyhow::Result<Vec<DiscoveredStore>> {
if let Some(from) = &args.from {
return Ok(vec![resolve_from(from)?]);
}
if !args.discover && args.root.is_empty() {
anyhow::bail!("name the stores: --discover, --root <path>, or --from <path>");
}
let mut roots = args.root.clone();
if args.discover {
roots.insert(
0,
dirs::home_dir().ok_or_else(|| anyhow::anyhow!("no home directory"))?,
);
}
let found = discover(&roots, args.max_depth);
for s in &found.skipped {
eprintln!(
"{} skipped {}: {}",
"ยท".dimmed(),
s.path.display(),
s.reason
);
}
Ok(found.stores)
}
async fn import_one(
env: &ImportEnv<'_>,
store: &DiscoveredStore,
args: &KuzuImportArgs,
) -> StoreReport {
let mut report = StoreReport {
store: store.dir.clone(),
palace: None,
source: None,
counts: StoreCounts::default(),
status: StoreStatus::Empty,
flush_error: None,
};
let fetched = fetch(env.runner, env.python, store).and_then(|export| {
let target = target_palace(store, args.palace.as_deref())?;
Ok((export, target))
});
let (export, target) = match fetched {
Ok(v) => v,
Err(e) => {
report.status = StoreStatus::Failed(e);
return report;
}
};
report.palace = Some(target.palace.clone());
report.source = Some(target.source);
if is_empty(&export) {
report.counts.unsupported_edges = export.other_edges.clone();
report.counts.unsupported_edges.retain(|_, n| *n > 0);
report.counts.unsupported_error = export.other_edges_error.clone();
return report;
}
let store_tag = store.dir.to_string_lossy();
let counts = if args.dry_run {
match open_snapshot_view(env.data_root, &target.palace) {
Ok(view) => run_plan(&export, &store_tag, Target::DryRun(&view), args.update).await,
Err(e) => {
report.status = StoreStatus::Failed(e);
return report;
}
}
} else {
match open_palace_for_write(env.data_root, &target.palace) {
Ok(handle) => {
let sink = HandleSink { handle };
let counts = run_plan(&export, &store_tag, Target::Write(&sink), args.update).await;
if let Err(e) = sink.handle.flush() {
report.flush_error = Some(format!("{e:#}"));
}
counts
}
Err(e) => {
report.status = StoreStatus::Failed(e);
return report;
}
}
};
report.counts = counts;
report.status = status_for(&report.counts, args.dry_run);
report
}
pub fn fetch(
runner: &dyn CommandRunner,
python: &Path,
store: &DiscoveredStore,
) -> Result<KuzuExport, KuzuImportError> {
export_store(runner, python, &store.db)
}
pub fn is_empty(export: &KuzuExport) -> bool {
export.memories.is_empty() && export.entities.is_empty() && export.edge_count() == 0
}
pub async fn run_plan(
export: &KuzuExport,
store: &str,
target: Target<'_>,
update: bool,
) -> StoreCounts {
match target {
Target::DryRun(view) => {
let ledger = Ledger::from_drawers(&view.drawers());
let plan = plan_store(export, store, &ledger, &store_is_live);
execute(plan, view, None, update).await
}
Target::Write(sink) => {
let ledger = Ledger::from_drawers(&sink.drawers());
let plan = plan_store(export, store, &ledger, &store_is_live);
execute(plan, sink, Some(sink), update).await
}
}
}
pub fn status_for(c: &StoreCounts, dry_run: bool) -> StoreStatus {
if c.failed_writes > 0 {
StoreStatus::Partial
} else if c.new_memories + c.updated + c.new_triples + c.retracted.len() == 0 {
StoreStatus::UpToDate
} else if dry_run {
StoreStatus::WouldImport
} else {
StoreStatus::Imported
}
}