use std::collections::BTreeMap;
use std::io::BufRead as _;
use std::path::PathBuf;
use anyhow::{Context, Result};
use clap::Parser;
#[cfg(test)]
use khive_mcp::serve::resolve_runtime_config;
use khive_mcp::serve::{
apply_env_output_format, build_server_multi_backend_with_db_anchor, config_discovery_db_anchor,
enforce_strict_actor_mode, RuntimeConfigInputs,
};
#[cfg(unix)]
use khive_mcp::server::compute_config_id;
use khive_mcp::server::KhiveMcpServer;
use khive_mcp::tools::request::RequestParams;
#[cfg(unix)]
use khive_runtime::{daemon::PROTOCOL_VERSION, DaemonRequestFrame};
use khive_runtime::{KhiveConfig, KhiveRuntime, Namespace, RuntimeConfig};
#[cfg(unix)]
type ForwardFuture<'a> = std::pin::Pin<
Box<dyn std::future::Future<Output = Option<Result<String, rmcp::ErrorData>>> + Send + 'a>,
>;
#[cfg(unix)]
type ForwardFnPtr = for<'a> fn(&'a DaemonRequestFrame) -> ForwardFuture<'a>;
#[cfg(unix)]
fn forward_or_spawn_boxed(frame: &DaemonRequestFrame) -> ForwardFuture<'_> {
Box::pin(khive_mcp::daemon::forward_or_spawn(frame))
}
use khive_mcp::pending_events;
#[cfg(unix)]
type LocalConstructionGuard = Option<khive_runtime::daemon::DaemonBootGuard>;
#[cfg(not(unix))]
type LocalConstructionGuard = Option<std::fs::File>;
#[cfg(unix)]
pub(crate) fn acquire_local_construction_guard(
cfg: &RuntimeConfig,
) -> Result<LocalConstructionGuard> {
if cfg.db_path.is_none() {
return Ok(None);
}
Ok(Some(
khive_runtime::daemon::acquire_daemon_boot_guard().context(
"acquire daemon boot/recovery guard for local kkernel exec construction \
(another process may be cold-booting the same database)",
)?,
))
}
#[cfg(not(unix))]
pub(crate) fn acquire_local_construction_guard(
cfg: &RuntimeConfig,
) -> Result<LocalConstructionGuard> {
if cfg.db_path.is_none() {
return Ok(None);
}
let path = khive_runtime::daemon::lock_path();
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).with_context(|| {
format!("create parent directory for construction guard lock file {path:?}")
})?;
}
let file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(&path)
.with_context(|| format!("open construction guard lock file {path:?}"))?;
file.lock().context(
"acquire local construction guard lock for kkernel exec construction \
(another process may be cold-booting the same database)",
)?;
Ok(Some(file))
}
const OPS_FILE_CHUNK_SIZE: usize = 100;
#[derive(Parser, Debug)]
pub struct ExecArgs {
pub ops: Option<String>,
#[arg(long, conflicts_with = "ops", conflicts_with = "ops_file")]
pub pending_events: bool,
#[arg(long, env = "KHIVE_DB")]
pub db: Option<String>,
#[arg(long, default_value = "local")]
pub namespace: String,
#[arg(long, default_value = "verbose")]
pub presentation: Option<String>,
#[arg(long, value_name = "FORMAT")]
pub output_format: Option<String>,
#[arg(long, short = 'v')]
pub verbose: bool,
#[arg(long)]
pub save_file: Option<String>,
#[arg(long, value_name = "PATH")]
pub ops_file: Option<PathBuf>,
#[arg(long, requires = "ops_file")]
pub dry_run: bool,
#[arg(long, requires = "ops_file")]
pub atomic: bool,
#[arg(long, requires = "atomic")]
pub atomic_max_ops: Option<usize>,
}
#[derive(Debug, Clone)]
pub(crate) struct OpsFileEntry {
pub(crate) tool: String,
pub(crate) args: serde_json::Value,
}
pub(crate) fn parse_ops_file(path: &PathBuf) -> Result<Vec<OpsFileEntry>> {
let file =
std::fs::File::open(path).with_context(|| format!("open ops-file {}", path.display()))?;
let reader = std::io::BufReader::new(file);
let mut ops: Vec<OpsFileEntry> = Vec::new();
for (line_idx, result) in reader.lines().enumerate() {
let line_num = line_idx + 1;
let raw = result.with_context(|| format!("read ops-file line {line_num}"))?;
let trimmed = raw.trim();
if trimmed.is_empty() {
continue;
}
let obj: serde_json::Value = serde_json::from_str(trimmed)
.map_err(|e| anyhow::anyhow!("ops-file line {line_num}: invalid JSON: {e}"))?;
let obj = obj.as_object().ok_or_else(|| {
anyhow::anyhow!(
"ops-file line {line_num}: expected a JSON object {{\"tool\":...,\"args\":...}}, \
got a non-object value"
)
})?;
let tool = obj
.get("tool")
.and_then(|v| v.as_str())
.ok_or_else(|| {
anyhow::anyhow!("ops-file line {line_num}: missing or non-string \"tool\" field")
})?
.to_owned();
let args = match obj.get("args") {
None => serde_json::Value::Object(serde_json::Map::new()),
Some(v) => {
if !v.is_object() {
anyhow::bail!(
"ops-file line {line_num}: \"args\" must be a JSON object, got {v}"
);
}
v.clone()
}
};
ops.push(OpsFileEntry { tool, args });
}
Ok(ops)
}
async fn apply_ops_file(
server: &KhiveMcpServer,
ops: Vec<OpsFileEntry>,
presentation: Option<String>,
) -> Result<()> {
let total = ops.len();
let mut total_succeeded: usize = 0;
let mut total_failed: usize = 0;
for (chunk_idx, chunk) in ops.chunks(OPS_FILE_CHUNK_SIZE).enumerate() {
let applied_before = (chunk_idx * OPS_FILE_CHUNK_SIZE).min(total);
let batch_arr: Vec<serde_json::Value> = chunk
.iter()
.map(|e| {
serde_json::json!({
"tool": e.tool,
"args": e.args,
})
})
.collect();
let batch_json = serde_json::to_string(&batch_arr).context("serialize chunk to JSON")?;
let params = RequestParams {
ops: batch_json,
presentation: presentation.clone(),
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server
.dispatch_request_local(params)
.await
.map_err(|e| anyhow::anyhow!("dispatch chunk {}: {}", chunk_idx + 1, e))?;
let parsed: serde_json::Value =
serde_json::from_str(&raw).context("parse dispatch result")?;
let chunk_succeeded = parsed["summary"]["succeeded"].as_u64().unwrap_or(0) as usize;
let chunk_failed = parsed["summary"]["failed"].as_u64().unwrap_or(0) as usize;
total_succeeded += chunk_succeeded;
total_failed += chunk_failed;
let applied_now = applied_before + chunk.len();
eprintln!("applied {applied_now}/{total} (ok={total_succeeded}, failed={total_failed})");
}
let summary = serde_json::json!({
"total": total,
"succeeded": total_succeeded,
"failed": total_failed,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize summary")
);
Ok(())
}
pub async fn run_exec(args: ExecArgs) -> Result<()> {
if args.pending_events {
let summary =
pending_events::run_pending_events(args.db.as_deref(), &args.namespace, args.verbose)
.await?;
pending_events::print_summary(&summary);
return Ok(());
}
let mode = match (&args.ops, &args.ops_file) {
(Some(_), Some(_)) => {
anyhow::bail!(
"cannot use both a positional ops string and --ops-file; supply exactly one"
);
}
(None, None) => {
anyhow::bail!(
"no ops provided; supply a DSL expression as a positional argument or use \
--ops-file <PATH>"
);
}
(Some(ops), None) => ExecMode::Inline(ops.clone()),
(None, Some(path)) => ExecMode::OpsFile(path.clone()),
};
let namespace = Namespace::parse(&args.namespace).map_err(|e| anyhow::anyhow!("{e}"))?;
let (cfg, db_anchor) =
khive_mcp::serve::resolve_runtime_config_with_db_anchor(RuntimeConfigInputs {
db: args.db.as_deref(),
config: None, namespace,
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})?;
khive_runtime::assert_captured_db_anchor_consistent(
cfg.db_path.as_deref(),
db_anchor.as_deref(),
)?;
let db_context = ExecDbContext {
raw: args.db,
anchor: db_anchor,
};
match mode {
ExecMode::Inline(ops) => {
run_exec_inline(
ops,
cfg,
args.presentation,
args.output_format,
args.save_file,
db_context,
)
.await
}
ExecMode::OpsFile(path) => {
run_exec_ops_file(
path,
cfg,
args.presentation,
args.dry_run,
db_context,
args.atomic,
args.atomic_max_ops,
)
.await
}
}
}
enum ExecMode {
Inline(String),
OpsFile(PathBuf),
}
#[derive(Default)]
struct ExecDbContext {
raw: Option<String>,
anchor: Option<PathBuf>,
}
async fn run_exec_inline(
ops: String,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
db_context: ExecDbContext,
) -> Result<()> {
#[cfg(unix)]
return run_exec_inline_with_forward(
ops,
cfg,
presentation,
output_format,
save_file,
db_context,
forward_or_spawn_boxed,
)
.await;
#[cfg(not(unix))]
return run_exec_inline_with_forward(
ops,
cfg,
presentation,
output_format,
save_file,
db_context,
)
.await;
}
#[cfg_attr(not(unix), allow(unused_variables))]
async fn run_exec_inline_with_forward(
ops: String,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
db_context: ExecDbContext,
#[cfg(unix)] forward_fn: ForwardFnPtr,
) -> Result<()> {
enforce_strict_actor_mode(cfg.actor_id.as_deref(), &cfg.packs)?;
let db_path_for_config = config_discovery_db_anchor(db_context.raw.as_deref());
let khive_cfg = KhiveConfig::load_with_home_fallback(None, db_path_for_config.as_deref())
.map_err(|e| anyhow::anyhow!("config error: {e}"))?
.unwrap_or_default();
#[cfg(unix)]
if save_file.is_none() {
let frame = DaemonRequestFrame {
ops: ops.clone(),
presentation: presentation.clone(),
presentation_per_op: None,
namespace: cfg.default_namespace.as_str().to_string(),
actor_id: cfg.actor_id.clone(),
visible_namespaces: cfg
.visible_namespaces
.iter()
.map(|ns| ns.as_str().to_string())
.collect(),
config_id: compute_config_id(&cfg, Some(&khive_cfg)),
protocol_version: PROTOCOL_VERSION,
probe_only: false,
metrics_only: false,
format: output_format.clone(),
format_per_op: None,
from_wire: false,
request_id: None,
};
if let Some(res) = forward_fn(&frame).await {
let output = res.map_err(|e| anyhow::anyhow!("{}", e.message))?;
println!("{output}");
return Ok(());
}
}
let server = build_local_fallback_server(
cfg,
&khive_cfg,
db_context.raw.as_deref(),
db_context.anchor.as_deref(),
)?;
let params = RequestParams {
ops,
presentation,
presentation_per_op: None,
save_to: save_file,
format: output_format,
format_per_op: None,
request_id: None,
};
let output = server
.dispatch_request_local(params)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
println!("{output}");
Ok(())
}
fn build_local_fallback_server(
cfg: RuntimeConfig,
khive_cfg: &KhiveConfig,
cli_db_override: Option<&str>,
db_anchor: Option<&std::path::Path>,
) -> Result<KhiveMcpServer> {
let _boot_guard = acquire_local_construction_guard(&cfg)?;
if khive_cfg.backends.is_empty() {
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
let env_fmt = apply_env_output_format(khive_cfg.runtime.default_output_format);
Ok(KhiveMcpServer::new(rt)
.map_err(|e| anyhow::anyhow!("{e}"))?
.with_default_output_format(env_fmt))
} else {
build_server_multi_backend_with_db_anchor(cfg, khive_cfg, cli_db_override, db_anchor)
}
}
async fn run_exec_ops_file(
path: PathBuf,
cfg: RuntimeConfig,
presentation: Option<String>,
dry_run: bool,
db_context: ExecDbContext,
atomic: bool,
atomic_max_ops: Option<usize>,
) -> Result<()> {
let ops = parse_ops_file(&path)?;
if ops.is_empty() {
anyhow::bail!("ops-file is empty (no non-blank lines): {}", path.display());
}
if dry_run {
let mut per_verb: BTreeMap<String, usize> = BTreeMap::new();
for op in &ops {
*per_verb.entry(op.tool.clone()).or_insert(0) += 1;
}
let summary = serde_json::json!({
"dry_run": true,
"total": ops.len(),
"per_verb": per_verb,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize dry-run summary")
);
return Ok(());
}
enforce_strict_actor_mode(cfg.actor_id.as_deref(), &cfg.packs)?;
let db_path_for_config = config_discovery_db_anchor(db_context.raw.as_deref());
let khive_cfg = KhiveConfig::load_with_home_fallback(None, db_path_for_config.as_deref())
.map_err(|e| anyhow::anyhow!("config error: {e}"))?
.unwrap_or_default();
if atomic {
let max_ops = atomic_max_ops.unwrap_or(khive_types::pack::ATOMIC_MAX_OPS_DEFAULT);
let envelope = crate::atomic_apply::execute_atomic_ops_file(ops, cfg, &khive_cfg, max_ops)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
println!(
"{}",
serde_json::to_string_pretty(&envelope).expect("serialize atomic envelope")
);
return Ok(());
}
let server = build_local_fallback_server(
cfg,
&khive_cfg,
db_context.raw.as_deref(),
db_context.anchor.as_deref(),
)?;
apply_ops_file(&server, ops, presentation).await
}
#[cfg(test)]
mod tests {
use super::*;
use clap::Parser;
use serial_test::serial;
use tempfile::NamedTempFile;
fn isolate_home_for_test() -> (Option<std::ffi::OsString>, tempfile::TempDir) {
let prev = std::env::var_os("HOME");
let dir = tempfile::tempdir().expect("tempdir for isolated HOME");
std::env::set_var("HOME", dir.path());
(prev, dir)
}
fn restore_home(prev: Option<std::ffi::OsString>) {
match prev {
Some(v) => std::env::set_var("HOME", v),
None => std::env::remove_var("HOME"),
}
}
#[test]
#[serial(local_exec_boot_guard)]
fn acquire_local_construction_guard_is_noop_for_in_memory_db() {
let dir = tempfile::tempdir().expect("tempdir");
std::env::set_var("KHIVE_LOCK", dir.path().join("khived.recovery.lock"));
let cfg = RuntimeConfig {
db_path: None,
..RuntimeConfig::default()
};
let guard = acquire_local_construction_guard(&cfg).expect("in-memory db needs no guard");
assert!(
guard.is_none(),
"an in-memory database has no shared file to serialize construction against"
);
std::env::remove_var("KHIVE_LOCK");
}
#[cfg(unix)]
#[test]
#[serial(local_exec_boot_guard)]
fn acquire_local_construction_guard_serializes_concurrent_file_backed_callers() {
acquire_local_construction_guard_serializes_concurrent_file_backed_callers_impl();
}
#[cfg(not(unix))]
#[test]
#[serial(local_exec_boot_guard)]
fn acquire_local_construction_guard_serializes_concurrent_file_backed_callers_nonunix() {
acquire_local_construction_guard_serializes_concurrent_file_backed_callers_impl();
}
fn acquire_local_construction_guard_serializes_concurrent_file_backed_callers_impl() {
let dir = tempfile::tempdir().expect("tempdir");
std::env::set_var("KHIVE_LOCK", dir.path().join("khived.recovery.lock"));
let db_path = dir.path().join("cold.db3");
let concurrent = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let max_observed = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let spawn_one = |label: &'static str| {
let db_path = db_path.clone();
let concurrent = concurrent.clone();
let max_observed = max_observed.clone();
std::thread::spawn(move || {
let cfg = RuntimeConfig {
db_path: Some(db_path),
..RuntimeConfig::default()
};
let guard = acquire_local_construction_guard(&cfg)
.unwrap_or_else(|e| panic!("{label} must acquire the guard: {e}"));
let now = concurrent.fetch_add(1, std::sync::atomic::Ordering::SeqCst) + 1;
max_observed.fetch_max(now, std::sync::atomic::Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(50));
concurrent.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
drop(guard);
})
};
let t_a = spawn_one("writer-a");
let t_b = spawn_one("writer-b");
t_a.join().expect("writer-a thread must not panic");
t_b.join().expect("writer-b thread must not panic");
assert_eq!(
max_observed.load(std::sync::atomic::Ordering::SeqCst),
1,
"the two guarded critical sections must never overlap — the guard \
failed to serialize concurrent local-construction callers"
);
std::env::remove_var("KHIVE_LOCK");
}
#[test]
#[serial]
fn khive_db_env_binds_to_db_arg() {
std::env::set_var("KHIVE_DB", "/tmp/kkernel-exec-env.db");
let args = ExecArgs::parse_from(["exec", "stats()"]);
std::env::remove_var("KHIVE_DB");
assert_eq!(args.db.as_deref(), Some("/tmp/kkernel-exec-env.db"));
}
#[test]
fn pending_events_flag_sets_mode() {
let args = ExecArgs::parse_from(["exec", "--pending-events"]);
assert!(args.pending_events);
assert!(args.ops.is_none());
}
#[test]
fn pending_events_conflicts_with_ops() {
let result = ExecArgs::try_parse_from(["exec", "--pending-events", "stats()"]);
assert!(
result.is_err(),
"--pending-events and positional ops must conflict"
);
}
#[test]
fn pending_events_conflicts_with_ops_file() {
let result =
ExecArgs::try_parse_from(["exec", "--pending-events", "--ops-file", "/tmp/x.jsonl"]);
assert!(
result.is_err(),
"--pending-events and --ops-file must conflict"
);
}
#[test]
fn ops_positional_is_optional() {
let args = ExecArgs::parse_from(["exec", "--ops-file", "/tmp/batch.jsonl"]);
assert!(args.ops.is_none());
assert_eq!(
args.ops_file.as_deref(),
Some(std::path::Path::new("/tmp/batch.jsonl"))
);
}
#[test]
fn ops_positional_works_without_pending_events() {
let args = ExecArgs::parse_from(["exec", "stats()"]);
assert_eq!(args.ops.as_deref(), Some("stats()"));
assert!(!args.pending_events);
}
#[test]
fn presentation_defaults_to_verbose_when_flag_omitted() {
let args = ExecArgs::parse_from(["exec", "stats()"]);
assert_eq!(args.presentation.as_deref(), Some("verbose"));
}
#[test]
fn presentation_agent_flag_still_selects_agent() {
let args = ExecArgs::parse_from(["exec", "stats()", "--presentation", "agent"]);
assert_eq!(args.presentation.as_deref(), Some("agent"));
}
#[test]
fn presentation_human_flag_still_selects_human() {
let args = ExecArgs::parse_from(["exec", "stats()", "--presentation", "human"]);
assert_eq!(args.presentation.as_deref(), Some("human"));
}
#[test]
fn dry_run_requires_ops_file() {
let result = ExecArgs::try_parse_from(["exec", "stats()", "--dry-run"]);
assert!(
result.is_err(),
"dry-run without --ops-file should be rejected by clap"
);
}
fn isolated_server(db_path: &str) -> KhiveMcpServer {
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(db_path)),
embedding_model: None,
additional_embedding_models: vec![],
..Default::default()
};
let rt = KhiveRuntime::new(cfg).expect("runtime on temp db");
KhiveMcpServer::new(rt).expect("server on temp db")
}
#[test]
#[serial]
fn exec_config_id_matches_serve_config_id_for_project_toml_actor() {
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::remove_var("KHIVE_ACTOR");
let dir = tempfile::tempdir().expect("tempdir");
let khive_dir = dir.path().join(".khive");
std::fs::create_dir_all(&khive_dir).expect("mkdir .khive");
std::fs::write(
khive_dir.join("config.toml"),
r#"
[actor]
id = "lambda:test-actor"
[[engines]]
name = "primary"
model = "bge-small-en-v1.5"
default = true
"#,
)
.expect("write config.toml");
let db_path = khive_dir.join("exec-parity-test.db");
let db_str = db_path.to_str().expect("utf8 path").to_string();
let ns = Namespace::parse("local").expect("ns");
let exec_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_str),
config: None,
namespace: ns.clone(),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve exec-shaped config");
let serve_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&db_str),
config: None,
namespace: ns,
namespace_explicit: false,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve serve-shaped config");
assert_eq!(exec_cfg.actor_id.as_deref(), Some("lambda:test-actor"));
assert_eq!(serve_cfg.actor_id.as_deref(), Some("lambda:test-actor"));
assert!(
exec_cfg
.visible_namespaces
.contains(&Namespace::parse("lambda:test-actor").expect("ns")),
"actor.id must fold into visible_namespaces (ADR-007 Rev 4 Rule 3b)"
);
assert!(
exec_cfg.embedding_model.is_some(),
"config-file [[engines]] must resolve an embedding model, not env/default"
);
assert_eq!(
format!("{:?}", exec_cfg.embedding_model),
format!("{:?}", serve_cfg.embedding_model),
);
assert_eq!(
compute_config_id(&exec_cfg, None),
compute_config_id(&serve_cfg, None),
"exec-path config_id must match the serve/daemon-path config_id for the same db"
);
}
#[test]
#[serial]
fn namespace_explicit_changes_actor_id_fill_but_not_config_id() {
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::remove_var("KHIVE_ACTOR");
let missing_config =
std::path::PathBuf::from("/nonexistent/khive-exec-parity-test/config.toml");
let ns = Namespace::parse("lambda:custom-ns").expect("ns");
let with_explicit_true = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&missing_config),
namespace: ns.clone(),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve with namespace_explicit=true");
let with_explicit_false = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&missing_config),
namespace: ns,
namespace_explicit: false,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve with namespace_explicit=false");
assert_eq!(
with_explicit_true.actor_id.as_deref(),
Some("lambda:custom-ns"),
"namespace_explicit=true + non-local namespace + no config actor.id \
must fill actor_id from the namespace (ADR-057)"
);
assert_eq!(
with_explicit_false.actor_id, None,
"namespace_explicit=false must NOT fill actor_id"
);
assert_eq!(
compute_config_id(&with_explicit_true, None),
compute_config_id(&with_explicit_false, None),
"namespace_explicit must not affect the daemon-forwarded config_id"
);
}
#[test]
#[serial]
fn exec_config_id_matches_serve_config_id_for_multi_backend_topology() {
use khive_runtime::{BackendConfig, BackendKind, PackConfig};
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::remove_var("KHIVE_ACTOR");
let missing_config = std::path::PathBuf::from(
"/nonexistent/khive-exec-parity-test/multi-backend-config.toml",
);
let ns = Namespace::parse("local").expect("ns");
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Sqlite,
path: Some(std::path::PathBuf::from("/tmp/khive-parity-main.db")),
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "sessions".to_string(),
kind: BackendKind::Sqlite,
path: Some(std::path::PathBuf::from("/tmp/khive-parity-sessions.db")),
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 exec_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&missing_config),
namespace: ns.clone(),
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve exec-shaped config");
let serve_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&missing_config),
namespace: ns,
namespace_explicit: false,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve serve-shaped config");
assert_ne!(
compute_config_id(&exec_cfg, None),
compute_config_id(&serve_cfg, Some(&khive_cfg)),
"pre-fix exec computation (None) must diverge from the daemon computation \
(Some) for a non-empty backends topology — proves this test catches the \
real divergence, not a tautology"
);
assert_eq!(
compute_config_id(&exec_cfg, Some(&khive_cfg)),
compute_config_id(&serve_cfg, Some(&khive_cfg)),
"exec-path config_id must match the daemon-path config_id for the same \
multi-backend topology (D1 fix acceptance gate)"
);
}
#[tokio::test]
#[serial]
async fn build_local_fallback_server_routes_through_multi_backend_when_backends_declared() {
use khive_runtime::{BackendConfig, BackendKind, PackConfig};
let main_db = NamedTempFile::new().expect("main db tempfile");
let secondary_db = NamedTempFile::new().expect("secondary db tempfile");
let main_path = main_db.path().to_path_buf();
let secondary_path = secondary_db.path().to_path_buf();
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Sqlite,
path: Some(main_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Sqlite,
path: Some(secondary_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"kg".to_string(),
PackConfig {
backend: "secondary".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let cfg = RuntimeConfig {
db_path: khive_runtime::resolve_db_anchor(None),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
actor_id: Some("actor-routing-test".to_string()),
..RuntimeConfig::default()
};
let db_anchor = cfg.db_path.clone();
let server = build_local_fallback_server(cfg, &khive_cfg, None, db_anchor.as_deref())
.expect("multi-backend local fallback must build");
let create = server
.dispatch_request_local(RequestParams {
ops: r#"create(kind="entity", entity_kind="concept", name="routed-via-secondary")"#
.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("create must dispatch");
let create_resp: serde_json::Value = serde_json::from_str(&create).expect("valid JSON");
assert_eq!(
create_resp["results"][0]["ok"].as_bool(),
Some(true),
"create must succeed through the multi-backend fallback server: {create_resp}"
);
async fn count_concepts(db_path: &std::path::Path) -> usize {
let cfg = RuntimeConfig {
db_path: Some(db_path.to_path_buf()),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
let rt = KhiveRuntime::new(cfg).expect("runtime on backend file");
let probe = KhiveMcpServer::new(rt).expect("server on backend file");
let raw = probe
.dispatch_request_local(RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("list must dispatch");
let resp: serde_json::Value = serde_json::from_str(&raw).expect("valid JSON");
resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0)
}
let main_count = count_concepts(&main_path).await;
let secondary_count = count_concepts(&secondary_path).await;
assert_eq!(
main_count, 0,
"kg pack must NOT write into the `main` backend file when pinned to \
`secondary` (D1-R2: a silent single-backend fallback would have written \
it here instead)"
);
assert_eq!(
secondary_count, 1,
"kg pack write must land in its declared `secondary` backend file"
);
}
#[test]
#[serial]
fn build_local_fallback_server_uses_captured_anchor_after_home_changes() {
let (previous_home, _first_home) = isolate_home_for_test();
let cfg = RuntimeConfig {
db_path: khive_runtime::resolve_db_anchor(None),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
let db_anchor = cfg.db_path.clone();
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,
}],
..KhiveConfig::default()
};
let second_home = tempfile::tempdir().expect("second HOME");
std::env::set_var("HOME", second_home.path());
let result = build_local_fallback_server(cfg, &khive_cfg, None, db_anchor.as_deref());
restore_home(previous_home);
assert!(
result.is_ok(),
"exec fallback must use the anchor captured with RuntimeConfig after HOME changes: {}",
result.err().unwrap()
);
}
#[cfg(unix)]
fn run_one_guarded_daemon_boot(
db_path: std::path::PathBuf,
writer_label: &'static str,
count: usize,
) {
let guard =
khive_runtime::daemon::acquire_recovery_lock().expect("acquire daemon boot guard");
let rt_handle = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build per-thread tokio runtime");
rt_handle.block_on(async {
let rt = KhiveRuntime::new(RuntimeConfig {
db_path: Some(db_path),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
})
.expect("cold-boot migrations succeed");
let token = rt.authorize(Namespace::local()).expect("authorize local");
for i in 0..count {
rt.create_note(
&token,
"memo",
None,
&format!("{writer_label} note {i} — boot race marker"),
None,
None,
vec![],
)
.await
.expect("note write must succeed inside the guarded boot window");
}
});
drop(guard);
}
#[cfg(unix)]
fn run_one_local_exec_construction(
db_path: std::path::PathBuf,
writer_label: &'static str,
count: usize,
) {
let rt_handle = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build per-thread tokio runtime");
rt_handle.block_on(async {
let cfg = RuntimeConfig {
db_path: Some(db_path),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
};
let khive_cfg = KhiveConfig::default();
let server = build_local_fallback_server(cfg, &khive_cfg, None, None)
.expect("guarded local-exec construction must succeed");
for i in 0..count {
let params = RequestParams {
ops: format!(
r#"create(kind="observation", content="{writer_label} note {i} — boot race marker")"#
),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server
.dispatch_request_local(params)
.await
.expect("dispatch must succeed inside the guarded construction window");
let resp: serde_json::Value = serde_json::from_str(&raw).expect("valid JSON");
assert_eq!(
resp["results"][0]["ok"],
serde_json::json!(true),
"write must succeed: {resp}"
);
}
});
}
#[cfg(unix)]
#[test]
#[serial(local_exec_boot_guard)]
fn build_local_fallback_server_blocks_while_recovery_lock_is_held() {
let dir = tempfile::tempdir().expect("tempdir");
let lock_file = dir.path().join("khived.recovery.lock");
std::env::set_var("KHIVE_LOCK", &lock_file);
let db_path = dir.path().join("guard_block_test.db3");
let held_guard =
khive_runtime::daemon::acquire_recovery_lock().expect("acquire recovery lock in test");
let (tx, rx) = std::sync::mpsc::channel();
let handle = std::thread::spawn(move || {
let rt_handle = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build per-thread tokio runtime");
let cfg = RuntimeConfig {
db_path: Some(db_path),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
};
let khive_cfg = KhiveConfig::default();
let result = rt_handle
.block_on(async { build_local_fallback_server(cfg, &khive_cfg, None, None) });
let _ = tx.send(());
result
});
let completed_while_locked = rx
.recv_timeout(std::time::Duration::from_millis(500))
.is_ok();
assert!(
!completed_while_locked,
"build_local_fallback_server must NOT complete while the boot/recovery \
lock is held by another holder — if this fires, the guard at its \
production call site has been removed or stopped acquiring the shared lock"
);
drop(held_guard);
handle
.join()
.expect("construction thread must not panic")
.expect("construction must succeed once the lock is released");
std::env::remove_var("KHIVE_LOCK");
}
#[cfg(unix)]
#[test]
#[serial(local_exec_boot_guard)]
fn local_exec_construction_races_guarded_daemon_boot_without_fts_corruption() {
let dir = tempfile::tempdir().expect("tempdir");
let lock_file = dir.path().join("khived.recovery.lock");
std::env::set_var("KHIVE_LOCK", &lock_file);
let db_path = dir.path().join("local_exec_boot_race.db3");
const PER_WRITER: usize = 10;
let path_a = db_path.clone();
let path_b = db_path.clone();
let t_a = std::thread::spawn(move || {
run_one_guarded_daemon_boot(path_a, "daemon-boot", PER_WRITER)
});
let t_b = std::thread::spawn(move || {
run_one_local_exec_construction(path_b, "local-exec", PER_WRITER)
});
t_a.join().expect("daemon-boot thread must not panic");
t_b.join().expect("local-exec thread must not panic");
let rt_handle = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("build verification tokio runtime");
rt_handle.block_on(async {
let verify_rt = KhiveRuntime::new(RuntimeConfig {
db_path: Some(db_path.clone()),
embedding_model: None,
additional_embedding_models: vec![],
..RuntimeConfig::default()
})
.expect("post-race runtime opens cleanly");
let token = verify_rt
.authorize(Namespace::local())
.expect("authorize local");
let hits = verify_rt
.search_notes(
&token,
"boot race marker",
None,
100,
None,
false,
&[],
None,
)
.await
.expect("FTS search over notes must succeed, not error on a corrupted index");
assert_eq!(
hits.len(),
PER_WRITER * 2,
"every planted note from both writers must be present and \
FTS-searchable — a corrupted/partial index would drop or \
duplicate rows: {hits:?}"
);
});
std::env::remove_var("KHIVE_LOCK");
}
#[test]
fn parse_ops_file_skips_blank_lines() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
f.write_all(b"\n").unwrap(); f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
let ops = parse_ops_file(&f.path().to_path_buf()).unwrap();
assert_eq!(ops.len(), 2);
}
#[test]
fn parse_ops_file_reports_line_number_on_malformed() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
f.write_all(b"not-json\n").unwrap(); let err = parse_ops_file(&f.path().to_path_buf()).unwrap_err();
let msg = format!("{err:#}");
assert!(
msg.contains("line 2"),
"error should name the bad line number, got: {msg}"
);
}
#[test]
fn parse_ops_file_missing_tool_field() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"notool\":\"x\",\"args\":{}}\n").unwrap();
let err = parse_ops_file(&f.path().to_path_buf()).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("line 1"), "should report line number: {msg}");
}
#[tokio::test]
async fn ops_file_applies_ops_and_summary_matches() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let server = isolated_server(&db_path);
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
for name in ["Alpha", "Beta", "Gamma"] {
let line = format!(
"{{\"tool\":\"create\",\"args\":{{\"kind\":\"concept\",\"name\":\"{name}\"}}}}\n"
);
f.write_all(line.as_bytes()).unwrap();
}
let ops = parse_ops_file(&f.path().to_path_buf()).unwrap();
assert_eq!(ops.len(), 3);
apply_ops_file(&server, ops, None).await.unwrap();
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(
count, 3,
"all 3 entities should be present after apply\nraw: {resp}"
);
}
#[tokio::test]
async fn non_atomic_dispatch_envelope_shape_is_unchanged_by_adr099_b1() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let server = isolated_server(&db_path);
async fn dispatch(server: &KhiveMcpServer, ops: &str) -> serde_json::Value {
let params = RequestParams {
ops: ops.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server
.dispatch_request_local(params)
.await
.unwrap_or_else(|e| panic!("dispatch {ops:?} failed: {e}"));
serde_json::from_str(&raw).expect("valid JSON")
}
let created = dispatch(
&server,
r#"create(kind="concept", name="ADR-099-B1-inertness")"#,
)
.await;
assert_golden_envelope_shape(&created, "create");
let entity_id = created["results"][0]["result"]["id"]
.as_str()
.expect("create must return an id")
.to_string();
let updated = dispatch(
&server,
&format!(r#"update(id="{entity_id}", description="updated by inertness test")"#),
)
.await;
assert_golden_envelope_shape(&updated, "update");
let target = dispatch(&server, r#"create(kind="concept", name="link-target")"#).await;
let target_id = target["results"][0]["result"]["id"]
.as_str()
.expect("create must return an id")
.to_string();
let linked = dispatch(
&server,
&format!(
r#"link(source_id="{entity_id}", target_id="{target_id}", relation="extends")"#
),
)
.await;
assert_golden_envelope_shape(&linked, "link");
let got = dispatch(&server, &format!(r#"get(id="{entity_id}")"#)).await;
assert_golden_envelope_shape(&got, "get");
}
fn assert_golden_envelope_shape(resp: &serde_json::Value, expected_tool: &str) {
let top_level_keys: std::collections::BTreeSet<&str> = resp
.as_object()
.expect("response must be a JSON object")
.keys()
.map(String::as_str)
.collect();
assert_eq!(
top_level_keys,
std::collections::BTreeSet::from(["results", "summary"]),
"non-atomic envelope must carry exactly results+summary, no `atomic` block: {resp}"
);
let summary_keys: std::collections::BTreeSet<&str> = resp["summary"]
.as_object()
.expect("summary must be an object")
.keys()
.map(String::as_str)
.collect();
assert_eq!(
summary_keys,
std::collections::BTreeSet::from(["total", "succeeded", "failed", "aborted"]),
"summary shape must be unchanged: {resp}"
);
assert_eq!(resp["summary"]["total"], serde_json::json!(1));
assert_eq!(resp["summary"]["succeeded"], serde_json::json!(1));
assert_eq!(resp["summary"]["failed"], serde_json::json!(0));
assert_eq!(resp["results"][0]["ok"], serde_json::json!(true));
assert_eq!(resp["results"][0]["tool"], serde_json::json!(expected_tool));
assert!(
resp["results"][0].get("result").is_some(),
"results[0] must carry a `result` field: {resp}"
);
}
#[tokio::test]
async fn ops_file_dry_run_writes_nothing() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
for name in ["DryA", "DryB"] {
let line = format!(
"{{\"tool\":\"create\",\"args\":{{\"kind\":\"concept\",\"name\":\"{name}\"}}}}\n"
);
f.write_all(line.as_bytes()).unwrap();
}
let path = f.path().to_path_buf();
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(&db_path)),
..Default::default()
};
run_exec_ops_file(
path.clone(),
cfg.clone(),
None,
true,
ExecDbContext::default(),
false,
None,
)
.await
.unwrap();
let server = isolated_server(&db_path);
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(count, 0, "dry-run must not write any entities");
}
#[tokio::test]
#[serial]
async fn strict_mode_allows_exec_when_actor_configured() {
let prev_strict = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
let prev_no_daemon = std::env::var("KHIVE_NO_DAEMON").ok();
let (prev_home, _home_dir) = isolate_home_for_test();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
std::env::set_var("KHIVE_NO_DAEMON", "1");
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
actor_id: Some("lambda:tenant-x".to_string()), ..RuntimeConfig::default()
};
let result = run_exec_inline(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
)
.await;
match prev_strict {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
match prev_no_daemon {
Some(v) => std::env::set_var("KHIVE_NO_DAEMON", v),
None => std::env::remove_var("KHIVE_NO_DAEMON"),
}
restore_home(prev_home);
assert!(
result.is_ok(),
"run_exec_inline must succeed under strict mode when actor IS configured; got: {result:?}"
);
}
#[tokio::test]
#[serial]
async fn strict_mode_off_exec_inline_passes_with_no_actor() {
let prev_strict = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
let prev_no_daemon = std::env::var("KHIVE_NO_DAEMON").ok();
let (prev_home, _home_dir) = isolate_home_for_test();
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"); std::env::set_var("KHIVE_NO_DAEMON", "1");
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
actor_id: None,
..RuntimeConfig::default()
};
let result = run_exec_inline(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
)
.await;
match prev_strict {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
match prev_no_daemon {
Some(v) => std::env::set_var("KHIVE_NO_DAEMON", v),
None => std::env::remove_var("KHIVE_NO_DAEMON"),
}
restore_home(prev_home);
assert!(
result.is_ok(),
"run_exec_inline must NOT reject when strict mode is OFF (OSS default); got: {result:?}"
);
}
#[cfg(unix)]
std::thread_local! {
static SPY_CAPTURED_CONFIG_ID: std::cell::RefCell<Option<String>> =
const { std::cell::RefCell::new(None) };
}
#[cfg(unix)]
fn spy_capture_config_id(frame: &DaemonRequestFrame) -> super::ForwardFuture<'_> {
SPY_CAPTURED_CONFIG_ID.with(|c| *c.borrow_mut() = Some(frame.config_id.clone()));
Box::pin(async { None })
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn exec_frame_config_id_matches_daemon_config_id_for_multi_backend_project_toml() {
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::remove_var("KHIVE_ACTOR");
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
let (prev_home, home_dir) = isolate_home_for_test();
SPY_CAPTURED_CONFIG_ID.with(|c| *c.borrow_mut() = None);
let khive_dir = home_dir.path().join(".khive");
std::fs::create_dir_all(&khive_dir).expect("mkdir .khive");
let main_backend_path = khive_dir.join("main-backend.db");
let sessions_backend_path = khive_dir.join("sessions-backend.db");
std::fs::write(
khive_dir.join("config.toml"),
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "sessions"
kind = "sqlite"
path = "{}"
[packs.session]
backend = "sessions"
"#,
main_backend_path.display(),
sessions_backend_path.display(),
),
)
.expect("write multi-backend config.toml");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: None,
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
spy_capture_config_id,
)
.await;
assert!(result.is_ok(), "exec dispatch must succeed: {result:?}");
let captured = SPY_CAPTURED_CONFIG_ID
.with(|c| c.borrow_mut().take())
.expect("spy must have captured a forwarded frame");
let serve_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: None,
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
actor_explicit: false,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve serve-shaped config");
let khive_cfg = KhiveConfig::load_with_home_fallback(None, serve_cfg.db_path.as_deref())
.expect("load multi-backend config.toml")
.expect("config.toml must be found at tier 3");
assert!(
!khive_cfg.backends.is_empty(),
"sanity: the written config.toml must actually resolve with a non-empty \
backends list, or this test proves nothing"
);
let daemon_config_id = compute_config_id(&serve_cfg, Some(&khive_cfg));
restore_home(prev_home);
assert_eq!(
captured, daemon_config_id,
"the config_id in the ACTUAL frame run_exec_inline_with_forward sends to the \
daemon must be byte-identical to what the daemon computes for the same \
multi-backend config.toml (D1 acceptance gate, exercised end-to-end through \
the real call site rather than a standalone compute_config_id comparison)"
);
}
#[tokio::test]
async fn ops_file_malformed_line_aborts_before_writes() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
f.write_all(
b"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"ShouldNotExist\"}}\n",
)
.unwrap();
f.write_all(b"INVALID JSON LINE\n").unwrap();
let path = f.path().to_path_buf();
let err = parse_ops_file(&path).unwrap_err();
let msg = format!("{err:#}");
assert!(
msg.contains("line 2"),
"should report line 2 as malformed: {msg}"
);
let server = isolated_server(&db_path);
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(
count, 0,
"nothing should be written when any line fails to parse"
);
}
fn atomic_op(tool: &str, args: serde_json::Value) -> OpsFileEntry {
OpsFileEntry {
tool: tool.to_string(),
args,
}
}
async fn dispatch_json(server: &KhiveMcpServer, ops: &str) -> serde_json::Value {
let params = RequestParams {
ops: ops.to_string(),
presentation: Some("verbose".to_string()),
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
serde_json::from_str(&raw).unwrap()
}
fn atomic_cfg(db_path: &str) -> RuntimeConfig {
RuntimeConfig {
db_path: Some(PathBuf::from(db_path)),
embedding_model: None,
additional_embedding_models: vec![],
..Default::default()
}
}
#[tokio::test]
async fn atomic_ops_file_success_commits_all_ops() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (x_id, y_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="AtomicX"), create(kind="concept", name="AtomicY")]"#,
)
.await;
let x_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("x id")
.to_string();
let y_id = resp["results"][1]["result"]["id"]
.as_str()
.expect("y id")
.to_string();
(x_id, y_id)
};
let ops = vec![
atomic_op(
"update",
serde_json::json!({"id": x_id, "name": "AtomicX-renamed"}),
),
atomic_op(
"update",
serde_json::json!({"id": y_id, "name": "AtomicY-renamed"}),
),
];
let khive_cfg = KhiveConfig::default();
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("atomic run must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let server = isolated_server(&db_path);
let x_resp = dispatch_json(&server, &format!(r#"get(id="{x_id}")"#)).await;
let y_resp = dispatch_json(&server, &format!(r#"get(id="{y_id}")"#)).await;
assert_eq!(x_resp["results"][0]["result"]["name"], "AtomicX-renamed");
assert_eq!(y_resp["results"][0]["result"]["name"], "AtomicY-renamed");
}
#[tokio::test]
async fn atomic_ops_file_mid_unit_failure_rolls_back_whole_unit() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (x_id, y_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="RollbackX"), create(kind="concept", name="RollbackY")]"#,
)
.await;
let x_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("x id")
.to_string();
let y_id = resp["results"][1]["result"]["id"]
.as_str()
.expect("y id")
.to_string();
(x_id, y_id)
};
let ops = vec![
atomic_op("delete", serde_json::json!({"id": x_id, "hard": true})),
atomic_op(
"link",
serde_json::json!({
"source_id": y_id,
"target_id": x_id,
"relation": "extends",
}),
),
];
let khive_cfg = KhiveConfig::default();
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("the seam call itself must not error — the unit rolls back cleanly");
assert_eq!(
envelope["atomic"]["rolled_back"], true,
"envelope: {envelope}"
);
assert_eq!(
envelope["atomic"]["failed_op_index"], 1,
"envelope: {envelope}"
);
let server = isolated_server(&db_path);
let x_resp = dispatch_json(&server, &format!(r#"get(id="{x_id}")"#)).await;
assert!(
x_resp["results"][0]["result"]["deleted_at"].is_null(),
"x must NOT be deleted — the whole unit must have rolled back: {x_resp}"
);
}
#[tokio::test]
async fn atomic_symmetric_update_absorbs_into_same_unit_link_and_renders_correct_id() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (a_id, b_id, x_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="LinkRaceA"), create(kind="concept", name="LinkRaceB")]"#,
)
.await;
let a_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("a id")
.to_string();
let b_id = resp["results"][1]["result"]["id"]
.as_str()
.expect("b id")
.to_string();
let link_resp = dispatch_json(
&server,
&format!(
r#"link(source_id="{a_id}", target_id="{b_id}", relation="extends", weight=0.2)"#
),
)
.await;
let x_id = link_resp["results"][0]["result"]["id"]
.as_str()
.expect("x id")
.to_string();
(a_id, b_id, x_id)
};
let ops = vec![
atomic_op(
"link",
serde_json::json!({
"source_id": a_id,
"target_id": b_id,
"relation": "competes_with",
"weight": 0.6,
}),
),
atomic_op(
"update",
serde_json::json!({"id": x_id, "relation": "competes_with", "weight": 0.9}),
),
];
let khive_cfg = KhiveConfig::default();
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("atomic run must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let linked_id = envelope["results"][0]["result"]["id"]
.as_str()
.expect("link result id")
.to_string();
let rendered_update_id = envelope["results"][1]["result"]["id"]
.as_str()
.expect("update result id")
.to_string();
assert_ne!(
rendered_update_id, x_id,
"the update's rendered result must NOT be X's stale requested id: {envelope}"
);
assert_eq!(
rendered_update_id, linked_id,
"the update's rendered result must be the surviving (just-linked) row: {envelope}"
);
assert_eq!(
envelope["results"][1]["result"]["weight"], 0.9,
"the surviving row must carry the update's patch: {envelope}"
);
let server = isolated_server(&db_path);
let surviving_resp = dispatch_json(&server, &format!(r#"get(id="{linked_id}")"#)).await;
assert_eq!(
surviving_resp["results"][0]["result"]["weight"], 0.9,
"the committed row itself must carry the patch: {surviving_resp}"
);
}
#[tokio::test]
async fn atomic_cli_boundary_rejections_happen_before_any_write() {
let khive_cfg = KhiveConfig::default();
{
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![atomic_op(
"create",
serde_json::json!({"kind": "concept", "name": "ShouldNotLand"}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("embedding-bearing verb must be rejected");
assert!(
format!("{err:#}").contains("embedding-bearing"),
"error: {err:#}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, r#"list(kind="entity")"#).await;
assert_eq!(resp["results"][0]["result"].as_array().unwrap().len(), 0);
}
{
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![atomic_op("search", serde_json::json!({"query": "x"}))];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("read verbs must be rejected");
assert!(format!("{err:#}").contains("read"), "error: {err:#}");
}
{
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![atomic_op("not_a_real_verb", serde_json::json!({}))];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("unlisted verbs must be rejected");
assert!(
format!("{err:#}").contains("not on the v1 atomic-admissible"),
"error: {err:#}"
);
}
{
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![
atomic_op(
"update",
serde_json::json!({"id": uuid::Uuid::new_v4().to_string()}),
),
atomic_op(
"update",
serde_json::json!({"id": uuid::Uuid::new_v4().to_string()}),
),
atomic_op(
"update",
serde_json::json!({"id": uuid::Uuid::new_v4().to_string()}),
),
];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
2,
)
.await
.expect_err("exceeding max_ops must be rejected");
assert!(
format!("{err:#}").contains("exceeds the configured maximum"),
"error: {err:#}"
);
}
for verb in ["propose", "review", "withdraw"] {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![atomic_op(
verb,
serde_json::json!({"title": "x", "description": "y", "changeset": {}}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err(&format!("{verb:?} must be rejected before any write"));
assert!(
format!("{err:#}").contains("no --atomic prepare/apply seam"),
"error for {verb:?}: {err:#}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, r#"list(kind="entity")"#).await;
assert_eq!(
resp["results"][0]["result"].as_array().unwrap().len(),
0,
"no write must have landed for {verb:?}"
);
}
{
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let ops = vec![atomic_op(
"merge",
serde_json::json!({
"into_id": uuid::Uuid::new_v4().to_string(),
"from_id": uuid::Uuid::new_v4().to_string(),
}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("merge must be rejected before any write");
assert!(
format!("{err:#}").contains("use the non-atomic merge verb instead"),
"error: {err:#}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, r#"list(kind="entity")"#).await;
assert_eq!(
resp["results"][0]["result"].as_array().unwrap().len(),
0,
"no write must have landed for merge"
);
}
}
#[tokio::test]
async fn atomic_update_entity_unknown_field_is_rejected_and_does_not_mutate_row() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (entity_id, updated_at_before) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"create(kind="concept", name="TypoGuardX", description="original")"#,
)
.await;
let id = resp["results"][0]["result"]["id"]
.as_str()
.expect("id")
.to_string();
let get_resp = dispatch_json(&server, &format!(r#"get(id="{id}")"#)).await;
let updated_at = get_resp["results"][0]["result"]["updated_at"].clone();
(id, updated_at)
};
let ops = vec![atomic_op(
"update",
serde_json::json!({"id": entity_id, "conten": "hello"}),
)];
let khive_cfg = KhiveConfig::default();
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("typo'd `conten` must be rejected, not silently dropped");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let server = isolated_server(&db_path);
let get_resp = dispatch_json(&server, &format!(r#"get(id="{entity_id}")"#)).await;
assert_eq!(
get_resp["results"][0]["result"]["description"], "original",
"a rejected op must not have mutated description: {get_resp}"
);
assert_eq!(
get_resp["results"][0]["result"]["updated_at"], updated_at_before,
"a rejected op must not bump updated_at (no write happened): {get_resp}"
);
}
#[tokio::test]
async fn atomic_update_note_unknown_field_rejected_well_formed_succeeds() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let note_id = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"create(kind="observation", content="original note")"#,
)
.await;
resp["results"][0]["result"]["id"]
.as_str()
.expect("id")
.to_string()
};
let khive_cfg = KhiveConfig::default();
let ops = vec![atomic_op(
"update",
serde_json::json!({"id": note_id, "conten": "typo'd"}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("typo'd `conten` on a note update must be rejected");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let ops = vec![atomic_op(
"update",
serde_json::json!({"id": note_id, "content": "updated note"}),
)];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("a well-formed note update must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let server = isolated_server(&db_path);
let get_resp = dispatch_json(&server, &format!(r#"get(id="{note_id}")"#)).await;
assert_eq!(
get_resp["results"][0]["result"]["content"], "updated note",
"the well-formed update must have landed: {get_resp}"
);
}
#[tokio::test]
async fn atomic_delete_unknown_field_rejected_well_formed_succeeds() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let entity_id = {
let server = isolated_server(&db_path);
let resp =
dispatch_json(&server, r#"create(kind="concept", name="DeleteTypoGuard")"#).await;
resp["results"][0]["result"]["id"]
.as_str()
.expect("id")
.to_string()
};
let khive_cfg = KhiveConfig::default();
let ops = vec![atomic_op(
"delete",
serde_json::json!({"id": entity_id, "hardd": true}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("typo'd `hardd` must be rejected");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let server = isolated_server(&db_path);
let get_resp = dispatch_json(&server, &format!(r#"get(id="{entity_id}")"#)).await;
assert!(
get_resp["results"][0]["result"]["deleted_at"].is_null(),
"a rejected delete must not have deleted the entity: {get_resp}"
);
let ops = vec![atomic_op("delete", serde_json::json!({"id": entity_id}))];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("a well-formed delete must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
}
#[tokio::test]
async fn atomic_link_unknown_field_rejected_well_formed_succeeds() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (a_id, b_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="LinkTypoA"), create(kind="concept", name="LinkTypoB")]"#,
)
.await;
let a_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("a id")
.to_string();
let b_id = resp["results"][1]["result"]["id"]
.as_str()
.expect("b id")
.to_string();
(a_id, b_id)
};
let khive_cfg = KhiveConfig::default();
let ops = vec![atomic_op(
"link",
serde_json::json!({
"source_id": a_id,
"target_id": b_id,
"relation": "extends",
"relatoin": "extends",
}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("typo'd `relatoin` must be rejected");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let ops = vec![atomic_op(
"link",
serde_json::json!({"source_id": a_id, "target_id": b_id, "relation": "extends"}),
)];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("a well-formed link must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
}
#[tokio::test]
async fn atomic_delete_rejects_kind_mismatch_and_accepts_matching_or_omitted_kind() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let khive_cfg = KhiveConfig::default();
let (mismatch_id, matching_id, omitted_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="KindMismatch"), create(kind="concept", name="KindMatching"), create(kind="concept", name="KindOmitted")]"#,
)
.await;
let id = |i: usize| {
resp["results"][i]["result"]["id"]
.as_str()
.expect("id")
.to_string()
};
(id(0), id(1), id(2))
};
let ops = vec![atomic_op(
"delete",
serde_json::json!({"id": mismatch_id, "kind": "note"}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("delete(kind=\"note\") on an entity must be rejected");
assert!(
format!("{err:#}").contains("not found"),
"expected a NotFound-shaped rejection, error: {err:#}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, &format!(r#"get(id="{mismatch_id}")"#)).await;
assert!(
resp["results"][0]["result"]["deleted_at"].is_null(),
"entity must NOT be deleted after a kind-mismatch rejection: {resp}"
);
let ops = vec![atomic_op(
"delete",
serde_json::json!({"id": matching_id, "kind": "entity"}),
)];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("delete(kind=\"entity\") on an entity must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let ops = vec![atomic_op("delete", serde_json::json!({"id": omitted_id}))];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("delete with kind omitted must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
}
#[tokio::test]
async fn atomic_update_null_and_type_semantics_match_canonical_no_op_behavior() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let khive_cfg = KhiveConfig::default();
let entity_id = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"create(kind="concept", name="NullSemantics", description="orig-desc", properties={"k": "v"}, tags=["a", "b"])"#,
)
.await;
resp["results"][0]["result"]["id"]
.as_str()
.expect("id")
.to_string()
};
let ops = vec![atomic_op(
"update",
serde_json::json!({"id": entity_id, "name": 123}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("name: 123 (non-null, non-string) must be rejected");
assert!(
format!("{err:#}").contains("name must be a string"),
"error: {err:#}"
);
let ops = vec![atomic_op(
"update",
serde_json::json!({
"id": entity_id,
"name": null,
"description": null,
"properties": null,
"tags": null,
}),
)];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("an all-null update must be a no-op success, not a rejection");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, &format!(r#"get(id="{entity_id}")"#)).await;
let row = &resp["results"][0]["result"];
assert_eq!(
row["name"], "NullSemantics",
"name must be unchanged: {row}"
);
assert_eq!(
row["description"], "orig-desc",
"description must be unchanged: {row}"
);
assert_eq!(
row["properties"]["k"], "v",
"properties must be unchanged: {row}"
);
assert_eq!(
row["tags"],
serde_json::json!(["a", "b"]),
"tags must be unchanged: {row}"
);
}
#[tokio::test]
async fn atomic_update_and_delete_accept_8_hex_prefix_ids() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let khive_cfg = KhiveConfig::default();
let (entity_full_id, doomed_full_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="PrefixEntity"), create(kind="concept", name="PrefixDoomed")]"#,
)
.await;
let entity_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("entity id")
.to_string();
let doomed_id = resp["results"][1]["result"]["id"]
.as_str()
.expect("doomed id")
.to_string();
(entity_id, doomed_id)
};
let entity_prefix = &entity_full_id[..8];
let doomed_prefix = &doomed_full_id[..8];
let ops = vec![
atomic_op(
"update",
serde_json::json!({"id": entity_prefix, "name": "PrefixEntity-renamed"}),
),
atomic_op("delete", serde_json::json!({"id": doomed_prefix})),
];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("8-hex-prefix ids must resolve identically to canonical");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let server = isolated_server(&db_path);
let resp = dispatch_json(&server, &format!(r#"get(id="{entity_full_id}")"#)).await;
assert_eq!(
resp["results"][0]["result"]["name"], "PrefixEntity-renamed",
"prefix-addressed update must have landed: {resp}"
);
let ops = vec![atomic_op(
"update",
serde_json::json!({"id": "deadbeef", "name": "should not resolve"}),
)];
let err = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("a non-existent prefix must be rejected");
assert!(
format!("{err:#}").contains("no record matches prefix"),
"error: {err:#}"
);
}
#[tokio::test]
async fn atomic_success_results_carry_canonical_shaped_result_per_op() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let khive_cfg = KhiveConfig::default();
let (entity_id, doomed_id, source_id, target_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="ResultUpdate"), create(kind="concept", name="ResultDelete"), create(kind="concept", name="ResultLinkSource"), create(kind="concept", name="ResultLinkTarget")]"#,
)
.await;
let id = |i: usize| {
resp["results"][i]["result"]["id"]
.as_str()
.expect("id")
.to_string()
};
(id(0), id(1), id(2), id(3))
};
let ops = vec![
atomic_op(
"update",
serde_json::json!({"id": entity_id, "name": "ResultUpdate-renamed"}),
),
atomic_op("delete", serde_json::json!({"id": doomed_id})),
atomic_op(
"link",
serde_json::json!({
"source_id": source_id,
"target_id": target_id,
"relation": "extends",
}),
),
];
let envelope = crate::atomic_apply::execute_atomic_ops_file(
ops,
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("three v1-admissible verbs must commit as one unit");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let results = envelope["results"].as_array().expect("results array");
assert_eq!(results.len(), 3, "envelope: {envelope}");
assert_eq!(
results[0]["result"]["name"], "ResultUpdate-renamed",
"update result must carry the updated name: {envelope}"
);
assert_eq!(
results[1]["result"]["deleted"], true,
"delete result: {envelope}"
);
assert_eq!(
results[1]["result"]["id"], doomed_id,
"delete result must echo the caller's id: {envelope}"
);
assert_eq!(
results[2]["result"]["relation"], "extends",
"link result must carry the edge's relation: {envelope}"
);
assert_eq!(
results[2]["result"]["source_id"], source_id,
"link result must carry source_id: {envelope}"
);
assert_eq!(
results[2]["result"]["target_id"], target_id,
"link result must carry target_id: {envelope}"
);
}
}