use std::path::PathBuf;
use std::sync::Arc;
use anyhow::{Context, Result};
use clap::{Parser, Subcommand};
use khive_runtime::{BackendId, KhiveConfig, KhiveRuntime, RuntimeConfig};
use kkernel::{
code_ingest,
coordinator::{BackendRegistry, SubstrateCoordinator, SubstrateCoordinatorService},
engine, exec, git_ingest, kg, pack_introspect, reindex, sync, vector,
};
#[derive(Parser, Debug)]
#[command(
name = "kkernel",
version,
about = "khive kernel — admin/management Rust binary"
)]
struct Args {
#[arg(long, env = "KHIVE_LOG", default_value = "warn", global = true)]
log: String,
#[arg(short = 'e', long = "exec", value_name = "OPS")]
exec: Option<String>,
#[command(subcommand)]
command: Option<Command>,
}
#[derive(Subcommand, Debug)]
enum Command {
Sync(SyncArgs),
#[command(subcommand)]
Pack(PackCommand),
#[command(subcommand)]
Kg(kg::KgCommand),
#[command(subcommand)]
Db(DbCommand),
#[command(subcommand)]
Engine(engine::EngineCommand),
#[command(subcommand)]
Vector(vector::VectorCommand),
Reindex(reindex::ReindexArgs),
Exec(exec::ExecArgs),
Mcp(khive_mcp::args::Args),
#[command(subcommand)]
Backend(BackendCommand),
GitIngest(git_ingest::GitIngestArgs),
CodeIngest(code_ingest::CodeIngestArgs),
}
#[derive(Subcommand, Debug)]
enum DbCommand {
Migrate(DbMigrateArgs),
Check(DbCheckArgs),
}
#[derive(clap::Parser, Debug)]
struct DbMigrateArgs {
#[arg(long)]
db: Option<PathBuf>,
#[arg(long)]
backend: Option<String>,
#[arg(long)]
dry_run: bool,
#[arg(long)]
check: bool,
#[arg(long)]
human: bool,
}
#[derive(clap::Parser, Debug)]
struct DbCheckArgs {
#[arg(long)]
db: Option<PathBuf>,
#[arg(long)]
strict: bool,
#[arg(long)]
human: bool,
}
#[derive(Parser, Debug)]
struct SyncArgs {
#[arg(long, default_value = ".")]
repo: PathBuf,
#[arg(long)]
db: PathBuf,
#[arg(long, default_value = "local")]
namespace: String,
}
#[derive(Subcommand, Debug)]
enum PackCommand {
List {
#[arg(long)]
human: bool,
},
Handler {
name: String,
#[arg(long)]
human: bool,
},
}
#[derive(Subcommand, Debug)]
enum BackendCommand {
List {
#[arg(long)]
human: bool,
},
Info {
name: String,
#[arg(long)]
human: bool,
},
}
fn load_khive_dotenv() {
let Some(home) = std::env::var_os("HOME") else {
return;
};
let path = std::path::Path::new(&home).join(".khive/.env");
match dotenvy::from_path(&path) {
Ok(()) => {}
Err(e) if e.not_found() => {}
Err(e) => eprintln!("warning: failed to load {}: {e}", path.display()),
}
}
#[tokio::main]
async fn main() -> Result<()> {
load_khive_dotenv();
let args = Args::parse();
init_tracing(&args.log);
let command = resolve_command(args.exec, args.command);
match command {
Command::Sync(s) => cmd_sync(s).await,
Command::Pack(p) => cmd_pack(p),
Command::Kg(k) => kg::run_kg(k).await,
Command::Db(d) => cmd_db(d).await,
Command::Engine(e) => engine::run_engine(e).await,
Command::Vector(v) => vector::run_vector(v),
Command::Reindex(r) => reindex::run_reindex(r).await,
Command::Exec(e) => exec::run_exec(e).await,
Command::Mcp(a) => {
let transport_registry = khive_mcp::transport::TransportRegistry::with_builtins();
let db_path_hint = khive_mcp::serve::config_discovery_db_anchor(a.db.as_deref());
let khive_cfg =
KhiveConfig::load_with_home_fallback(a.config.as_deref(), db_path_hint.as_deref())
.unwrap_or_default()
.unwrap_or_default();
if khive_cfg.backends.len() <= 1 {
khive_mcp::serve::run(a, &transport_registry).await
} else {
let (cli_ns_explicit, cli_ns) = khive_mcp::args::resolve_cli_namespace(&a)
.map_err(|e| anyhow::anyhow!("{e}"))?;
let base_cfg = khive_mcp::serve::resolve_runtime_config(
khive_mcp::serve::RuntimeConfigInputs {
db: a.db.as_deref(),
config: a.config.as_deref(),
namespace: cli_ns,
namespace_explicit: cli_ns_explicit,
actor_explicit: cli_ns_explicit,
no_embed: a.no_embed,
packs: if a.pack.is_empty() {
None
} else {
Some(a.pack.clone())
},
brain_profile: a.brain_profile.clone(),
},
)?;
#[cfg(unix)]
let boot_guard = if a.daemon {
Some(khive_runtime::daemon::acquire_daemon_boot_guard()?)
} else {
khive_runtime::daemon::acquire_recovery_lock()
};
#[cfg(not(unix))]
let boot_guard: Option<std::fs::File> = None;
let (server, schedule_rt) = build_multi_backend_server_with_coordinator(
base_cfg,
&khive_cfg,
a.db.as_deref(),
)?;
khive_mcp::serve::serve_server(
server,
&a,
&transport_registry,
boot_guard,
schedule_rt,
)
.await
}
}
Command::Backend(b) => cmd_backend(b),
Command::GitIngest(a) => git_ingest::run_git_ingest(a).await,
Command::CodeIngest(a) => code_ingest::run_code_ingest(a).await,
}
}
#[derive(Debug, PartialEq, Eq)]
enum ResolveCommandError {
Missing,
Conflict,
}
fn resolve_command_result(
exec: Option<String>,
command: Option<Command>,
) -> Result<Command, ResolveCommandError> {
match (exec, command) {
(Some(ops), None) => Ok(Command::Exec(exec::ExecArgs::parse_from([
"exec", "--", &ops,
]))),
(None, Some(cmd)) => Ok(cmd),
(None, None) => Err(ResolveCommandError::Missing),
(Some(_), Some(_)) => Err(ResolveCommandError::Conflict),
}
}
fn resolve_command(exec: Option<String>, command: Option<Command>) -> Command {
use clap::{error::ErrorKind, CommandFactory};
match resolve_command_result(exec, command) {
Ok(cmd) => cmd,
Err(ResolveCommandError::Missing) => Args::command()
.error(
ErrorKind::MissingRequiredArgument,
"either provide -e/--exec <OPS> or a subcommand",
)
.exit(),
Err(ResolveCommandError::Conflict) => Args::command()
.error(
ErrorKind::ArgumentConflict,
"the argument '-e/--exec <OPS>' cannot be used with a subcommand",
)
.exit(),
}
}
fn build_multi_backend_server_with_coordinator(
base_cfg: RuntimeConfig,
khive_cfg: &KhiveConfig,
cli_db_override: Option<&str>,
) -> Result<(khive_mcp::server::KhiveMcpServer, Option<KhiveRuntime>)> {
let multi =
khive_mcp::serve::build_registry_for_multi_backend(base_cfg, khive_cfg, cli_db_override)?;
let schedule_rt = multi
.per_pack_runtimes
.get("schedule")
.map(|rt| (**rt).clone());
let mut backend_reg = BackendRegistry::new();
for (pack_name, rt) in &multi.per_pack_runtimes {
let backend_name = khive_cfg
.packs
.get(pack_name.as_str())
.map(|pc| pc.backend.as_str())
.unwrap_or(BackendId::MAIN);
let backend_id = BackendId::new(backend_name);
backend_reg.register(backend_id, Arc::clone(rt));
}
let note_kinds: std::collections::HashSet<String> = multi
.registry
.all_note_kinds()
.into_iter()
.map(str::to_string)
.collect();
let coord =
SubstrateCoordinatorService::new(SubstrateCoordinator::new(backend_reg), note_kinds);
let server = khive_mcp::serve::build_server_from_multi_backend_registry(
multi,
khive_cfg,
Some(Arc::new(coord) as Arc<dyn khive_mcp::coordinator::CoordinatorService>),
);
Ok((server, schedule_rt))
}
async fn cmd_db(cmd: DbCommand) -> Result<()> {
match cmd {
DbCommand::Migrate(args) => cmd_db_migrate(args).await,
DbCommand::Check(args) => cmd_db_check(args).await,
}
}
async fn cmd_db_migrate(args: DbMigrateArgs) -> Result<()> {
let mut cfg = RuntimeConfig::default();
if let Some(ref db) = args.db {
cfg.db_path = Some(db.clone());
}
if args.dry_run || args.check {
return cmd_db_check(DbCheckArgs {
db: args.db,
strict: args.check,
human: args.human,
})
.await;
}
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
let latest = khive_db::MIGRATIONS.len() as u32;
let sql = rt.sql();
let mut reader = sql
.reader()
.await
.context("open SQL reader after migration")?;
use khive_storage::types::{SqlStatement, SqlValue};
let rows = reader
.query_all(SqlStatement {
sql: "SELECT COALESCE(MAX(version), 0) FROM _schema_migrations".into(),
params: vec![],
label: Some("db_migrate_version".into()),
})
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
let applied: u32 = rows
.first()
.and_then(|r| match r.get("COALESCE(MAX(version), 0)") {
Some(SqlValue::Integer(v)) => Some(*v as u32),
_ => None,
})
.unwrap_or(latest);
if args.human {
println!("schema migrated: version {applied} of {latest} (current)");
} else {
let json = serde_json::json!({
"applied_version": applied,
"latest_version": latest,
"current": applied == latest,
});
println!("{}", serde_json::to_string(&json).expect("serialize"));
}
Ok(())
}
async fn cmd_db_check(args: DbCheckArgs) -> Result<()> {
let latest = khive_db::MIGRATIONS.len() as u32;
let resolved: Option<PathBuf> = match args.db {
Some(p) => Some(p),
None => std::env::var("HOME")
.ok()
.map(|h| PathBuf::from(h).join(".khive/khive.db")),
};
let current_version: u32 = match resolved {
Some(ref p) if p.exists() => {
khive_db::inspect_schema_version(p).map_err(|e| anyhow::anyhow!("{e}"))?
}
_ => 0,
};
let is_current = current_version == latest;
let ahead = current_version > latest;
if args.human {
let state = if ahead {
"ahead — predates the consolidated baseline (ADR-015) or written by a newer build; recreate it"
} else if is_current {
"current"
} else {
"behind — run: kkernel db migrate"
};
println!("main: V{current_version} ({state})");
} else {
let json = serde_json::json!({
"current_version": current_version,
"latest_version": latest,
"current": is_current,
"ahead": ahead,
"pending": latest.saturating_sub(current_version),
});
println!("{}", serde_json::to_string(&json).expect("serialize"));
}
if args.strict && !is_current {
if ahead {
anyhow::bail!(
"schema version {current_version} is ahead of the latest known migration {latest} — \
this database predates the consolidated baseline (ADR-015) or was written by a newer \
build; recreate it from the current schema"
);
}
anyhow::bail!(
"schema is behind: V{current_version} applied, V{latest} is current — \
run `kkernel db migrate` to bring the schema up to date"
);
}
Ok(())
}
fn init_tracing(level: &str) {
let filter = format!("{level},lattice_inference=error");
tracing_subscriber::fmt()
.with_writer(std::io::stderr)
.with_env_filter(filter)
.with_ansi(false)
.init();
}
async fn cmd_sync(args: SyncArgs) -> Result<()> {
let report = sync::run_sync(&args.repo, &args.db, &args.namespace)
.await
.with_context(|| {
format!(
"sync failed for repo={} db={}",
args.repo.display(),
args.db.display()
)
})?;
let json = serde_json::to_string(&report).expect("serialize SyncReport");
println!("{json}");
Ok(())
}
fn cmd_pack(cmd: PackCommand) -> Result<()> {
match cmd {
PackCommand::List { human } => {
let packs = pack_introspect::list_packs()?;
if human {
for p in &packs {
println!("# {} ({} verbs)", p.name, p.verbs.len());
if !p.requires.is_empty() {
println!(" requires: {}", p.requires.join(", "));
}
if !p.note_kinds.is_empty() {
println!(" note_kinds: {}", p.note_kinds.join(", "));
}
if !p.entity_kinds.is_empty() {
println!(" entity_kinds: {}", p.entity_kinds.join(", "));
}
for v in &p.verbs {
println!(" {:<20} {}", v.name, v.description);
}
println!();
}
} else {
let json = serde_json::to_string(&packs).expect("serialize PackInfo[]");
println!("{json}");
}
Ok(())
}
PackCommand::Handler { name, human } => {
let info = pack_introspect::pack_handler(&name)?;
let info = info.with_context(|| format!("pack {name:?} is not registered"))?;
if human {
println!("# {} ({} verbs)", info.name, info.verbs.len());
if !info.requires.is_empty() {
println!("requires: {}", info.requires.join(", "));
}
if !info.note_kinds.is_empty() {
println!("note_kinds: {}", info.note_kinds.join(", "));
}
if !info.entity_kinds.is_empty() {
println!("entity_kinds: {}", info.entity_kinds.join(", "));
}
for v in &info.verbs {
println!(" {:<20} {}", v.name, v.description);
}
} else {
let json = serde_json::to_string(&info).expect("serialize PackInfo");
println!("{json}");
}
Ok(())
}
}
}
fn cmd_backend(cmd: BackendCommand) -> Result<()> {
let default_config = RuntimeConfig::default();
let default_id = default_config.backend_id.clone();
let default_path = default_config
.db_path
.as_ref()
.map(|p| p.display().to_string())
.unwrap_or_else(|| ":memory:".to_string());
let mut registry = BackendRegistry::new();
let rt = KhiveRuntime::new(default_config).map_err(|e| anyhow::anyhow!("{e}"))?;
registry.register(default_id.clone(), std::sync::Arc::new(rt));
match cmd {
BackendCommand::List { human } => {
let ids: Vec<_> = registry.ids();
if human {
println!("Registered backends ({}):", ids.len());
for id in &ids {
let entry = registry.get(id).unwrap();
let primary_marker = if registry.primary().map(|p| p.id == *id).unwrap_or(false)
{
" [primary]"
} else {
""
};
println!(" {}{}", id.as_str(), primary_marker);
let _ = entry; }
} else {
let names: Vec<&str> = ids.iter().map(|id| id.as_str()).collect();
let json = serde_json::json!({
"backends": names,
"primary": registry.primary().map(|e| e.id.as_str()),
"count": ids.len(),
});
println!("{}", serde_json::to_string(&json).expect("serialize"));
}
Ok(())
}
BackendCommand::Info { name, human } => {
let id = BackendId::new(&name);
let entry = registry
.get(&id)
.with_context(|| format!("backend {name:?} is not registered"))?;
if human {
let is_primary = registry
.primary()
.map(|p| p.id == entry.id)
.unwrap_or(false);
println!("backend: {}", entry.id.as_str());
println!(" primary: {is_primary}");
println!(" path: {default_path}");
} else {
let json = serde_json::json!({
"name": entry.id.as_str(),
"path": default_path,
"primary": registry.primary().map(|p| p.id == entry.id).unwrap_or(false),
});
println!("{}", serde_json::to_string(&json).expect("serialize"));
}
Ok(())
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use serial_test::serial;
use tempfile::TempDir;
#[tokio::test]
async fn db_check_does_not_create_missing_file() {
let tmp = TempDir::new().expect("temp dir");
let path = tmp.path().join("missing.db");
assert!(!path.exists());
cmd_db_check(DbCheckArgs {
db: Some(path.clone()),
strict: false,
human: false,
})
.await
.expect("db check succeeds on a missing file");
assert!(!path.exists(), "db check must not create the database file");
}
#[tokio::test]
async fn db_check_does_not_mutate_existing_db() {
let tmp = TempDir::new().expect("temp dir");
let path = tmp.path().join("real.db");
cmd_db_migrate(DbMigrateArgs {
db: Some(path.clone()),
backend: None,
dry_run: false,
check: false,
human: false,
})
.await
.expect("migrate creates the database");
let before = std::fs::read(&path).expect("read db before check");
cmd_db_check(DbCheckArgs {
db: Some(path.clone()),
strict: true,
human: false,
})
.await
.expect("db check passes on a current db");
let after = std::fs::read(&path).expect("read db after check");
assert_eq!(before, after, "db check must not mutate the database");
}
#[test]
fn exec_shortcut_short_flag_parses_ops() {
let args = Args::parse_from(["kkernel", "-e", "stats()"]);
assert_eq!(args.exec.as_deref(), Some("stats()"));
assert!(args.command.is_none());
}
#[test]
fn exec_shortcut_long_flag_parses_ops() {
let args = Args::parse_from(["kkernel", "--exec", "stats()"]);
assert_eq!(args.exec.as_deref(), Some("stats()"));
assert!(args.command.is_none());
}
#[test]
fn exec_shortcut_conflicts_with_subcommand() {
let args = Args::parse_from(["kkernel", "-e", "stats()", "exec", "other()"]);
assert!(args.exec.is_some());
assert!(args.command.is_some());
let result = resolve_command_result(args.exec, args.command);
assert!(matches!(result, Err(ResolveCommandError::Conflict)));
}
#[test]
fn exec_shortcut_conflicts_with_subcommand_reverse_order() {
let bare = Args::try_parse_from(["kkernel", "pack", "list"]);
assert!(bare.is_ok(), "a bare subcommand alone must still parse");
let result = Args::try_parse_from(["kkernel", "exec", "other()", "-e", "stats()"]);
assert!(
result.is_err(),
"-e after a subcommand's own args must be rejected"
);
}
#[test]
fn resolve_command_result_missing_when_neither_given() {
let args = Args::parse_from(["kkernel"]);
let result = resolve_command_result(args.exec, args.command);
assert!(matches!(result, Err(ResolveCommandError::Missing)));
}
#[test]
fn resolve_command_result_exec_only_maps_to_exec_command() {
let args = Args::parse_from(["kkernel", "-e", "stats()"]);
let result = resolve_command_result(args.exec, args.command);
match result {
Ok(Command::Exec(e)) => assert_eq!(e.ops.as_deref(), Some("stats()")),
other => panic!("expected Ok(Command::Exec), got {other:?}"),
}
}
#[test]
fn exec_shortcut_maps_to_same_ops_as_exec_subcommand() {
let via_shortcut = match resolve_command_result(Some("stats()".into()), None) {
Ok(Command::Exec(e)) => e,
other => panic!("expected Ok(Command::Exec), got {other:?}"),
};
let via_subcommand = match Args::parse_from(["kkernel", "exec", "stats()"]).command {
Some(Command::Exec(e)) => e,
other => panic!("expected Command::Exec, got {other:?}"),
};
assert_eq!(via_shortcut.ops, via_subcommand.ops);
assert_eq!(via_shortcut.db, via_subcommand.db);
assert_eq!(via_shortcut.namespace, via_subcommand.namespace);
assert_eq!(via_shortcut.presentation, via_subcommand.presentation);
}
#[test]
fn exec_shortcut_flag_like_ops_binds_as_ops_not_as_exec_flag() {
let resolved = match resolve_command_result(Some("--pending-events".into()), None) {
Ok(Command::Exec(e)) => e,
other => panic!("expected Ok(Command::Exec), got {other:?}"),
};
assert_eq!(resolved.ops.as_deref(), Some("--pending-events"));
assert!(!resolved.pending_events);
}
#[test]
fn bare_invocation_without_exec_or_subcommand_is_not_a_valid_parse_state() {
let args = Args::parse_from(["kkernel"]);
assert!(args.exec.is_none());
assert!(args.command.is_none());
}
fn base_multi_backend_runtime_config() -> RuntimeConfig {
use khive_runtime::Namespace;
RuntimeConfig {
db_path: khive_runtime::resolve_db_anchor(None),
default_namespace: Namespace::parse("local").expect("valid namespace"),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
backend_id: BackendId::main(),
..RuntimeConfig::default()
}
}
fn single_main_backend_config(
kind: khive_runtime::BackendKind,
path: Option<PathBuf>,
) -> KhiveConfig {
KhiveConfig {
backends: vec![khive_runtime::BackendConfig {
name: "main".to_string(),
kind,
path,
cache_mb: None,
journal_mode: None,
read_only: false,
}],
..KhiveConfig::default()
}
}
#[test]
fn multi_backend_boot_paths_share_identical_wiring_surface_file_backed() {
let dir = TempDir::new().expect("temp dir");
let main_path = dir.path().join("main.db");
let khive_cfg =
single_main_backend_config(khive_runtime::BackendKind::Sqlite, Some(main_path));
let plain_server = khive_mcp::serve::build_server_multi_backend(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("plain multi-backend boot must succeed");
let (coordinator_server, _schedule_rt) = build_multi_backend_server_with_coordinator(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("kkernel coordinator-attached multi-backend boot must succeed");
let plain_surface = khive_mcp::serve::WiringSurface::capture(&plain_server);
let coordinator_surface = khive_mcp::serve::WiringSurface::capture(&coordinator_server);
assert_eq!(
plain_surface, coordinator_surface,
"the plain multi-backend boot path and kkernel's coordinator-attached \
boot path must produce an identical wiring surface for the same config"
);
assert!(
plain_surface.has_checkpoint_pool,
"file-backed main must wire a checkpoint pool on both paths"
);
}
#[test]
fn multi_backend_boot_paths_share_identical_wiring_surface_in_memory() {
let khive_cfg = single_main_backend_config(khive_runtime::BackendKind::Memory, None);
let plain_server = khive_mcp::serve::build_server_multi_backend(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("plain multi-backend boot must succeed");
let (coordinator_server, _schedule_rt) = build_multi_backend_server_with_coordinator(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("kkernel coordinator-attached multi-backend boot must succeed");
let plain_surface = khive_mcp::serve::WiringSurface::capture(&plain_server);
let coordinator_surface = khive_mcp::serve::WiringSurface::capture(&coordinator_server);
assert_eq!(
plain_surface, coordinator_surface,
"the plain multi-backend boot path and kkernel's coordinator-attached \
boot path must produce an identical wiring surface for the same config"
);
assert!(
!plain_surface.has_checkpoint_pool,
"in-memory main must never carry a checkpoint pool on either path"
);
}
#[test]
#[serial]
fn multi_backend_boot_paths_share_identical_non_default_output_format() {
struct OutputFormatEnvGuard {
prev: Option<String>,
}
impl OutputFormatEnvGuard {
fn clear() -> Self {
let prev = std::env::var("KHIVE_OUTPUT_FORMAT").ok();
std::env::remove_var("KHIVE_OUTPUT_FORMAT");
Self { prev }
}
}
impl Drop for OutputFormatEnvGuard {
fn drop(&mut self) {
match &self.prev {
Some(v) => std::env::set_var("KHIVE_OUTPUT_FORMAT", v),
None => std::env::remove_var("KHIVE_OUTPUT_FORMAT"),
}
}
}
let _env_guard = OutputFormatEnvGuard::clear();
let mut khive_cfg = single_main_backend_config(khive_runtime::BackendKind::Memory, None);
khive_cfg.runtime.default_output_format = Some(khive_runtime::OutputFormat::Table);
let plain_server = khive_mcp::serve::build_server_multi_backend(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("plain multi-backend boot must succeed");
let (coordinator_server, _schedule_rt) = build_multi_backend_server_with_coordinator(
base_multi_backend_runtime_config(),
&khive_cfg,
None,
)
.expect("kkernel coordinator-attached multi-backend boot must succeed");
let plain_surface = khive_mcp::serve::WiringSurface::capture(&plain_server);
let coordinator_surface = khive_mcp::serve::WiringSurface::capture(&coordinator_server);
assert_eq!(
plain_surface, coordinator_surface,
"the plain multi-backend boot path and kkernel's coordinator-attached \
boot path must produce an identical wiring surface for the same config"
);
assert_eq!(
plain_surface.output_format,
khive_runtime::OutputFormat::Table,
"both paths must resolve the configured non-default [runtime].default_output_format \
(Table), not silently fall back to the builtin Json default — this is the exact \
ADR-078 regression class the parity test exists to catch"
);
}
#[test]
fn coordinator_boundary_rejects_diverging_db_path() {
let args_db = "/tmp/khive-coordinator-guard-real.db";
let wrong_path = std::path::PathBuf::from("/tmp/khive-coordinator-guard-wrong.db");
let base_cfg = RuntimeConfig {
db_path: Some(wrong_path.clone()),
..base_multi_backend_runtime_config()
};
let khive_cfg = KhiveConfig::default();
let result =
build_multi_backend_server_with_coordinator(base_cfg, &khive_cfg, Some(args_db));
let err = match result {
Ok(_) => panic!(
"a resolved db_path diverging from the canonical anchor must be rejected \
at the coordinator-attached construction boundary"
),
Err(e) => e,
};
let msg = err.to_string();
let anchor =
khive_runtime::resolve_db_anchor(Some(args_db)).expect("explicit path always anchors");
assert!(
msg.contains(&wrong_path.display().to_string()),
"error must name the resolved (wrong) path: {msg}"
);
assert!(
msg.contains(&anchor.display().to_string()),
"error must name the canonical anchor path: {msg}"
);
}
#[tokio::test]
async fn coordinator_link_annotates_resolves_edge_target_like_get() {
use khive_mcp::tools::request::RequestParams;
use khive_runtime::PackConfig;
let khive_cfg = KhiveConfig {
backends: vec![
khive_runtime::BackendConfig {
name: "main".to_string(),
kind: khive_runtime::BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
khive_runtime::BackendConfig {
name: "sessions".to_string(),
kind: khive_runtime::BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"session".to_string(),
PackConfig {
backend: "sessions".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let base_cfg = RuntimeConfig {
packs: vec!["kg".to_string(), "session".to_string()],
..base_multi_backend_runtime_config()
};
let (server, _schedule_rt) =
build_multi_backend_server_with_coordinator(base_cfg, &khive_cfg, None)
.expect("coordinator-attached multi-backend boot must succeed");
let dispatch = |ops: String| {
let server = &server;
async move {
let resp = server
.dispatch_request_local(RequestParams {
ops,
presentation: Some("verbose".to_string()),
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.expect("dispatch must not error");
serde_json::from_str::<serde_json::Value>(&resp).expect("valid JSON")
}
};
let a = dispatch(r#"create(kind="concept", name="edge-endpoint-a")"#.to_string()).await;
let a_id = a["results"][0]["result"]["id"]
.as_str()
.expect("create must return an id")
.to_string();
let b = dispatch(r#"create(kind="concept", name="edge-endpoint-b")"#.to_string()).await;
let b_id = b["results"][0]["result"]["id"]
.as_str()
.expect("create must return an id")
.to_string();
let edge_resp = dispatch(format!(
r#"link(source_id="{a_id}", target_id="{b_id}", relation="extends")"#
))
.await;
assert_eq!(
edge_resp["results"][0]["ok"].as_bool(),
Some(true),
"seed edge creation must succeed: {edge_resp}"
);
let edge_id = edge_resp["results"][0]["result"]["id"]
.as_str()
.expect("link must return an edge id")
.to_string();
let note_resp =
dispatch(r#"create(kind="observation", content="annotates source")"#.to_string()).await;
let note_id = note_resp["results"][0]["result"]["id"]
.as_str()
.expect("create must return an id")
.to_string();
let got_edge = dispatch(format!(r#"get(id="{edge_id}")"#)).await;
assert_eq!(
got_edge["results"][0]["ok"].as_bool(),
Some(true),
"get(<edge_uuid>) must succeed: {got_edge}"
);
assert_eq!(
got_edge["results"][0]["result"]["kind"].as_str(),
Some("edge"),
"get must resolve the UUID as an edge: {got_edge}"
);
let annotate_resp = dispatch(format!(
r#"link(source_id="{note_id}", target_id="{edge_id}", relation="annotates")"#
))
.await;
assert_eq!(
annotate_resp["results"][0]["ok"].as_bool(),
Some(true),
"note->edge annotates link must succeed, proving get/link resolution parity \
for an edge-substrate UUID under multi-backend pack bindings: {annotate_resp}"
);
assert_eq!(
annotate_resp["results"][0]["result"]["target_id"].as_str(),
got_edge["results"][0]["result"]["id"].as_str(),
"link's resolved annotates target must be the exact same edge UUID get() resolved: \
annotate_resp={annotate_resp} got_edge={got_edge}"
);
}
}