use std::collections::BTreeMap;
use std::io::{BufRead as _, Read as _, Seek as _, Write as _};
use std::path::{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,
build_single_backend_runtime, config_discovery_db_anchor, enforce_strict_actor_mode,
normalize_redundant_db_override_with_source, reject_conflicting_db_override_with_source,
validate_declared_backend_access_modes, RuntimeConfigInputs,
};
use khive_mcp::server::KhiveMcpServer;
#[cfg(unix)]
use khive_mcp::server::{compute_config_id, compute_config_id_with_storage_mode};
use khive_mcp::tools::request::RequestParams;
#[cfg(test)]
use khive_runtime::KhiveRuntime;
#[cfg(unix)]
use khive_runtime::{daemon::PROTOCOL_VERSION, DaemonRequestFrame};
use khive_runtime::{KhiveConfig, Namespace, RuntimeConfig};
use khive_types::RefusalReason;
mod plan;
const REFUSAL_PREFIX: &str = "kkernel-refusal: ";
#[derive(Debug)]
struct ExecRefusal {
reason: RefusalReason,
message: String,
}
impl std::fmt::Display for ExecRefusal {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.message)
}
}
impl std::error::Error for ExecRefusal {}
fn refusal_error(reason: RefusalReason, message: impl Into<String>) -> anyhow::Error {
anyhow::Error::new(ExecRefusal {
reason,
message: message.into(),
})
}
fn emit_refusal(reason: RefusalReason) {
eprintln!("{REFUSAL_PREFIX}{reason}");
}
fn refusal_envelope_for_tools(
tools: Vec<String>,
chain: bool,
reason: RefusalReason,
message: &str,
) -> serde_json::Value {
debug_assert!(
!tools.is_empty(),
"per-operation refusal envelopes require at least one parsed operation"
);
let total = tools.len();
let results: Vec<serde_json::Value> = tools
.into_iter()
.enumerate()
.map(|(index, tool)| {
if chain && index > 0 {
serde_json::json!({
"ok": false,
"tool": tool,
"aborted": true,
"message": message,
"reason": reason.as_str(),
})
} else {
serde_json::json!({
"ok": false,
"tool": tool,
"error": message,
"reason": reason.as_str(),
})
}
})
.collect();
let aborted = if chain { total.saturating_sub(1) } else { 0 };
let failed = total - aborted;
serde_json::json!({
"results": results,
"summary": {
"total": total,
"succeeded": 0,
"failed": failed,
"aborted": aborted,
},
"status": "partial",
})
}
fn invocation_refusal_envelope(reason: RefusalReason, message: &str) -> serde_json::Value {
let code = if reason == RefusalReason::ParseError {
"invalid_params"
} else {
"invocation_refused"
};
serde_json::json!({
"error": {
"code": code,
"message": message,
"reason": reason.as_str(),
},
"invocation": {"started": false},
})
}
fn report_unscoped_refusal(reason: RefusalReason, message: impl Into<String>) -> anyhow::Error {
let message = message.into();
emit_refusal(reason);
println!(
"{}",
serde_json::to_string(&invocation_refusal_envelope(reason, &message))
.expect("invocation refusal envelope is serializable")
);
anyhow::anyhow!(message)
}
fn report_invocation_refusal(
raw_ops: Option<&str>,
reason: RefusalReason,
error: impl std::fmt::Display,
) -> anyhow::Error {
let message = error.to_string();
let (tools, chain) = raw_ops
.and_then(|ops| khive_request::parse_request(ops).ok())
.map(|parsed| {
let chain = parsed.mode == khive_request::ExecutionMode::Chain;
let tools: Vec<String> = parsed.ops.into_iter().map(|op| op.tool).collect();
(tools, chain)
})
.unwrap_or_default();
if tools.is_empty() {
report_unscoped_refusal(reason, message)
} else {
report_tools_refusal(tools, chain, reason, message)
}
}
fn report_tools_refusal(
tools: Vec<String>,
chain: bool,
reason: RefusalReason,
message: impl Into<String>,
) -> anyhow::Error {
let message = message.into();
if tools.is_empty() {
return report_unscoped_refusal(reason, message);
}
emit_refusal(reason);
let envelope = refusal_envelope_for_tools(tools, chain, reason, &message);
println!(
"{}",
serde_json::to_string(&envelope).expect("refusal envelope is serializable")
);
anyhow::anyhow!(message)
}
#[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,
Option<PathBuf>,
Option<&'a str>,
Vec<String>,
) -> ForwardFuture<'a>;
#[cfg(unix)]
fn forward_or_spawn_boxed<'a>(
frame: &'a DaemonRequestFrame,
config: Option<PathBuf>,
db: Option<&'a str>,
packs: Vec<String>,
) -> ForwardFuture<'a> {
Box::pin(async move {
khive_mcp::daemon::forward_or_spawn_with_config_and_packs(
frame,
config.as_deref(),
db,
Some(&packs),
)
.await
})
}
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;
const OPS_FILE_CHUNK_MAX_BYTES: usize = 32 * 1024 * 1024;
const MAX_OPS_FILE_LINE_BYTES: usize = 96 * 1024 * 1024;
const MAX_OPS_FILE_BYTES: u64 = 512 * 1024 * 1024;
const MAX_OPS_FILE_FAILURE_DETAILS: usize = 1_000;
const MAX_OPS_FILE_FAILURE_ERROR_BYTES: usize = 4 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OpsFileDispatchMode {
BoundedParallel,
Serial,
}
impl OpsFileDispatchMode {
const fn is_serial(self) -> bool {
matches!(self, Self::Serial)
}
}
#[derive(Parser, Debug)]
pub struct ExecArgs {
pub ops: Option<String>,
#[arg(
long,
requires = "ops",
conflicts_with_all = [
"presentation", "strict", "output_format", "save_file", "ops_file",
"pending_events", "dry_run", "serial", "atomic", "atomic_max_ops",
"verbose", "actor", "expect_actor", "namespace"
]
)]
pub plan: bool,
#[arg(long, conflicts_with = "ops", conflicts_with = "ops_file")]
pub pending_events: bool,
#[arg(long, env = "KHIVE_DB")]
pub db: Option<String>,
#[arg(long, env = "KHIVE_CONFIG")]
pub config: Option<PathBuf>,
#[arg(long, default_value = "local")]
pub namespace: String,
#[arg(long, value_name = "ACTOR", conflicts_with = "pending_events")]
pub actor: Option<String>,
#[arg(long, value_name = "ACTOR", conflicts_with = "pending_events")]
pub expect_actor: Option<String>,
#[arg(
long,
default_value = "verbose",
default_value_if("plan", "true", None)
)]
pub presentation: Option<String>,
#[arg(long, value_name = "FORMAT")]
pub output_format: Option<String>,
#[arg(long, short = 'v')]
pub verbose: bool,
#[arg(long, conflicts_with = "dry_run")]
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",
conflicts_with_all = ["atomic", "ops"]
)]
pub serial: bool,
#[arg(long, requires = "ops_file")]
pub atomic: bool,
#[arg(long, requires = "atomic")]
pub atomic_max_ops: Option<usize>,
#[arg(long)]
pub strict: bool,
}
#[derive(Debug, Clone, serde::Serialize)]
pub(crate) struct OpsFileEntry {
pub(crate) tool: String,
pub(crate) args: serde_json::Value,
}
#[derive(Debug)]
struct ValidatedOpsFile {
snapshot: std::fs::File,
total: usize,
per_verb: BTreeMap<String, usize>,
}
fn parse_ops_file_line(raw: &str, line_num: usize) -> Result<Option<OpsFileEntry>> {
let trimmed = raw.trim();
if trimmed.is_empty() {
return Ok(None);
}
let obj: serde_json::Value = serde_json::from_str(trimmed).map_err(|error| {
refusal_error(
RefusalReason::ParseError,
format!("ops-file line {line_num}: invalid JSON: {error}"),
)
})?;
let obj = obj.as_object().ok_or_else(|| {
refusal_error(
RefusalReason::ParseError,
format!(
"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(|| {
refusal_error(
RefusalReason::ParseError,
format!("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() => v.clone(),
Some(v) => {
return Err(refusal_error(
RefusalReason::ParseError,
format!("ops-file line {line_num}: \"args\" must be a JSON object, got {v}"),
))
}
};
Ok(Some(OpsFileEntry { tool, args }))
}
fn read_bounded_ops_line<R: std::io::BufRead>(
reader: &mut R,
line_num: usize,
) -> Result<Option<String>> {
read_bounded_ops_line_with_limit(reader, line_num, MAX_OPS_FILE_LINE_BYTES)
}
fn read_bounded_ops_line_with_limit<R: std::io::BufRead>(
reader: &mut R,
line_num: usize,
limit: usize,
) -> Result<Option<String>> {
let mut bytes = Vec::new();
let read = {
let mut limited = (&mut *reader).take((limit + 1) as u64);
limited
.read_until(b'\n', &mut bytes)
.with_context(|| format!("read ops-file line {line_num}"))?
};
if read == 0 {
return Ok(None);
}
if read > limit {
anyhow::bail!("ops-file line {line_num} exceeds the {limit}-byte physical-line limit");
}
if bytes.last() == Some(&b'\n') {
bytes.pop();
if bytes.last() == Some(&b'\r') {
bytes.pop();
}
}
String::from_utf8(bytes)
.map(Some)
.map_err(|error| anyhow::anyhow!("ops-file line {line_num} is not valid UTF-8: {error}"))
}
fn validate_ops_file(path: &Path) -> Result<ValidatedOpsFile> {
let file =
std::fs::File::open(path).with_context(|| format!("open ops-file {}", path.display()))?;
let metadata_len = file
.metadata()
.with_context(|| format!("stat ops-file {}", path.display()))?
.len();
if metadata_len > MAX_OPS_FILE_BYTES {
anyhow::bail!(
"ops-file {} is {metadata_len} bytes, exceeding the {MAX_OPS_FILE_BYTES}-byte total limit",
path.display()
);
}
let mut reader = std::io::BufReader::new(file);
let mut snapshot = tempfile::tempfile().context("create validated ops-file snapshot")?;
let mut total = 0_usize;
let mut total_bytes = 0_u64;
let mut per_verb = BTreeMap::new();
let mut line_num = 1_usize;
while let Some(raw) = read_bounded_ops_line(&mut reader, line_num)? {
total_bytes = total_bytes
.checked_add(raw.len() as u64)
.and_then(|value| value.checked_add(1))
.ok_or_else(|| anyhow::anyhow!("ops-file byte count overflow"))?;
if total_bytes > MAX_OPS_FILE_BYTES {
anyhow::bail!(
"ops-file exceeds the {MAX_OPS_FILE_BYTES}-byte total limit while reading line {line_num}"
);
}
if let Some(op) = parse_ops_file_line(&raw, line_num)? {
*per_verb.entry(op.tool).or_insert(0) += 1;
snapshot
.write_all(raw.trim().as_bytes())
.context("write validated ops-file snapshot")?;
snapshot
.write_all(b"\n")
.context("write validated ops-file snapshot newline")?;
total += 1;
}
line_num += 1;
}
snapshot
.rewind()
.context("rewind validated ops-file snapshot")?;
Ok(ValidatedOpsFile {
snapshot,
total,
per_verb,
})
}
fn parse_validated_snapshot<R>(snapshot: &mut R) -> Result<Vec<OpsFileEntry>>
where
R: std::io::Read + std::io::Seek,
{
snapshot
.rewind()
.context("rewind validated ops-file snapshot")?;
let mut reader = std::io::BufReader::new(snapshot);
let mut ops = Vec::new();
let mut line_num = 1_usize;
while let Some(raw) = read_bounded_ops_line(&mut reader, line_num)? {
if let Some(op) = parse_ops_file_line(&raw, line_num)? {
ops.push(op);
}
line_num += 1;
}
Ok(ops)
}
fn preflight_typed_validated_snapshot<R>(snapshot: &mut R, expected_total: usize) -> Result<()>
where
R: std::io::Read + std::io::Seek,
{
snapshot
.rewind()
.context("rewind validated ops-file snapshot for typed preflight")?;
let result: Result<()> = (|| {
let mut reader = std::io::BufReader::new(&mut *snapshot);
let mut processed = 0_usize;
let mut line_number = 1_usize;
let mut chunk_number = 0_usize;
let mut pending: Option<(OpsFileEntry, usize)> = None;
let mut eof = false;
while !eof {
let mut chunk = Vec::with_capacity(OPS_FILE_CHUNK_SIZE);
let mut chunk_bytes = 0_usize;
if let Some((op, bytes)) = pending.take() {
chunk_bytes = bytes;
chunk.push(op);
}
while chunk.len() < OPS_FILE_CHUNK_SIZE {
let Some(raw) = read_bounded_ops_line(&mut reader, line_number)? else {
eof = true;
break;
};
let physical_bytes = raw.len().saturating_add(1);
if let Some(op) = parse_ops_file_line(&raw, line_number)? {
if should_defer_chunk_entry(chunk.len(), chunk_bytes, physical_bytes) {
pending = Some((op, physical_bytes));
line_number += 1;
break;
}
chunk_bytes = chunk_bytes.saturating_add(physical_bytes);
chunk.push(op);
}
line_number += 1;
}
if chunk.is_empty() {
break;
}
chunk_number += 1;
let chunk_len = chunk.len();
let typed_ops = chunk
.into_iter()
.map(|op| {
let serde_json::Value::Object(args) = op.args else {
unreachable!("validated ops-file args are always JSON objects")
};
khive_request::TypedJsonOp {
tool: op.tool,
args,
}
})
.collect();
khive_request::parse_typed_json_batch(typed_ops).map_err(|error| {
refusal_error(
RefusalReason::ParseError,
format!("ops-file typed preflight chunk {chunk_number}: {error}"),
)
})?;
processed += chunk_len;
}
if processed != expected_total {
anyhow::bail!(
"validated ops-file snapshot changed during typed preflight: expected \
{expected_total} ops, read {processed}"
);
}
Ok(())
})();
snapshot
.rewind()
.context("rewind typed-preflighted ops-file snapshot")?;
result
}
fn validated_tool_names<R>(snapshot: &mut R) -> Result<Vec<String>>
where
R: std::io::Read + std::io::Seek,
{
snapshot
.rewind()
.context("rewind validated ops-file snapshot")?;
let mut reader = std::io::BufReader::new(snapshot);
let mut tools = Vec::new();
let mut line_num = 1_usize;
while let Some(raw) = read_bounded_ops_line(&mut reader, line_num)? {
if let Some(op) = parse_ops_file_line(&raw, line_num)? {
tools.push(op.tool);
}
line_num += 1;
}
Ok(tools)
}
fn parse_atomic_validated_snapshot<R>(
snapshot: &mut R,
total: usize,
max_ops: usize,
) -> Result<Vec<OpsFileEntry>>
where
R: std::io::Read + std::io::Seek,
{
if total > max_ops {
anyhow::bail!(
"--atomic op count {total} exceeds the configured maximum {max_ops}; \
split the file or raise --atomic-max-ops"
);
}
parse_validated_snapshot(snapshot)
}
#[cfg(test)]
pub(crate) fn parse_ops_file(path: &Path) -> Result<Vec<OpsFileEntry>> {
let mut validated = validate_ops_file(path)?;
parse_validated_snapshot(&mut validated.snapshot)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum OpsFileReportMode {
LegacyNoSave,
BoundedSave,
}
fn collect_op_failures(
parsed: &serde_json::Value,
applied_before: usize,
mode: OpsFileReportMode,
) -> Vec<serde_json::Value> {
let Some(results) = parsed["results"].as_array() else {
return Vec::new();
};
results
.iter()
.enumerate()
.filter(|(_, entry)| entry["ok"].as_bool() == Some(false))
.map(|(i, entry)| {
let error = match &entry["error"] {
serde_json::Value::Null => serde_json::Value::from("unknown error"),
other if mode == OpsFileReportMode::BoundedSave => bounded_failure_error(other),
other => other.clone(),
};
let mut failure = serde_json::json!({
"op_index": applied_before + i,
"tool": entry["tool"].as_str().unwrap_or("?"),
"error": error,
});
if mode == OpsFileReportMode::BoundedSave {
failure["aborted"] =
serde_json::Value::Bool(entry["aborted"].as_bool().unwrap_or(false));
if let Some(reason) = entry["reason"].as_str().and_then(RefusalReason::from_token) {
failure["reason"] = serde_json::json!(reason.as_str());
}
}
failure
})
.collect()
}
fn retain_failure_detail(
mode: OpsFileReportMode,
failure: serde_json::Value,
failures: &mut Vec<serde_json::Value>,
omitted: &mut usize,
) -> bool {
if mode == OpsFileReportMode::BoundedSave && failures.len() >= MAX_OPS_FILE_FAILURE_DETAILS {
*omitted += 1;
false
} else {
failures.push(failure);
true
}
}
fn ops_file_progress_line(
mode: OpsFileReportMode,
applied: usize,
total: usize,
succeeded: usize,
failed: usize,
aborted: usize,
) -> String {
match mode {
OpsFileReportMode::LegacyNoSave => {
format!("applied {applied}/{total} (ok={succeeded}, failed={failed})")
}
OpsFileReportMode::BoundedSave => format!(
"applied {applied}/{total} (ok={succeeded}, failed={failed}, aborted={aborted})"
),
}
}
#[allow(clippy::too_many_arguments)]
fn ops_file_summary(
mode: OpsFileReportMode,
total: usize,
succeeded: usize,
failed: usize,
aborted: usize,
failures: Vec<serde_json::Value>,
failure_details_omitted: usize,
) -> serde_json::Value {
let mut summary = match mode {
OpsFileReportMode::LegacyNoSave => serde_json::json!({
"total": total,
"succeeded": succeeded,
"failed": failed,
}),
OpsFileReportMode::BoundedSave => serde_json::json!({
"total": total,
"succeeded": succeeded,
"failed": failed,
"aborted": aborted,
}),
};
if !failures.is_empty() {
summary["failures"] = serde_json::Value::Array(failures);
}
if mode == OpsFileReportMode::BoundedSave && failure_details_omitted > 0 {
summary["failure_details_omitted"] = serde_json::json!(failure_details_omitted);
}
summary
}
fn bounded_failure_error(error: &serde_json::Value) -> serde_json::Value {
let mut writer = CountingWriter::default();
if serde_json::to_writer(&mut writer, error).is_ok()
&& writer.bytes <= MAX_OPS_FILE_FAILURE_ERROR_BYTES
{
error.clone()
} else {
serde_json::Value::String(format!(
"error detail omitted: exceeds {MAX_OPS_FILE_FAILURE_ERROR_BYTES}-byte ops-file diagnostic limit"
))
}
}
fn required_summary_count(parsed: &serde_json::Value, field: &str) -> Result<usize> {
let value = parsed
.pointer(&format!("/summary/{field}"))
.and_then(serde_json::Value::as_u64)
.ok_or_else(|| anyhow::anyhow!("dispatch result is missing integer summary.{field}"))?;
usize::try_from(value).context("dispatch summary count does not fit usize")
}
fn classify_ordered_chunk(
expected_tools: &[String],
results: &[serde_json::Value],
) -> Result<(usize, usize, usize)> {
if results.len() != expected_tools.len() {
anyhow::bail!(
"ordered chunk result count {} does not match input count {}",
results.len(),
expected_tools.len()
);
}
let mut succeeded = 0_usize;
let mut failed = 0_usize;
let mut aborted = 0_usize;
for (index, (expected_tool, row)) in expected_tools.iter().zip(results).enumerate() {
let object = row
.as_object()
.ok_or_else(|| anyhow::anyhow!("dispatch result row {index} is not a JSON object"))?;
let returned_tool = object
.get("tool")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| anyhow::anyhow!("dispatch result row {index} has no string tool"))?;
if returned_tool != expected_tool {
anyhow::bail!(
"dispatch result row {index} tool mismatch: expected {:?}, got {:?}",
expected_tool,
returned_tool
);
}
let ok = object
.get("ok")
.and_then(serde_json::Value::as_bool)
.ok_or_else(|| anyhow::anyhow!("dispatch result row {index} has no boolean ok"))?;
let row_aborted = match object.get("aborted") {
None => false,
Some(value) => value.as_bool().ok_or_else(|| {
anyhow::anyhow!("dispatch result row {index} has non-boolean aborted")
})?,
};
match (ok, row_aborted) {
(true, false) => {
if !object.contains_key("result") {
anyhow::bail!(
"dispatch result row {index} is successful but has no result field"
);
}
if object.contains_key("error") {
anyhow::bail!(
"dispatch result row {index} is successful but also has an error field"
);
}
succeeded += 1;
}
(false, false) => {
if !object.contains_key("error") {
anyhow::bail!("dispatch result row {index} failed but has no error field");
}
if object.contains_key("result") {
anyhow::bail!("dispatch result row {index} failed but also has a result field");
}
failed += 1;
}
(false, true) => {
if !object.contains_key("error") {
anyhow::bail!("dispatch result row {index} aborted but has no error field");
}
if object.contains_key("result") {
anyhow::bail!(
"dispatch result row {index} aborted but also has a result field"
);
}
aborted += 1;
}
(true, true) => {
anyhow::bail!("dispatch result row {index} cannot be both successful and aborted")
}
}
}
Ok((succeeded, failed, aborted))
}
fn validate_ordered_chunk_envelope(
expected_tools: &[String],
parsed: &serde_json::Value,
chunk_number: usize,
) -> Result<(usize, usize, usize)> {
let results = parsed
.get("results")
.and_then(serde_json::Value::as_array)
.ok_or_else(|| {
anyhow::anyhow!("dispatch chunk {chunk_number} returned no results array")
})?;
let chunk_total = required_summary_count(parsed, "total")?;
let chunk_succeeded = required_summary_count(parsed, "succeeded")?;
let chunk_failed = required_summary_count(parsed, "failed")?;
let chunk_aborted = required_summary_count(parsed, "aborted")?;
let (derived_succeeded, derived_failed, derived_aborted) =
classify_ordered_chunk(expected_tools, results)?;
if chunk_total != expected_tools.len()
|| chunk_succeeded != derived_succeeded
|| chunk_failed != derived_failed
|| chunk_aborted != derived_aborted
{
anyhow::bail!(
"dispatch chunk {chunk_number} summary disagrees with ordered rows: expected total {}, summary total {}, derived/summary succeeded {derived_succeeded}/{chunk_succeeded}, failed {derived_failed}/{chunk_failed}, aborted {derived_aborted}/{chunk_aborted}",
expected_tools.len(),
chunk_total,
);
}
let status = parsed
.get("status")
.and_then(serde_json::Value::as_str)
.ok_or_else(|| {
anyhow::anyhow!("dispatch chunk {chunk_number} returned no string status")
})?;
let expected_status = if derived_failed == 0 && derived_aborted == 0 {
"success"
} else {
"partial"
};
if status != expected_status {
anyhow::bail!(
"dispatch chunk {chunk_number} status disagrees with ordered rows: expected {expected_status:?}, got {status:?}"
);
}
Ok((chunk_succeeded, chunk_failed, chunk_aborted))
}
fn should_defer_chunk_entry(current_count: usize, current_bytes: usize, next_bytes: usize) -> bool {
current_count > 0
&& (current_count >= OPS_FILE_CHUNK_SIZE
|| current_bytes.saturating_add(next_bytes) > OPS_FILE_CHUNK_MAX_BYTES)
}
#[derive(Default)]
struct CountingWriter {
bytes: usize,
}
impl std::io::Write for CountingWriter {
fn write(&mut self, buffer: &[u8]) -> std::io::Result<usize> {
self.bytes = self.bytes.checked_add(buffer.len()).ok_or_else(|| {
std::io::Error::new(std::io::ErrorKind::FileTooLarge, "JSON byte count overflow")
})?;
Ok(buffer.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[derive(Debug)]
struct AbortedOpsFileError {
message: String,
#[cfg_attr(not(test), allow(dead_code))]
manifest: serde_json::Value,
}
impl std::fmt::Display for AbortedOpsFileError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str(&self.message)
}
}
impl std::error::Error for AbortedOpsFileError {}
#[allow(clippy::too_many_arguments)]
fn emit_aborted_ops_file_manifest(
error: anyhow::Error,
save_path: &str,
requested_total: usize,
confirmed_ops: usize,
committed_chunks: &[usize],
dispatched_chunk: Option<usize>,
summary: serde_json::Value,
) -> anyhow::Error {
let message = format!("{error:#}");
let mut manifest = serde_json::json!({
"status": "aborted",
"path": save_path,
"file_published": false,
"requested_total": requested_total,
"confirmed_ops": confirmed_ops,
"unconfirmed_ops": requested_total.saturating_sub(confirmed_ops),
"committed_chunks": committed_chunks,
"summary": summary,
"error": message.clone(),
});
if let Some(chunk_number) = dispatched_chunk {
manifest["dispatched_chunk"] = serde_json::json!(chunk_number);
}
println!(
"{}",
serde_json::to_string(&manifest).expect("serialize aborted ops-file manifest")
);
anyhow::Error::new(AbortedOpsFileError { message, manifest })
}
#[allow(clippy::too_many_arguments)]
#[cfg(test)]
async fn apply_ops_file_reader<R: std::io::BufRead>(
server: &KhiveMcpServer,
reader: R,
total: usize,
presentation: Option<String>,
_output_format: Option<String>,
save_file: Option<String>,
strict: bool,
) -> Result<serde_json::Value> {
apply_ops_file_reader_with_dispatch_mode(
server,
reader,
total,
presentation,
_output_format,
save_file,
strict,
OpsFileDispatchMode::BoundedParallel,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn apply_ops_file_reader_with_dispatch_mode<R: std::io::BufRead>(
server: &KhiveMcpServer,
reader: R,
total: usize,
presentation: Option<String>,
_output_format: Option<String>,
save_file: Option<String>,
strict: bool,
dispatch_mode: OpsFileDispatchMode,
) -> Result<serde_json::Value> {
apply_ops_file_reader_with_response_transform_and_dispatch_mode(
server,
reader,
total,
presentation,
_output_format,
save_file,
strict,
dispatch_mode,
|_, raw| raw,
)
.await
}
#[allow(clippy::too_many_arguments)]
#[cfg(test)]
async fn apply_ops_file_reader_with_response_transform<R, F>(
server: &KhiveMcpServer,
reader: R,
total: usize,
presentation: Option<String>,
_output_format: Option<String>,
save_file: Option<String>,
strict: bool,
response_transform: F,
) -> Result<serde_json::Value>
where
R: std::io::BufRead,
F: FnMut(usize, String) -> String,
{
apply_ops_file_reader_with_response_transform_and_dispatch_mode(
server,
reader,
total,
presentation,
_output_format,
save_file,
strict,
OpsFileDispatchMode::BoundedParallel,
response_transform,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn apply_ops_file_reader_with_response_transform_and_dispatch_mode<R, F>(
server: &KhiveMcpServer,
mut reader: R,
total: usize,
presentation: Option<String>,
_output_format: Option<String>,
save_file: Option<String>,
strict: bool,
dispatch_mode: OpsFileDispatchMode,
mut response_transform: F,
) -> Result<serde_json::Value>
where
R: std::io::BufRead,
F: FnMut(usize, String) -> String,
{
let report_mode = if save_file.is_some() {
OpsFileReportMode::BoundedSave
} else {
OpsFileReportMode::LegacyNoSave
};
let mut total_succeeded: usize = 0;
let mut total_failed: usize = 0;
let mut total_aborted: usize = 0;
let mut failures: Vec<serde_json::Value> = Vec::new();
let mut failure_details_omitted = 0_usize;
let save_path = save_file.clone();
let mut save_sink = save_file
.as_deref()
.map(|path| khive_mcp::save_sink::JsonlSaveSink::new(Path::new(path), false))
.transpose()?;
let mut processed = 0_usize;
let mut snapshot_line = 1_usize;
let mut chunk_idx = 0_usize;
let mut eof = false;
let mut pending: Option<(OpsFileEntry, usize)> = None;
let mut confirmed_ops = 0_usize;
let mut committed_chunks = Vec::new();
let mut dispatched_chunk = None;
let execution_result: Result<()> = async {
while !eof {
let mut chunk = Vec::with_capacity(OPS_FILE_CHUNK_SIZE);
let mut chunk_bytes = 0_usize;
if let Some((op, bytes)) = pending.take() {
chunk_bytes = bytes;
chunk.push(op);
}
while chunk.len() < OPS_FILE_CHUNK_SIZE {
let Some(raw) = read_bounded_ops_line(&mut reader, snapshot_line)? else {
eof = true;
break;
};
let physical_bytes = raw.len().saturating_add(1);
if let Some(op) = parse_ops_file_line(&raw, snapshot_line)? {
if should_defer_chunk_entry(chunk.len(), chunk_bytes, physical_bytes) {
pending = Some((op, physical_bytes));
snapshot_line += 1;
break;
}
chunk_bytes = chunk_bytes.saturating_add(physical_bytes);
chunk.push(op);
}
snapshot_line += 1;
}
if chunk.is_empty() {
break;
}
let applied_before = processed;
let chunk_len = chunk.len();
let expected_tools: Vec<String> = chunk.iter().map(|op| op.tool.clone()).collect();
let typed_ops: Vec<khive_request::TypedJsonOp> = chunk
.into_iter()
.map(|op| {
let serde_json::Value::Object(args) = op.args else {
unreachable!("validated ops-file args are always JSON objects")
};
khive_request::TypedJsonOp {
tool: op.tool,
args,
}
})
.collect();
let chunk_number = chunk_idx + 1;
dispatched_chunk = Some(chunk_number);
let output_format = save_sink.is_some().then(|| "json".to_string());
let raw = if dispatch_mode.is_serial() {
server
.dispatch_typed_json_batch_serial_local_for_exec(
typed_ops,
presentation.clone(),
output_format,
strict,
)
.await
} else {
server
.dispatch_typed_json_batch_local_for_exec(
typed_ops,
presentation.clone(),
output_format,
strict,
)
.await
}
.map_err(|error| anyhow::anyhow!("dispatch chunk {chunk_number}: {error}"))?;
let raw = response_transform(chunk_number, raw);
let mut parsed: serde_json::Value =
serde_json::from_str(&raw).context("parse dispatch result")?;
annotate_and_emit_refusals(&mut parsed, strict);
let (chunk_succeeded, chunk_failed, chunk_aborted) =
validate_ordered_chunk_envelope(&expected_tools, &parsed, chunk_number)?;
let chunk_results = parsed["results"]
.as_array()
.expect("validated ordered results array");
total_succeeded += chunk_succeeded;
total_failed += chunk_failed;
total_aborted += chunk_aborted;
confirmed_ops += chunk_len;
committed_chunks.push(chunk_number);
dispatched_chunk = None;
for failure in collect_op_failures(&parsed, applied_before, report_mode) {
let reason = match &failure["error"] {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
};
let op_index = failure["op_index"].clone();
let tool = failure["tool"].as_str().unwrap_or("?").to_string();
if retain_failure_detail(
report_mode,
failure,
&mut failures,
&mut failure_details_omitted,
) {
eprintln!("op {} ({}) failed: {reason}", op_index, tool);
}
}
if let Some(save_sink) = save_sink.as_mut() {
for row in chunk_results {
save_sink.write_row(row)?;
}
}
processed += chunk_len;
let applied_now = processed;
eprintln!(
"{}",
ops_file_progress_line(
report_mode,
applied_now,
total,
total_succeeded,
total_failed,
total_aborted,
)
);
chunk_idx += 1;
}
if processed != total {
anyhow::bail!(
"validated ops-file snapshot changed: expected {total} ops, read {}",
processed
);
}
Ok(())
}
.await;
if let Err(error) = execution_result {
if let Some(path) = save_path
.as_deref()
.filter(|_| !committed_chunks.is_empty() || dispatched_chunk.is_some())
{
drop(save_sink.take());
let summary = ops_file_summary(
report_mode,
confirmed_ops,
total_succeeded,
total_failed,
total_aborted,
failures,
failure_details_omitted,
);
return Err(emit_aborted_ops_file_manifest(
error,
path,
total,
confirmed_ops,
&committed_chunks,
dispatched_chunk,
summary,
));
}
return Err(error);
}
let summary = ops_file_summary(
report_mode,
total,
total_succeeded,
total_failed,
total_aborted,
failures,
failure_details_omitted,
);
let output = if let Some(save_sink) = save_sink {
let manifest = match save_sink.finish(summary.clone()) {
Ok(manifest) => manifest,
Err(error) => {
let path = save_path
.as_deref()
.expect("save sink exists only when save path exists");
return Err(emit_aborted_ops_file_manifest(
error,
path,
total,
confirmed_ops,
&committed_chunks,
None,
summary,
));
}
};
println!(
"{}",
serde_json::to_string(&manifest).expect("serialize save manifest")
);
manifest
} else {
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize summary")
);
summary
};
if total > 0 && total_succeeded == 0 {
match report_mode {
OpsFileReportMode::LegacyNoSave => anyhow::bail!(
"every op failed: {total_failed} op(s) failed out of {total}, 0 succeeded (see printed summary above)"
),
OpsFileReportMode::BoundedSave => anyhow::bail!(
"every op failed: {total_failed} failed, {total_aborted} aborted out of {total}, 0 succeeded (see printed output above)"
),
}
}
if strict {
match report_mode {
OpsFileReportMode::LegacyNoSave if total_failed > 0 => anyhow::bail!(
"--strict: {total_failed} op(s) failed out of {total} (see printed summary above)"
),
OpsFileReportMode::BoundedSave if total_failed > 0 || total_aborted > 0 => {
anyhow::bail!(
"--strict: {total_failed} op(s) failed, {total_aborted} op(s) aborted out of {total} (see printed output above)"
)
}
_ => {}
}
}
Ok(output)
}
#[cfg(test)]
async fn apply_ops_file(
server: &KhiveMcpServer,
ops: Vec<OpsFileEntry>,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
strict: bool,
) -> Result<serde_json::Value> {
let total = ops.len();
let mut encoded = Vec::new();
for op in ops {
serde_json::to_writer(
&mut encoded,
&serde_json::json!({"tool": op.tool, "args": op.args}),
)
.context("serialize test ops-file entry")?;
encoded.push(b'\n');
}
apply_ops_file_reader(
server,
std::io::Cursor::new(encoded),
total,
presentation,
output_format,
save_file,
strict,
)
.await
}
#[cfg(test)]
#[allow(clippy::too_many_arguments)]
async fn apply_ops_file_with_dispatch_mode(
server: &KhiveMcpServer,
ops: Vec<OpsFileEntry>,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
strict: bool,
dispatch_mode: OpsFileDispatchMode,
) -> Result<serde_json::Value> {
let total = ops.len();
let mut encoded = Vec::new();
for op in ops {
serde_json::to_writer(
&mut encoded,
&serde_json::json!({"tool": op.tool, "args": op.args}),
)
.context("serialize test ops-file entry")?;
encoded.push(b'\n');
}
apply_ops_file_reader_with_dispatch_mode(
server,
std::io::Cursor::new(encoded),
total,
presentation,
output_format,
save_file,
strict,
dispatch_mode,
)
.await
}
#[cfg(test)]
#[allow(clippy::too_many_arguments)]
async fn apply_ops_file_with_response_transform<F>(
server: &KhiveMcpServer,
ops: Vec<OpsFileEntry>,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
strict: bool,
response_transform: F,
) -> Result<serde_json::Value>
where
F: FnMut(usize, String) -> String,
{
let total = ops.len();
let mut encoded = Vec::new();
for op in ops {
serde_json::to_writer(
&mut encoded,
&serde_json::json!({"tool": op.tool, "args": op.args}),
)
.context("serialize test ops-file entry")?;
encoded.push(b'\n');
}
apply_ops_file_reader_with_response_transform(
server,
std::io::Cursor::new(encoded),
total,
presentation,
output_format,
save_file,
strict,
response_transform,
)
.await
}
pub async fn run_exec(args: ExecArgs) -> Result<()> {
if args.plan {
let result = plan::run(&args).await?;
writeln!(std::io::stdout().lock(), "{result}")?;
return Ok(());
}
if args.serial && (args.ops_file.is_none() || args.ops.is_some() || args.atomic) {
anyhow::bail!(
"--serial requires --ops-file and conflicts with positional ops and --atomic"
);
}
if args.pending_events {
let summary = pending_events::run_pending_events_with_config(
args.db.as_deref(),
args.config.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()),
};
preflight_exec_mode(&mode)?;
let namespace = Namespace::parse(&args.namespace).map_err(|e| anyhow::anyhow!("{e}"))?;
let (mut cfg, db_anchor) =
khive_mcp::serve::resolve_runtime_config_with_db_anchor(RuntimeConfigInputs {
db: args.db.as_deref(),
config: args.config.as_deref(),
namespace,
namespace_explicit: true,
actor_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})?;
debug_assert!(cfg
.events_split
.as_ref()
.is_none_or(|split| split.socket_path.is_none()));
if let Err(error) = apply_actor_pin_and_expectation(
&mut cfg,
args.actor.as_deref(),
args.expect_actor.as_deref(),
) {
if let Some(refusal) = error.downcast_ref::<ExecRefusal>() {
return Err(report_mode_refusal(&mode, refusal.reason, &refusal.message));
}
return Err(error);
}
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,
config: args.config,
};
match mode {
ExecMode::Inline(ops) => {
run_exec_inline(
ops,
cfg,
args.presentation,
args.output_format,
args.save_file,
db_context,
args.strict,
)
.await
}
ExecMode::OpsFile(path) => {
run_exec_ops_file(
path,
cfg,
args.presentation,
args.output_format,
args.save_file,
args.dry_run,
db_context,
args.serial,
args.atomic,
args.atomic_max_ops,
args.strict,
)
.await
}
}
}
fn apply_actor_pin_and_expectation(
cfg: &mut RuntimeConfig,
actor: Option<&str>,
expect_actor: Option<&str>,
) -> Result<()> {
if let Some(raw) = actor {
let parsed =
Namespace::parse(raw).map_err(|e| anyhow::anyhow!("invalid --actor {raw:?}: {e}"))?;
if let Some(prev) = cfg.actor_id.as_deref() {
if let Ok(prev_ns) = Namespace::parse(prev) {
cfg.visible_namespaces.retain(|ns| *ns != prev_ns);
}
}
cfg.actor_id = if parsed == Namespace::local() {
None
} else {
if !cfg.visible_namespaces.contains(&parsed) {
cfg.visible_namespaces.push(parsed.clone());
}
Some(parsed.as_str().to_owned())
};
}
if let Some(raw_expected) = expect_actor {
let expected = Namespace::parse(raw_expected)
.map_err(|e| anyhow::anyhow!("invalid --expect-actor {raw_expected:?}: {e}"))?;
let actual = khive_runtime::resolve_actor(cfg.actor_id.as_deref());
if actual.id != expected.as_str() {
return Err(refusal_error(
RefusalReason::ExpectActorMismatch,
format!(
"--expect-actor mismatch: expected {:?}, resolved {:?}",
expected.as_str(),
actual.id
),
));
}
}
Ok(())
}
fn enforce_strict_batch_result(raw: &str, strict: bool) -> Result<()> {
let Ok(parsed) = serde_json::from_str::<serde_json::Value>(raw) else {
return Ok(());
};
let succeeded = parsed["summary"]["succeeded"].as_u64().unwrap_or(0);
let failed = parsed
.get("summary")
.and_then(|summary| summary.get("failed"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
let aborted = parsed
.get("summary")
.and_then(|summary| summary.get("aborted"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(0);
if succeeded == 0 && (failed > 0 || aborted > 0) {
anyhow::bail!(
"every op failed: {failed} failed, {aborted} aborted, 0 succeeded (see printed output above)"
);
}
if strict && (failed > 0 || aborted > 0) {
anyhow::bail!(
"--strict: {failed} op(s) failed, {aborted} op(s) aborted (see printed output above)"
);
}
Ok(())
}
fn annotate_and_emit_refusals(parsed: &mut serde_json::Value, strict: bool) -> bool {
if let Some(reason) = parsed
.get("error")
.and_then(|error| error.get("reason"))
.and_then(serde_json::Value::as_str)
.and_then(RefusalReason::from_token)
{
emit_refusal(reason);
return false;
}
let failed = parsed["summary"]["failed"].as_u64().unwrap_or(0);
let aborted = parsed["summary"]["aborted"].as_u64().unwrap_or(0);
let strict_refusal = strict && (failed > 0 || aborted > 0);
let entries = parsed.as_object_mut().and_then(|object| {
if object
.get("results")
.is_some_and(serde_json::Value::is_array)
{
object
.get_mut("results")
.and_then(serde_json::Value::as_array_mut)
} else {
object
.get_mut("failures")
.and_then(serde_json::Value::as_array_mut)
}
});
let mut changed = false;
let mut emitted = 0usize;
if let Some(entries) = entries {
for entry in entries {
if entry["ok"].as_bool() == Some(true) {
continue;
}
let specific = entry["reason"].as_str().and_then(RefusalReason::from_token);
let reason = specific.or(strict_refusal.then_some(RefusalReason::StrictOpFailure));
if let Some(reason) = reason {
if specific.is_none() {
if let Some(object) = entry.as_object_mut() {
object.insert("reason".to_string(), serde_json::json!(reason.as_str()));
changed = true;
}
}
emit_refusal(reason);
emitted += 1;
}
}
}
if strict_refusal && emitted == 0 {
emit_refusal(RefusalReason::StrictOpFailure);
}
changed
}
fn prepare_exec_output(raw: &str, strict: bool) -> String {
let Ok(mut parsed) = serde_json::from_str::<serde_json::Value>(raw) else {
return raw.to_owned();
};
if annotate_and_emit_refusals(&mut parsed, strict) {
serde_json::to_string(&parsed).expect("serde_json::Value is serializable")
} else {
raw.to_owned()
}
}
enum ExecMode {
Inline(String),
OpsFile(PathBuf),
}
fn preflight_inline_ops(ops: &str) -> Result<()> {
if let Err(error) = khive_request::parse_request(ops) {
let error = rmcp::ErrorData::invalid_params(error.to_string(), None);
return Err(report_unscoped_refusal(
RefusalReason::ParseError,
error.to_string(),
));
}
Ok(())
}
fn preflight_exec_mode(mode: &ExecMode) -> Result<()> {
match mode {
ExecMode::Inline(ops) => preflight_inline_ops(ops),
ExecMode::OpsFile(path) => match validate_ops_file(path) {
Ok(_) => Ok(()),
Err(error) => {
if let Some(refusal) = error.downcast_ref::<ExecRefusal>() {
Err(report_unscoped_refusal(
refusal.reason,
refusal.message.as_str(),
))
} else {
Err(error)
}
}
},
}
}
fn report_mode_refusal(
mode: &ExecMode,
reason: RefusalReason,
error: impl std::fmt::Display,
) -> anyhow::Error {
let message = error.to_string();
let parsed = match mode {
ExecMode::Inline(raw) => khive_request::parse_request(raw).ok().map(|request| {
let chain = request.mode == khive_request::ExecutionMode::Chain;
let tools = request.ops.into_iter().map(|op| op.tool).collect();
(tools, chain)
}),
ExecMode::OpsFile(path) => validate_ops_file(path).ok().and_then(|mut validated| {
validated_tool_names(&mut validated.snapshot)
.ok()
.map(|tools| (tools, false))
}),
};
match parsed {
Some((tools, chain)) if !tools.is_empty() => {
report_tools_refusal(tools, chain, reason, message)
}
_ => report_unscoped_refusal(reason, message),
}
}
fn disclose_resolved_database(cfg: &RuntimeConfig, khive_cfg: &KhiveConfig) {
use std::io::Write;
let line =
khive_mcp::serve::resolved_database_disclosure(cfg.db_path.as_deref(), &khive_cfg.backends);
let _ = writeln!(std::io::stderr(), "{line}");
}
fn disclose_resolved_actor(cfg: &RuntimeConfig) {
use std::io::Write;
let line = khive_mcp::serve::resolved_actor_disclosure(cfg.actor_id.as_deref());
let _ = writeln!(std::io::stderr(), "{line}");
}
#[cfg(unix)]
fn disclose_daemon_execution() {
use std::io::Write;
let _ = writeln!(
std::io::stderr(),
"execution: answered by daemon; --log and KHIVE_LOG set the client process log level only; \
the daemon log level is fixed at startup"
);
}
#[derive(Default)]
struct ExecDbContext {
raw: Option<String>,
anchor: Option<PathBuf>,
config: Option<PathBuf>,
}
fn load_exec_config(db_context: &ExecDbContext) -> Result<(KhiveConfig, Option<PathBuf>)> {
let db_path_for_config = config_discovery_db_anchor(db_context.raw.as_deref());
let loaded = KhiveConfig::load_with_home_fallback_and_source(
db_context.config.as_deref(),
db_path_for_config.as_deref(),
)
.map_err(|e| anyhow::anyhow!("config error: {e}"))?;
Ok(match loaded {
Some((config, source)) => (config, Some(source)),
None => (KhiveConfig::default(), None),
})
}
async fn run_exec_inline(
ops: String,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
db_context: ExecDbContext,
strict: bool,
) -> Result<()> {
#[cfg(unix)]
return run_exec_inline_with_forward(
ops,
cfg,
presentation,
output_format,
save_file,
db_context,
strict,
forward_or_spawn_boxed,
)
.await;
#[cfg(not(unix))]
return run_exec_inline_with_forward(
ops,
cfg,
presentation,
output_format,
save_file,
db_context,
strict,
)
.await;
}
const EXEC_STORAGE_SETTLE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
async fn settle_exec_storage_before_return(server: KhiveMcpServer) {
if let Err(reason) = server.shutdown_audit_batch().await {
tracing::warn!(
reason = ?reason,
"exec: audit-batch drain did not complete cleanly before storage settle"
);
}
let mut writer_task_joins = Vec::new();
if let Some(pool) = server.pool() {
if let Some(join) = pool.take_writer_task_join() {
writer_task_joins.push(join);
}
}
for pool in server.secondary_pools() {
if let Some(join) = pool.take_writer_task_join() {
writer_task_joins.push(join);
}
}
drop(server);
for join in writer_task_joins {
match tokio::time::timeout(EXEC_STORAGE_SETTLE_TIMEOUT, join).await {
Ok(Ok(())) => {}
Ok(Err(error)) => tracing::warn!(
%error,
"exec: a writer task panicked while settling storage before return"
),
Err(_) => tracing::warn!(
timeout = ?EXEC_STORAGE_SETTLE_TIMEOUT,
"exec: a writer task did not settle before returning"
),
}
}
}
#[cfg_attr(not(unix), allow(unused_variables))]
#[allow(clippy::too_many_arguments)]
async fn run_exec_inline_with_forward(
ops: String,
mut cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
mut db_context: ExecDbContext,
strict: bool,
#[cfg(unix)] forward_fn: ForwardFnPtr,
) -> Result<()> {
preflight_inline_ops(&ops)?;
if let Err(error) = enforce_strict_actor_mode(cfg.actor_id.as_deref(), &cfg.packs) {
return Err(report_invocation_refusal(
Some(&ops),
RefusalReason::AnonymousActor,
error,
));
}
let (khive_cfg, config_source) = load_exec_config(&db_context)?;
let force_memory = if khive_cfg.backends.is_empty() {
false
} else {
let force_memory = normalize_redundant_db_override_with_source(
&mut cfg,
db_context.raw.as_deref(),
&khive_cfg.backends,
config_source.as_deref(),
)?;
db_context.anchor = cfg.db_path.clone();
force_memory
};
if !force_memory {
validate_declared_backend_access_modes(&khive_cfg.backends)?;
}
disclose_resolved_database(&cfg, &khive_cfg);
disclose_resolved_actor(&cfg);
#[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(),
process_ref: khive_runtime::process_ref_from_env(),
visible_namespaces: cfg
.visible_namespaces
.iter()
.map(|ns| ns.as_str().to_string())
.collect(),
config_id: if force_memory {
compute_config_id_with_storage_mode(&cfg, Some(&khive_cfg), false)
} else {
compute_config_id(&cfg, Some(&khive_cfg))
},
protocol_version: PROTOCOL_VERSION,
plan: false,
probe_only: false,
metrics_only: false,
format: output_format.clone(),
format_per_op: None,
from_wire: false,
request_id: None,
};
let spawn_db = match db_context.raw.as_deref() {
Some(":memory:") => Some(":memory:"),
Some(concrete) if khive_cfg.backends.is_empty() => Some(concrete),
_ => None,
};
let spawn_config = match (&db_context.config, db_context.raw.as_deref()) {
(Some(explicit), _) => Some(explicit.clone()),
(None, Some(raw)) if raw != ":memory:" && !khive_cfg.backends.is_empty() => {
config_source.clone()
}
_ => None,
};
let spawn_packs = cfg.packs.clone();
if let Some(res) = forward_fn(&frame, spawn_config, spawn_db, spawn_packs).await {
let output = res.map_err(|e| anyhow::anyhow!("{}", e.message))?;
disclose_daemon_execution();
let output = prepare_exec_output(&output, strict);
println!("{output}");
enforce_strict_batch_result(&output, strict)?;
return Ok(());
}
}
let server = build_local_fallback_server(
cfg,
&khive_cfg,
db_context.raw.as_deref(),
db_context.anchor.as_deref(),
)
.await?;
let params = RequestParams {
plan: None,
ops,
presentation,
presentation_per_op: None,
save_to: save_file,
format: output_format,
format_per_op: None,
request_id: None,
};
let dispatch_result = server
.dispatch_request_local_for_exec(params, strict)
.await
.map_err(|e| anyhow::anyhow!("{e}"));
let final_result = match dispatch_result {
Ok(raw_output) => {
let output = prepare_exec_output(&raw_output, strict);
println!("{output}");
enforce_strict_batch_result(&output, strict)
}
Err(error) => Err(error),
};
settle_exec_storage_before_return(server).await;
final_result
}
async 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 = build_single_backend_runtime(cfg, khive_cfg).await?;
let env_fmt = apply_env_output_format(khive_cfg.runtime.default_output_format);
Ok(KhiveMcpServer::new_with_mounts(rt)
.await
.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).await
}
}
struct AtomicSavePublishFailure {
stdout: String,
error: anyhow::Error,
}
fn render_atomic_output(
envelope: &mut serde_json::Value,
save_sink: Option<khive_mcp::save_sink::JsonlSaveSink>,
) -> std::result::Result<String, AtomicSavePublishFailure> {
let Some(save_sink) = save_sink else {
return Ok(serde_json::to_string_pretty(envelope).expect("serialize atomic envelope"));
};
match save_sink.write_envelope(envelope) {
Ok(manifest) => {
Ok(serde_json::to_string(&manifest).expect("serialize atomic save manifest"))
}
Err(error) => {
let committed = crate::atomic_apply::record_save_file_publish_failure(envelope, &error);
let stdout = serde_json::to_string_pretty(envelope)
.expect("serialize atomic save failure reconciliation envelope");
let error = if committed {
error.context(
"atomic database changes committed but --save-file publication failed; \
do not replay the mutation (inspect stdout for reconciliation details)",
)
} else {
error.context("atomic --save-file publication failed")
};
Err(AtomicSavePublishFailure { stdout, error })
}
}
}
#[allow(clippy::too_many_arguments)]
async fn run_exec_ops_file(
path: PathBuf,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
dry_run: bool,
db_context: ExecDbContext,
serial: bool,
atomic: bool,
atomic_max_ops: Option<usize>,
strict: bool,
) -> Result<()> {
if serial && atomic {
anyhow::bail!("--serial conflicts with --atomic");
}
let mut validated = match validate_ops_file(&path) {
Ok(validated) => validated,
Err(error) => {
if let Some(refusal) = error.downcast_ref::<ExecRefusal>() {
return Err(report_unscoped_refusal(
refusal.reason,
refusal.message.as_str(),
));
}
return Err(error);
}
};
if validated.total == 0 {
anyhow::bail!("ops-file is empty (no non-blank lines): {}", path.display());
}
if serial {
if let Err(error) =
preflight_typed_validated_snapshot(&mut validated.snapshot, validated.total)
{
if let Some(refusal) = error.downcast_ref::<ExecRefusal>() {
return Err(report_unscoped_refusal(
refusal.reason,
refusal.message.as_str(),
));
}
return Err(error);
}
}
if dry_run && !atomic {
let summary = serde_json::json!({
"dry_run": true,
"total": validated.total,
"per_verb": validated.per_verb,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize dry-run summary")
);
return Ok(());
}
if let Err(error) = enforce_strict_actor_mode(cfg.actor_id.as_deref(), &cfg.packs) {
let tools = validated_tool_names(&mut validated.snapshot)?;
return Err(report_tools_refusal(
tools,
false,
RefusalReason::AnonymousActor,
error.to_string(),
));
}
let (khive_cfg, config_source) = load_exec_config(&db_context)?;
if !khive_cfg.backends.is_empty() {
reject_conflicting_db_override_with_source(
db_context.raw.as_deref(),
&khive_cfg.backends,
config_source.as_deref(),
)?;
}
disclose_resolved_database(&cfg, &khive_cfg);
disclose_resolved_actor(&cfg);
if atomic {
let max_ops = atomic_max_ops.unwrap_or(khive_types::pack::ATOMIC_MAX_OPS_DEFAULT);
let ops =
parse_atomic_validated_snapshot(&mut validated.snapshot, validated.total, max_ops)?;
if dry_run {
if let Err(error) =
crate::atomic_apply::preflight_atomic_ops_file(&ops, &cfg, &khive_cfg, max_ops)
{
if let Some(failure) =
error.downcast_ref::<crate::atomic_apply::AtomicExecFailure>()
{
let mut envelope = failure.envelope();
annotate_and_emit_refusals(&mut envelope, strict);
println!(
"{}",
serde_json::to_string_pretty(&envelope)
.expect("serialize atomic dry-run refusal envelope")
);
}
return Err(error);
}
let summary = serde_json::json!({
"dry_run": true,
"atomic": true,
"total": validated.total,
"per_verb": validated.per_verb,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize atomic dry-run summary")
);
return Ok(());
}
let save_sink = save_file
.as_deref()
.map(|path| khive_mcp::save_sink::JsonlSaveSink::new(Path::new(path), false))
.transpose()?;
let mut envelope =
match crate::atomic_apply::execute_atomic_ops_file(ops, cfg, &khive_cfg, max_ops).await
{
Ok(envelope) => envelope,
Err(error) => {
if let Some(failure) =
error.downcast_ref::<crate::atomic_apply::AtomicExecFailure>()
{
let mut envelope = failure.envelope();
annotate_and_emit_refusals(&mut envelope, strict);
println!(
"{}",
serde_json::to_string_pretty(&envelope)
.expect("serialize atomic refusal envelope")
);
}
return Err(error);
}
};
annotate_and_emit_refusals(&mut envelope, strict);
let output = match render_atomic_output(&mut envelope, save_sink) {
Ok(output) => output,
Err(failure) => {
println!("{}", failure.stdout);
return Err(failure.error);
}
};
println!("{output}");
if envelope["atomic"]["rolled_back"].as_bool() == Some(true) {
anyhow::bail!("atomic unit rolled back; inspect stdout for the failed operation");
}
return Ok(());
}
let server = build_local_fallback_server(
cfg,
&khive_cfg,
db_context.raw.as_deref(),
db_context.anchor.as_deref(),
)
.await?;
validated
.snapshot
.rewind()
.context("rewind validated ops-file snapshot for dispatch")?;
let dispatch_result = apply_ops_file_reader_with_dispatch_mode(
&server,
std::io::BufReader::new(validated.snapshot),
validated.total,
presentation,
output_format,
save_file,
strict,
if serial {
OpsFileDispatchMode::Serial
} else {
OpsFileDispatchMode::BoundedParallel
},
)
.await
.map(|_| ());
settle_exec_storage_before_return(server).await;
dispatch_result
}
#[cfg(test)]
mod tests {
use super::*;
use clap::Parser;
use serial_test::serial;
use tempfile::NamedTempFile;
use uuid::Uuid;
#[test]
fn collect_op_failures_extracts_reason_with_global_index() {
let parsed = serde_json::json!({
"results": [
{"ok": true, "tool": "create", "result": {}},
{"ok": false, "tool": "create", "error": "content rejected: suspected credential material"},
{"ok": false, "tool": "link"},
],
"summary": {"total": 3, "succeeded": 1, "failed": 2}
});
let failures = collect_op_failures(&parsed, 500, OpsFileReportMode::LegacyNoSave);
assert_eq!(failures.len(), 2);
assert_eq!(failures[0]["op_index"], 501);
assert_eq!(failures[0]["tool"], "create");
assert_eq!(
failures[0]["error"],
"content rejected: suspected credential material"
);
assert_eq!(failures[1]["op_index"], 502);
assert_eq!(
failures[1]["error"], "unknown error",
"a failed entry with no error value still surfaces a placeholder"
);
}
#[test]
fn collect_op_failures_preserves_structured_error_payloads() {
let parsed = serde_json::json!({
"results": [
{"ok": false, "tool": "create",
"error": {"kind": "invalid_input", "message": "content rejected"}},
],
"summary": {"total": 1, "succeeded": 0, "failed": 1}
});
let failures = collect_op_failures(&parsed, 0, OpsFileReportMode::LegacyNoSave);
assert_eq!(
failures[0]["error"],
serde_json::json!({"kind": "invalid_input", "message": "content rejected"}),
"structured KhiveError payloads pass through as JSON, not a placeholder"
);
}
#[test]
fn collect_op_failures_preserves_stable_refusal_reason() {
let parsed = serde_json::json!({
"results": [
{
"ok": false,
"tool": "not_loaded",
"error": "unknown verb",
"reason": "verb-refused"
},
],
"summary": {"total": 1, "succeeded": 0, "failed": 1}
});
let failures = collect_op_failures(&parsed, 9, OpsFileReportMode::BoundedSave);
assert_eq!(failures[0]["reason"], "verb-refused");
assert_eq!(failures[0]["op_index"], 9);
let legacy = collect_op_failures(&parsed, 9, OpsFileReportMode::LegacyNoSave);
assert!(
legacy[0].get("reason").is_none(),
"legacy no-save summary must retain its pre-reason wire shape"
);
}
#[test]
fn collect_op_failures_empty_on_all_ok_or_missing_results() {
let all_ok = serde_json::json!({
"results": [{"ok": true, "tool": "stats", "result": {}}],
"summary": {"total": 1, "succeeded": 1, "failed": 0}
});
assert!(collect_op_failures(&all_ok, 0, OpsFileReportMode::LegacyNoSave).is_empty());
assert!(
collect_op_failures(&serde_json::json!({}), 0, OpsFileReportMode::LegacyNoSave)
.is_empty()
);
}
#[test]
fn no_save_reporting_matches_exact_legacy_golden() {
let parsed = serde_json::json!({
"results": [
{"ok": true, "tool": "create", "result": {}},
{"ok": false, "tool": "search", "error": "boom"},
],
});
let failures = collect_op_failures(&parsed, 0, OpsFileReportMode::LegacyNoSave);
assert_eq!(failures.len(), 1);
assert!(failures[0].get("aborted").is_none());
let summary = ops_file_summary(OpsFileReportMode::LegacyNoSave, 2, 1, 1, 0, failures, 0);
assert_eq!(
ops_file_progress_line(OpsFileReportMode::LegacyNoSave, 2, 2, 1, 1, 0),
"applied 2/2 (ok=1, failed=1)"
);
assert_eq!(
serde_json::to_string_pretty(&summary).unwrap(),
concat!(
"{\n",
" \"failed\": 1,\n",
" \"failures\": [\n",
" {\n",
" \"error\": \"boom\",\n",
" \"op_index\": 1,\n",
" \"tool\": \"search\"\n",
" }\n",
" ],\n",
" \"succeeded\": 1,\n",
" \"total\": 2\n",
"}"
)
);
assert!(summary.get("aborted").is_none());
assert!(summary.get("failure_details_omitted").is_none());
}
#[test]
fn no_save_reporting_retains_more_than_one_thousand_failures() {
let parsed = serde_json::json!({
"results": (0..=MAX_OPS_FILE_FAILURE_DETAILS)
.map(|index| serde_json::json!({
"ok": false,
"tool": "create",
"error": format!("failure-{index}"),
}))
.collect::<Vec<_>>(),
});
let mut retained = Vec::new();
let mut omitted = 0;
for failure in collect_op_failures(&parsed, 0, OpsFileReportMode::LegacyNoSave) {
assert!(retain_failure_detail(
OpsFileReportMode::LegacyNoSave,
failure,
&mut retained,
&mut omitted,
));
}
assert_eq!(retained.len(), MAX_OPS_FILE_FAILURE_DETAILS + 1);
assert_eq!(omitted, 0);
let summary = ops_file_summary(
OpsFileReportMode::LegacyNoSave,
retained.len(),
0,
retained.len(),
0,
retained,
omitted,
);
assert_eq!(
summary["failures"].as_array().unwrap().len(),
MAX_OPS_FILE_FAILURE_DETAILS + 1
);
assert!(summary.get("failure_details_omitted").is_none());
}
#[test]
fn no_save_reporting_retains_error_larger_than_four_kib() {
let large_error = "x".repeat(MAX_OPS_FILE_FAILURE_ERROR_BYTES + 1);
let parsed = serde_json::json!({
"results": [{"ok": false, "tool": "create", "error": large_error}],
});
let failures = collect_op_failures(&parsed, 0, OpsFileReportMode::LegacyNoSave);
assert_eq!(
failures[0]["error"].as_str().unwrap().len(),
MAX_OPS_FILE_FAILURE_ERROR_BYTES + 1
);
assert_eq!(failures[0]["error"], large_error);
}
#[test]
fn save_reporting_bounds_failure_count_and_error_detail() {
let large_error = "x".repeat(MAX_OPS_FILE_FAILURE_ERROR_BYTES + 1);
let parsed = serde_json::json!({
"results": (0..=MAX_OPS_FILE_FAILURE_DETAILS)
.map(|index| serde_json::json!({
"ok": false,
"tool": "create",
"aborted": index % 2 == 0,
"error": if index == 0 {
serde_json::Value::String(large_error.clone())
} else {
serde_json::Value::String(format!("failure-{index}"))
},
}))
.collect::<Vec<_>>(),
});
let mut retained = Vec::new();
let mut omitted = 0;
for failure in collect_op_failures(&parsed, 0, OpsFileReportMode::BoundedSave) {
retain_failure_detail(
OpsFileReportMode::BoundedSave,
failure,
&mut retained,
&mut omitted,
);
}
assert_eq!(retained.len(), MAX_OPS_FILE_FAILURE_DETAILS);
assert_eq!(omitted, 1);
assert_eq!(retained[0]["aborted"], true);
assert_eq!(
retained[0]["error"],
format!(
"error detail omitted: exceeds {MAX_OPS_FILE_FAILURE_ERROR_BYTES}-byte ops-file diagnostic limit"
)
);
let summary = ops_file_summary(
OpsFileReportMode::BoundedSave,
MAX_OPS_FILE_FAILURE_DETAILS + 1,
0,
MAX_OPS_FILE_FAILURE_DETAILS + 1,
0,
retained,
omitted,
);
assert_eq!(summary["failure_details_omitted"], 1);
assert_eq!(
summary["failures"].as_array().unwrap().len(),
MAX_OPS_FILE_FAILURE_DETAILS
);
}
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"),
}
}
const DAEMON_SPAWN_TEST_ENV_VARS: [&str; 8] = [
"KHIVE_EMBEDDING_MODEL",
"KHIVE_ADDITIONAL_EMBEDDING_MODELS",
"KHIVE_ACTOR",
"KHIVE_REQUIRE_ATTRIBUTED_ACTOR",
"KHIVE_DB",
"KHIVE_PACKS",
"KHIVE_LOCK",
"HOME",
];
struct EnvAndCwdGuard {
original_env: Vec<(&'static str, Option<std::ffi::OsString>)>,
original_cwd: std::path::PathBuf,
}
impl EnvAndCwdGuard {
fn capture() -> Self {
Self {
original_env: DAEMON_SPAWN_TEST_ENV_VARS
.into_iter()
.map(|name| (name, std::env::var_os(name)))
.collect(),
original_cwd: std::env::current_dir().expect("read cwd"),
}
}
}
impl Drop for EnvAndCwdGuard {
fn drop(&mut self) {
for (name, original) in self.original_env.drain(..) {
match original {
Some(value) => std::env::set_var(name, value),
None => std::env::remove_var(name),
}
}
let _ = std::env::set_current_dir(&self.original_cwd);
}
}
#[test]
#[serial]
fn daemon_spawn_env_guard_restores_every_mutated_variable() {
if crate::test_process::run_in_child() {
return;
}
let _restore_machine_env = EnvAndCwdGuard::capture();
for name in DAEMON_SPAWN_TEST_ENV_VARS {
std::env::set_var(name, format!("sentinel-{name}"));
}
{
let _guard = EnvAndCwdGuard::capture();
for name in DAEMON_SPAWN_TEST_ENV_VARS {
std::env::remove_var(name);
}
}
for name in DAEMON_SPAWN_TEST_ENV_VARS {
assert_eq!(
std::env::var(name).as_deref(),
Ok(format!("sentinel-{name}").as_str())
);
}
}
#[test]
#[serial]
fn acquire_local_construction_guard_is_noop_for_in_memory_db() {
if crate::test_process::run_in_child() {
return;
}
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]
fn acquire_local_construction_guard_serializes_concurrent_file_backed_callers() {
if crate::test_process::run_in_child() {
return;
}
acquire_local_construction_guard_serializes_concurrent_file_backed_callers_impl();
}
#[cfg(not(unix))]
#[test]
#[serial]
fn acquire_local_construction_guard_serializes_concurrent_file_backed_callers_nonunix() {
if crate::test_process::run_in_child() {
return;
}
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() {
if crate::test_process::run_in_child() {
return;
}
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]
#[serial]
fn config_flag_and_env_bind_with_flag_precedence() {
if crate::test_process::run_in_child() {
return;
}
let previous = std::env::var_os("KHIVE_CONFIG");
std::env::set_var("KHIVE_CONFIG", "/tmp/kkernel-exec-env-config.toml");
let from_env = ExecArgs::parse_from(["exec", "stats()"]);
assert_eq!(
from_env.config.as_deref(),
Some(std::path::Path::new("/tmp/kkernel-exec-env-config.toml"))
);
let from_flag = ExecArgs::parse_from([
"exec",
"stats()",
"--config",
"/tmp/kkernel-exec-flag-config.toml",
]);
assert_eq!(
from_flag.config.as_deref(),
Some(std::path::Path::new("/tmp/kkernel-exec-flag-config.toml"))
);
match previous {
Some(value) => std::env::set_var("KHIVE_CONFIG", value),
None => std::env::remove_var("KHIVE_CONFIG"),
}
}
#[test]
fn explicit_config_flag_parses_for_exec() {
let args = ExecArgs::parse_from([
"exec",
"stats()",
"--config",
"/tmp/kkernel-exec-config.toml",
]);
assert_eq!(
args.config.as_deref(),
Some(std::path::Path::new("/tmp/kkernel-exec-config.toml"))
);
}
#[test]
fn actor_and_expect_actor_flags_parse_together() {
let args = ExecArgs::parse_from([
"exec",
"stats()",
"--actor",
"lambda:worker",
"--expect-actor",
"lambda:worker",
]);
assert_eq!(args.actor.as_deref(), Some("lambda:worker"));
assert_eq!(args.expect_actor.as_deref(), Some("lambda:worker"));
}
#[test]
#[serial]
fn khive_actor_env_does_not_bind_to_explicit_actor_arg() {
if crate::test_process::run_in_child() {
return;
}
let previous = std::env::var("KHIVE_ACTOR").ok();
std::env::set_var("KHIVE_ACTOR", "lambda:env");
let args = ExecArgs::parse_from(["exec", "stats()"]);
match previous {
Some(value) => std::env::set_var("KHIVE_ACTOR", value),
None => std::env::remove_var("KHIVE_ACTOR"),
}
assert_eq!(args.actor, None, "the env fallback must not become tier 1");
}
#[test]
fn actor_flags_conflict_with_pending_events_mode() {
assert!(ExecArgs::try_parse_from(
["exec", "--pending-events", "--actor", "lambda:worker",]
)
.is_err());
assert!(ExecArgs::try_parse_from([
"exec",
"--pending-events",
"--expect-actor",
"lambda:worker",
])
.is_err());
}
#[test]
fn explicit_actor_overrides_fallback_without_changing_namespace() {
let mut cfg = RuntimeConfig {
default_namespace: Namespace::parse("project:data").unwrap(),
actor_id: Some("lambda:fallback".to_string()),
visible_namespaces: vec![Namespace::parse("lambda:fallback").unwrap()],
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, Some("lambda:cli"), Some("lambda:cli")).unwrap();
assert_eq!(cfg.actor_id.as_deref(), Some("lambda:cli"));
assert_eq!(cfg.default_namespace.as_str(), "project:data");
assert_eq!(
cfg.visible_namespaces,
vec![Namespace::parse("lambda:cli").unwrap()],
"pinning must drop the displaced actor's folded read visibility and add the pinned one"
);
}
#[test]
fn explicit_local_actor_authoritatively_clears_fallback() {
let mut cfg = RuntimeConfig {
actor_id: Some("lambda:fallback".to_string()),
visible_namespaces: vec![Namespace::parse("lambda:fallback").unwrap()],
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, Some("local"), Some("local")).unwrap();
assert_eq!(cfg.actor_id, None);
assert!(
cfg.visible_namespaces.is_empty(),
"pinning to local must drop the displaced fallback actor's read visibility \
without adding a replacement: {:?}",
cfg.visible_namespaces
);
}
#[test]
fn explicit_actor_pin_retains_unrelated_configured_visibility() {
let mut cfg = RuntimeConfig {
actor_id: Some("lambda:fallback".to_string()),
visible_namespaces: vec![
Namespace::parse("lambda:fallback").unwrap(),
Namespace::parse("project:shared").unwrap(),
],
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, Some("lambda:cli"), None).unwrap();
assert_eq!(
cfg.visible_namespaces,
vec![
Namespace::parse("project:shared").unwrap(),
Namespace::parse("lambda:cli").unwrap(),
],
"an explicitly configured extra visibility entry unrelated to the displaced \
actor must survive the pin: {:?}",
cfg.visible_namespaces
);
}
#[test]
fn expect_actor_alone_validates_resolved_identity() {
let mut cfg = RuntimeConfig {
actor_id: Some("lambda:project".to_string()),
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, None, Some("lambda:project")).unwrap();
let err = apply_actor_pin_and_expectation(&mut cfg, None, Some("lambda:other"))
.expect_err("a mismatched expectation must fail before dispatch");
assert!(err.to_string().contains("--expect-actor mismatch"));
assert!(err.to_string().contains("lambda:project"));
}
#[test]
fn actor_inputs_are_namespace_validated() {
let mut cfg = RuntimeConfig::default();
assert!(apply_actor_pin_and_expectation(&mut cfg, Some("bad actor"), None).is_err());
assert!(apply_actor_pin_and_expectation(&mut cfg, None, Some("bad actor")).is_err());
}
#[tokio::test]
#[serial]
async fn authorized_explicit_actor_is_used_for_write_attribution() {
if crate::test_process::run_in_child() {
return;
}
let (previous_home, _home_dir) = isolate_home_for_test();
let mut cfg = RuntimeConfig {
db_path: None,
actor_id: Some("lambda:fallback".to_string()),
gate: std::sync::Arc::new(khive_runtime::AllowAllGate),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string(), "comm".to_string()],
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, Some("lambda:pinned"), None).unwrap();
let server = build_local_fallback_server(cfg, &KhiveConfig::default(), None, None)
.await
.unwrap();
let raw = server
.dispatch_request_local(RequestParams {
plan: None,
ops: r#"comm.send(to="local", content="actor pin attribution")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.unwrap();
restore_home(previous_home);
let response: serde_json::Value = serde_json::from_str(&raw).unwrap();
assert_eq!(
response["results"][0]["ok"], true,
"comm.send dispatch failed: {raw}"
);
assert_eq!(response["results"][0]["result"]["from"], "lambda:pinned");
}
#[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"
);
}
#[test]
fn serial_requires_ops_file_conflicts_with_inline_and_atomic_and_defaults_off() {
let serial =
ExecArgs::try_parse_from(["exec", "--ops-file", "/tmp/batch.jsonl", "--serial"])
.expect("--serial must be accepted for a non-atomic ops-file");
assert!(serial.serial);
let default_parallel = ExecArgs::parse_from(["exec", "--ops-file", "/tmp/batch.jsonl"]);
assert!(!default_parallel.serial);
assert!(
ExecArgs::try_parse_from(["exec", "stats()", "--serial"]).is_err(),
"--serial without --ops-file must fail during CLI parsing"
);
assert!(
ExecArgs::try_parse_from([
"exec",
"stats()",
"--ops-file",
"/tmp/batch.jsonl",
"--serial",
])
.is_err(),
"--serial must not compose with inline positional ops"
);
assert!(
ExecArgs::try_parse_from([
"exec",
"--ops-file",
"/tmp/batch.jsonl",
"--atomic",
"--serial",
])
.is_err(),
"--serial and --atomic are distinct execution contracts and must conflict"
);
}
#[test]
fn atomic_takes_no_inline_batch_of_chains() {
for ops in ["[stats() | stats(), stats()]", "[stats(), stats()]"] {
let error = ExecArgs::try_parse_from(["exec", ops, "--atomic"])
.expect_err("--atomic with inline ops must fail during CLI parsing");
assert_eq!(
error.kind(),
clap::error::ErrorKind::MissingRequiredArgument,
"{ops}: {error}"
);
}
let atomic =
ExecArgs::try_parse_from(["exec", "--ops-file", "/tmp/batch.jsonl", "--atomic"])
.expect("--atomic must be accepted with an ops file");
assert!(atomic.atomic);
}
fn isolated_server(db_path: &str) -> KhiveMcpServer {
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(db_path)),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string(), "gtd".to_string()],
..Default::default()
};
let rt = KhiveRuntime::new(cfg).expect("runtime on temp db");
KhiveMcpServer::new(rt).expect("server on temp db")
}
struct OpsFileConcurrencyProbePack {
in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
max_in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
reject_overlap: bool,
}
impl khive_types::Pack for OpsFileConcurrencyProbePack {
const NAME: &'static str = "ops-file-concurrency-probe";
const NOTE_KINDS: &'static [&'static str] = &[];
const ENTITY_KINDS: &'static [&'static str] = &[];
const HANDLERS: &'static [khive_runtime::HandlerDef] = &[khive_runtime::HandlerDef {
name: "reader_probe",
description: "records test-only handler concurrency",
visibility: khive_runtime::Visibility::Verb,
category: khive_runtime::VerbCategory::Assertive,
params: &[],
}];
}
struct ProbeInFlightGuard {
in_flight: std::sync::Arc<std::sync::atomic::AtomicUsize>,
}
impl Drop for ProbeInFlightGuard {
fn drop(&mut self) {
self.in_flight
.fetch_sub(1, std::sync::atomic::Ordering::SeqCst);
}
}
#[async_trait::async_trait]
impl khive_runtime::PackRuntime for OpsFileConcurrencyProbePack {
fn name(&self) -> &str {
<Self as khive_types::Pack>::NAME
}
fn note_kinds(&self) -> &'static [&'static str] {
<Self as khive_types::Pack>::NOTE_KINDS
}
fn entity_kinds(&self) -> &'static [&'static str] {
<Self as khive_types::Pack>::ENTITY_KINDS
}
fn handlers(&self) -> &'static [khive_runtime::HandlerDef] {
<Self as khive_types::Pack>::HANDLERS
}
async fn dispatch(
&self,
_verb: &str,
params: serde_json::Value,
_registry: &khive_runtime::VerbRegistry,
_token: &khive_runtime::NamespaceToken,
) -> std::result::Result<serde_json::Value, khive_runtime::RuntimeError> {
let current = self
.in_flight
.fetch_add(1, std::sync::atomic::Ordering::SeqCst)
+ 1;
let _guard = ProbeInFlightGuard {
in_flight: self.in_flight.clone(),
};
self.max_in_flight
.fetch_max(current, std::sync::atomic::Ordering::SeqCst);
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
if params["fail"].as_bool().unwrap_or(false) {
return Err(khive_runtime::RuntimeError::Internal(
"requested probe failure".to_string(),
));
}
if self.reject_overlap && current > 1 {
return Err(khive_runtime::RuntimeError::Internal(
"sql_bridge.reader_open constrained-reader overlap".to_string(),
));
}
Ok(serde_json::json!({"sequence": params["sequence"]}))
}
}
fn concurrency_probe_server(
reject_overlap: bool,
) -> (
KhiveMcpServer,
std::sync::Arc<std::sync::atomic::AtomicUsize>,
std::sync::Arc<std::sync::atomic::AtomicUsize>,
) {
let in_flight = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let max_in_flight = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut builder = khive_runtime::VerbRegistryBuilder::new();
builder.register(OpsFileConcurrencyProbePack {
in_flight: in_flight.clone(),
max_in_flight: max_in_flight.clone(),
reject_overlap,
});
(
KhiveMcpServer::from_registry(builder.build().expect("probe registry")),
in_flight,
max_in_flight,
)
}
fn reader_probe_ops(count: usize) -> Vec<OpsFileEntry> {
(0..count)
.map(|sequence| OpsFileEntry {
tool: "reader_probe".to_string(),
args: serde_json::json!({"sequence": sequence}),
})
.collect()
}
#[test]
fn isolated_server_ignores_ambient_khive_packs_naming_unavailable_pack() {
const CHILD_MARKER: &str = "KKERNEL_KHIVE_PACKS_TEST_CHILD";
const TEST_NAME: &str =
"exec::tests::isolated_server_ignores_ambient_khive_packs_naming_unavailable_pack";
if std::env::var_os(CHILD_MARKER).is_none() {
let status = std::process::Command::new(
std::env::current_exe().expect("current test executable"),
)
.arg(TEST_NAME)
.arg("--exact")
.env("KHIVE_PACKS", "kg,gtd")
.env(CHILD_MARKER, "1")
.status()
.expect("spawn isolated KHIVE_PACKS test process");
assert!(status.success(), "isolated child test failed: {status}");
return;
}
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);
}
fn rerun_in_command_scoped_empty_home(child_marker: &str, test_name: &str) -> bool {
if std::env::var_os(child_marker).is_some() {
return false;
}
let calling_process_home = std::env::var_os("HOME");
let empty_home = tempfile::tempdir().expect("isolated child HOME");
let status =
std::process::Command::new(std::env::current_exe().expect("current test executable"))
.arg(test_name)
.arg("--exact")
.env("HOME", empty_home.path())
.env_remove("KHIVE_EMBEDDING_MODEL")
.env_remove("KHIVE_ADDITIONAL_EMBEDDING_MODELS")
.env_remove("KHIVE_ACTOR")
.env(child_marker, "1")
.status()
.expect("spawn isolated config-discovery test process");
assert_eq!(
std::env::var_os("HOME"),
calling_process_home,
"command-scoped HOME must not mutate the calling test process"
);
assert!(status.success(), "isolated child test failed: {status}");
true
}
#[test]
#[serial]
fn exec_config_id_matches_serve_config_id_for_project_toml_actor() {
const CHILD_MARKER: &str = "KKERNEL_EXEC_PROJECT_CONFIG_TEST_CHILD";
const TEST_NAME: &str =
"exec::tests::exec_config_id_matches_serve_config_id_for_project_toml_actor";
if rerun_in_command_scoped_empty_home(CHILD_MARKER, TEST_NAME) {
return;
}
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 pinned_packs = Some(vec!["kg".to_string()]);
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: pinned_packs.clone(),
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: pinned_packs,
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 actor_pin_rebuilds_visible_namespaces_dropping_displaced_fallback() {
const CHILD_MARKER: &str = "KKERNEL_ACTOR_PIN_CONFIG_TEST_CHILD";
const TEST_NAME: &str =
"exec::tests::actor_pin_rebuilds_visible_namespaces_dropping_displaced_fallback";
if rerun_in_command_scoped_empty_home(CHILD_MARKER, TEST_NAME) {
return;
}
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:fallback"
"#,
)
.expect("write config.toml");
let db_path = khive_dir.join("actor-pin-visibility-test.db");
let db_str = db_path.to_str().expect("utf8 path").to_string();
let pinned_packs = Some(vec!["kg".to_string()]);
let resolve = |db: &str| {
resolve_runtime_config(RuntimeConfigInputs {
db: Some(db),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: pinned_packs.clone(),
brain_profile: None,
})
.expect("resolve exec-shaped config")
};
let baseline = resolve(&db_str);
assert_eq!(baseline.actor_id.as_deref(), Some("lambda:fallback"));
assert!(baseline
.visible_namespaces
.contains(&Namespace::parse("lambda:fallback").expect("ns")));
let mut pinned_cfg = resolve(&db_str);
apply_actor_pin_and_expectation(&mut pinned_cfg, Some("lambda:pinned"), None).unwrap();
assert_eq!(pinned_cfg.actor_id.as_deref(), Some("lambda:pinned"));
assert!(
pinned_cfg
.visible_namespaces
.contains(&Namespace::parse("lambda:pinned").expect("ns")),
"pinned actor must be added to the default read scope: {:?}",
pinned_cfg.visible_namespaces
);
assert!(
!pinned_cfg
.visible_namespaces
.contains(&Namespace::parse("lambda:fallback").expect("ns")),
"the displaced fallback actor must not remain visible under the pinned \
identity: {:?}",
pinned_cfg.visible_namespaces
);
let mut local_cfg = resolve(&db_str);
apply_actor_pin_and_expectation(&mut local_cfg, Some("local"), None).unwrap();
assert_eq!(local_cfg.actor_id, None);
assert!(
local_cfg.visible_namespaces.is_empty(),
"pinning to local must leave only local visible, retaining neither the \
fallback actor nor adding a new one: {:?}",
local_cfg.visible_namespaces
);
}
#[test]
#[serial]
fn namespace_explicit_changes_actor_id_fill_but_not_config_id() {
if crate::test_process::run_in_child() {
return;
}
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::remove_var("KHIVE_ACTOR");
let empty_config_dir = tempfile::tempdir().expect("empty config tempdir");
let missing_config = empty_config_dir.path().join("config.toml");
std::fs::write(&missing_config, "").expect("write empty config");
let ns = Namespace::parse("lambda:custom-ns").expect("ns");
let pinned_packs = Some(vec!["kg".to_string()]);
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: pinned_packs.clone(),
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: pinned_packs,
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() {
if crate::test_process::run_in_child() {
return;
}
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 empty_config_dir = tempfile::tempdir().expect("empty config tempdir");
let missing_config = empty_config_dir.path().join("multi-backend-config.toml");
std::fs::write(&missing_config, "").expect("write empty config");
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,
served_kinds: 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,
served_kinds: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"session".to_string(),
PackConfig {
backend: "sessions".to_string(),
no_embed: false,
},
);
m
},
..KhiveConfig::default()
};
let pinned_packs = Some(vec!["kg".to_string()]);
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: pinned_packs.clone(),
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: pinned_packs,
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() {
if crate::test_process::run_in_child() {
return;
}
use khive_runtime::{BackendConfig, BackendKind, PackConfig};
let lock_dir = tempfile::tempdir().expect("private construction lock");
let _env = EnvAndCwdGuard::capture();
std::env::set_var("KHIVE_LOCK", lock_dir.path().join("khived.recovery.lock"));
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,
served_kinds: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Sqlite,
path: Some(secondary_path.clone()),
cache_mb: None,
journal_mode: None,
served_kinds: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "secondary".to_string(),
no_embed: false,
},
);
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(), "comm".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())
.await
.expect("multi-backend local fallback must build");
assert!(lock_dir.path().join("khived.recovery.lock").is_file());
let send = server
.dispatch_request_local(RequestParams {
plan: None,
ops: r#"comm.send(to="actor-routing-test", content="routed-via-secondary", self_send=true)"#
.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("comm.send must dispatch");
let send_resp: serde_json::Value = serde_json::from_str(&send).expect("valid JSON");
assert_eq!(
send_resp["results"][0]["ok"].as_bool(),
Some(true),
"comm.send must succeed through the multi-backend fallback server: {send_resp}"
);
async fn count_messages(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(), "comm".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 {
plan: None,
ops: r#"list(kind="message")"#.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"]["items"]
.as_array()
.map(|a| a.len())
.unwrap_or(0)
}
let main_count = count_messages(&main_path).await;
let secondary_count = count_messages(&secondary_path).await;
assert_eq!(
main_count, 0,
"comm 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, 2,
"comm pack write must land in its declared `secondary` backend file — \
`comm.send` dual-writes an outbound + inbound note copy per message \
(khive-pack-comm's message.rs), both via the SAME pack runtime, so a \
single self-send yields 2 `message` notes in whichever backend `comm` \
is pinned to"
);
}
#[tokio::test]
#[serial]
async fn build_local_fallback_server_uses_captured_anchor_after_home_changes() {
if crate::test_process::run_in_child() {
return;
}
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,
served_kinds: 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()).await;
restore_home(previous_home);
assert!(
result.is_ok(),
"exec fallback must use the anchor captured with RuntimeConfig after HOME changes: {}",
result.err().unwrap()
);
}
#[tokio::test]
#[serial]
async fn build_local_fallback_server_installs_blob_store_single_backend() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("tempdir");
let _env = EnvAndCwdGuard::capture();
std::env::set_var("KHIVE_LOCK", dir.path().join("khived.recovery.lock"));
let db_path = dir.path().join("exec_blob.db");
let cfg = RuntimeConfig {
db_path: Some(db_path),
embedding_model: None,
additional_embedding_models: vec![],
packs: RuntimeConfig::built_in_packs(),
..RuntimeConfig::default()
};
let khive_cfg = KhiveConfig::default();
let _server = build_local_fallback_server(cfg, &khive_cfg, None, None)
.await
.expect("single-backend local-exec construction must succeed");
assert!(dir.path().join("khived.recovery.lock").is_file());
assert!(
dir.path().join("blobs").is_dir(),
"default <db_dir>/blobs root must exist after construction, proving \
install_resolved_blob_store ran for the single-backend fallback path"
);
}
#[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![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
let khive_cfg = KhiveConfig::default();
let server = build_local_fallback_server(cfg, &khive_cfg, None, None)
.await
.expect("guarded local-exec construction must succeed");
for i in 0..count {
let params = RequestParams {
plan: None,
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]
fn build_local_fallback_server_blocks_while_recovery_lock_is_held() {
if crate::test_process::run_in_child() {
return;
}
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![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
let khive_cfg = KhiveConfig::default();
let result = rt_handle
.block_on(async { build_local_fallback_server(cfg, &khive_cfg, None, None).await });
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]
fn local_exec_construction_races_guarded_daemon_boot_without_fts_corruption() {
if crate::test_process::run_in_child() {
return;
}
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()).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()).unwrap_err();
assert_eq!(
err.downcast_ref::<ExecRefusal>().map(|error| error.reason),
Some(RefusalReason::ParseError)
);
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()).unwrap_err();
assert_eq!(
err.downcast_ref::<ExecRefusal>().map(|error| error.reason),
Some(RefusalReason::ParseError)
);
let msg = format!("{err:#}");
assert!(msg.contains("line 1"), "should report line number: {msg}");
}
#[test]
fn atomic_op_limit_is_checked_before_snapshot_materialization() {
struct PanicOnSnapshotAccess;
impl std::io::Read for PanicOnSnapshotAccess {
fn read(&mut self, _buf: &mut [u8]) -> std::io::Result<usize> {
panic!("over-limit atomic snapshot must not be read or materialized")
}
}
impl std::io::Seek for PanicOnSnapshotAccess {
fn seek(&mut self, _pos: std::io::SeekFrom) -> std::io::Result<u64> {
panic!("over-limit atomic snapshot must not be rewound or materialized")
}
}
let error = parse_atomic_validated_snapshot(&mut PanicOnSnapshotAccess, 2, 1)
.expect_err("the validated op count exceeds the configured atomic ceiling");
assert!(
error
.to_string()
.contains("op count 2 exceeds the configured maximum 1"),
"unexpected error: {error:#}"
);
}
#[test]
fn ops_file_physical_line_and_total_caps_are_fail_closed() {
let mut within = std::io::Cursor::new(b"1234567\n".to_vec());
assert_eq!(
read_bounded_ops_line_with_limit(&mut within, 1, 8)
.unwrap()
.unwrap(),
"1234567"
);
let mut over = std::io::Cursor::new(b"12345678\n".to_vec());
let error = read_bounded_ops_line_with_limit(&mut over, 7, 8).unwrap_err();
assert!(error.to_string().contains("line 7"));
let oversized = NamedTempFile::new().unwrap();
oversized.as_file().set_len(MAX_OPS_FILE_BYTES + 1).unwrap();
let error = validate_ops_file(oversized.path()).unwrap_err();
assert!(error.to_string().contains("total limit"));
}
#[test]
fn large_ops_file_payload_is_read_from_path_not_argv() {
let mut file = NamedTempFile::new().unwrap();
let payload = "x".repeat(1024 * 1024);
serde_json::to_writer(
&mut file,
&serde_json::json!({"tool":"stats","args":{"payload":payload}}),
)
.unwrap();
file.write_all(b"\n").unwrap();
let path = file.path().to_str().unwrap();
let args = ExecArgs::try_parse_from(["exec", "--ops-file", path]).unwrap();
assert!(args.ops.is_none());
assert_eq!(args.ops_file.as_deref(), Some(file.path()));
assert_eq!(validate_ops_file(file.path()).unwrap().total, 1);
}
#[test]
fn chunk_byte_boundary_is_exact() {
assert!(!should_defer_chunk_entry(
0,
0,
OPS_FILE_CHUNK_MAX_BYTES + 1
));
assert!(!should_defer_chunk_entry(
1,
OPS_FILE_CHUNK_MAX_BYTES - 1,
1
));
assert!(should_defer_chunk_entry(1, OPS_FILE_CHUNK_MAX_BYTES, 1));
}
#[test]
fn ordered_chunk_contract_rejects_tool_or_summary_drift() {
let tools = vec![
"first".to_string(),
"second".to_string(),
"third".to_string(),
];
let valid = serde_json::json!({
"results": [
{"ok":true,"tool":"first","result":{}},
{"ok":false,"tool":"second","error":"no"},
{"ok":false,"tool":"third","aborted":true,"error":"not attempted"}
],
"summary":{"total":3,"succeeded":1,"failed":1,"aborted":1},
"status":"partial"
});
assert_eq!(
validate_ordered_chunk_envelope(&tools, &valid, 1).unwrap(),
(1, 1, 1)
);
let mut wrong_tool = valid.clone();
wrong_tool["results"][1]["tool"] = serde_json::json!("third");
assert!(validate_ordered_chunk_envelope(&tools, &wrong_tool, 1).is_err());
let mut lying_summary = valid;
lying_summary["summary"]["succeeded"] = serde_json::json!(2);
lying_summary["summary"]["failed"] = serde_json::json!(0);
assert!(validate_ordered_chunk_envelope(&tools, &lying_summary, 1).is_err());
let mut missing_result = serde_json::json!({
"results": [
{"ok":true,"tool":"first"},
{"ok":false,"tool":"second","error":"no"},
{"ok":false,"tool":"third","aborted":true,"error":"not attempted"}
],
"summary":{"total":3,"succeeded":1,"failed":1,"aborted":1},
"status":"partial"
});
assert!(validate_ordered_chunk_envelope(&tools, &missing_result, 1).is_err());
missing_result["results"][0]["result"] = serde_json::Value::Null;
missing_result["results"][1]
.as_object_mut()
.unwrap()
.remove("error");
assert!(validate_ordered_chunk_envelope(&tools, &missing_result, 1).is_err());
let mut contradictory = serde_json::json!({
"results": [
{"ok":true,"tool":"first","result":null,"error":null},
{"ok":false,"tool":"second","error":"no","result":null},
{"ok":false,"tool":"third","aborted":true,"error":"not attempted"}
],
"summary":{"total":3,"succeeded":1,"failed":1,"aborted":1},
"status":"partial"
});
assert!(validate_ordered_chunk_envelope(&tools, &contradictory, 1).is_err());
contradictory["results"][0]
.as_object_mut()
.unwrap()
.remove("error");
assert!(validate_ordered_chunk_envelope(&tools, &contradictory, 1).is_err());
}
fn status_contract_fixture(status: &str) -> (Vec<String>, serde_json::Value) {
(
vec!["first".to_string(), "second".to_string()],
serde_json::json!({
"results": [
{"ok":true,"tool":"first","result":{}},
{"ok":false,"tool":"second","error":"no"}
],
"summary":{"total":2,"succeeded":1,"failed":1,"aborted":0},
"status":status
}),
)
}
#[test]
fn ordered_chunk_truthful_status_passes() {
let (ops, envelope) = status_contract_fixture("partial");
assert_eq!(
validate_ordered_chunk_envelope(&ops, &envelope, 1).unwrap(),
(1, 1, 0)
);
}
#[test]
fn ordered_chunk_contradicting_status_is_rejected() {
let (ops, envelope) = status_contract_fixture("success");
let error = validate_ordered_chunk_envelope(&ops, &envelope, 1).unwrap_err();
assert!(error.to_string().contains("status"));
}
#[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()).unwrap();
assert_eq!(ops.len(), 3);
let summary = apply_ops_file(&server, ops, None, None, None, false)
.await
.unwrap();
assert_eq!(summary["total"], 3);
assert_eq!(summary["succeeded"], 3);
assert_eq!(summary["failed"], 0);
assert!(summary.get("aborted").is_none());
assert!(summary.get("failure_details_omitted").is_none());
assert!(summary.get("results").is_none());
let params = RequestParams {
plan: None,
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"]["items"]
.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 serial_ops_file_max_in_flight_is_one_while_default_remains_parallel() {
let op_count = OPS_FILE_CHUNK_SIZE;
let (parallel_server, parallel_in_flight, parallel_max) = concurrency_probe_server(false);
let parallel_summary = apply_ops_file(
¶llel_server,
reader_probe_ops(op_count),
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
)
.await
.expect("the ordinary bounded-parallel batch must succeed");
assert_eq!(parallel_summary["succeeded"], op_count);
let observed_parallel_max = parallel_max.load(std::sync::atomic::Ordering::SeqCst);
assert_eq!(
observed_parallel_max, 8,
"default dispatch must retain the server's bounded parallelism"
);
assert_eq!(
parallel_in_flight.load(std::sync::atomic::Ordering::SeqCst),
0
);
let (serial_server, serial_in_flight, serial_max) = concurrency_probe_server(false);
let serial_summary = apply_ops_file_with_dispatch_mode(
&serial_server,
reader_probe_ops(op_count),
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
OpsFileDispatchMode::Serial,
)
.await
.expect("serial dispatch must succeed");
assert_eq!(serial_summary["succeeded"], op_count);
assert_eq!(
serial_max.load(std::sync::atomic::Ordering::SeqCst),
1,
"--serial must await every handler before starting the next"
);
assert_eq!(
serial_in_flight.load(std::sync::atomic::Ordering::SeqCst),
0
);
}
#[tokio::test]
async fn serial_ops_file_succeeds_with_a_reader_that_refuses_overlap() {
let op_count = 4;
let (parallel_server, _, parallel_max) = concurrency_probe_server(true);
let parallel_summary = apply_ops_file(
¶llel_server,
reader_probe_ops(op_count),
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
)
.await
.expect("one parallel probe succeeds, so partial failure remains in-band");
assert_eq!(parallel_summary["succeeded"], 1);
assert_eq!(parallel_summary["failed"], op_count - 1);
assert_eq!(
parallel_max.load(std::sync::atomic::Ordering::SeqCst),
op_count
);
assert!(parallel_summary["failures"]
.as_array()
.expect("failure rows")
.iter()
.all(|failure| failure["error"]["message"]
.as_str()
.expect("error.message is text")
.contains("sql_bridge.reader_open")));
let (serial_server, _, serial_max) = concurrency_probe_server(true);
let serial_summary = apply_ops_file_with_dispatch_mode(
&serial_server,
reader_probe_ops(op_count),
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
OpsFileDispatchMode::Serial,
)
.await
.expect("serial dispatch must avoid constrained-reader overlap");
assert_eq!(serial_summary["succeeded"], op_count);
assert_eq!(serial_summary["failed"], 0);
assert_eq!(serial_max.load(std::sync::atomic::Ordering::SeqCst), 1);
}
#[tokio::test]
async fn serial_ops_file_preserves_order_save_rows_and_strict_failure() {
let (server, _, max_in_flight) = concurrency_probe_server(false);
let mut ops = reader_probe_ops(3);
ops[1].args["fail"] = serde_json::json!(true);
let output_dir = tempfile::tempdir().expect("output dir");
let save_path = output_dir.path().join("serial-ordered.jsonl");
let error = apply_ops_file_with_dispatch_mode(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
true,
OpsFileDispatchMode::Serial,
)
.await
.expect_err("strict mode must report the middle handler failure");
assert!(error.to_string().contains("--strict"), "{error:#}");
assert_eq!(max_in_flight.load(std::sync::atomic::Ordering::SeqCst), 1);
let rows: Vec<serde_json::Value> = std::fs::read_to_string(save_path)
.expect("strict failure still publishes all confirmed ordered rows")
.lines()
.map(|line| serde_json::from_str(line).expect("JSON result row"))
.collect();
assert_eq!(rows.len(), 3);
assert_eq!(rows[0]["tool"], "reader_probe");
assert_eq!(rows[0]["ok"], true);
assert_eq!(rows[0]["result"]["sequence"], 0);
assert_eq!(rows[1]["tool"], "reader_probe");
assert_eq!(rows[1]["ok"], false);
assert_eq!(rows[1]["reason"], "strict-op-failure");
assert_eq!(rows[2]["tool"], "reader_probe");
assert_eq!(rows[2]["ok"], true);
assert_eq!(rows[2]["result"]["sequence"], 2);
}
#[tokio::test]
async fn serial_dispatch_consumes_the_preflighted_stable_snapshot() {
use std::io::{Seek as _, Write as _};
let mut source = NamedTempFile::new().expect("ops source");
for sequence in 0..2 {
serde_json::to_writer(
source.as_file_mut(),
&serde_json::json!({
"tool": "reader_probe",
"args": {"sequence": sequence},
}),
)
.expect("write source op");
source.write_all(b"\n").expect("write newline");
}
let mut validated = validate_ops_file(source.path()).expect("preflight source");
source.as_file_mut().set_len(0).expect("truncate source");
source
.as_file_mut()
.rewind()
.expect("rewind replacement source");
source
.write_all(b"{malformed replacement\n")
.expect("replace source after preflight");
source.flush().expect("flush replacement");
let (server, _, max_in_flight) = concurrency_probe_server(false);
let summary = apply_ops_file_reader_with_dispatch_mode(
&server,
std::io::BufReader::new(&mut validated.snapshot),
validated.total,
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
OpsFileDispatchMode::Serial,
)
.await
.expect("dispatch must consume the immutable validated snapshot");
assert_eq!(summary["total"], 2);
assert_eq!(summary["succeeded"], 2);
assert_eq!(max_in_flight.load(std::sync::atomic::Ordering::SeqCst), 1);
}
#[tokio::test]
async fn serial_typed_preflight_rejects_later_prev_before_first_write() {
if crate::test_process::run_in_child() {
return;
}
use std::io::Write as _;
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut source = NamedTempFile::new().expect("ops source");
source
.write_all(
b"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"must-not-land\"}}\n",
)
.expect("write valid first op");
source
.write_all(b"{\"tool\":\"stats\",\"args\":{\"probe\":\"$prev.id\"}}\n")
.expect("write typed-invalid later op");
let config_dir = tempfile::tempdir().expect("config dir");
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").expect("write empty config");
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(&db_path)),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
let error = run_exec_ops_file(
source.path().to_path_buf(),
cfg,
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
ExecDbContext {
raw: Some(db_path.clone()),
anchor: Some(PathBuf::from(&db_path)),
config: Some(config_path),
},
true,
false,
None,
false,
)
.await
.expect_err("later $prev must fail before the first serial handler starts");
assert!(error.to_string().contains("$prev"), "{error:#}");
let server = isolated_server(&db_path);
let response = dispatch_json(&server, r#"list(kind="concept")"#).await;
assert_eq!(
response["results"][0]["result"]["items"],
serde_json::json!([]),
"whole-snapshot typed preflight must prevent the valid first write"
);
}
#[tokio::test]
async fn serial_whole_snapshot_preflight_does_not_change_default_chunk_commit_parity() {
if crate::test_process::run_in_child() {
return;
}
use std::io::Write as _;
let mut source = NamedTempFile::new().expect("ops source");
for sequence in 0..OPS_FILE_CHUNK_SIZE {
serde_json::to_writer(
source.as_file_mut(),
&serde_json::json!({
"tool": "create",
"args": {
"kind": "concept",
"name": format!("default-first-chunk-{sequence}"),
},
}),
)
.expect("write valid first-chunk op");
source.write_all(b"\n").expect("write newline");
}
source
.write_all(b"{\"tool\":\"stats\",\"args\":{\"probe\":\"$prev.id\"}}\n")
.expect("write typed-invalid second-chunk op");
let config_dir = tempfile::tempdir().expect("config dir");
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").expect("write empty config");
let default_db = NamedTempFile::new().expect("default db");
let default_db_path = default_db.path().to_str().expect("utf8").to_string();
let default_error = run_exec_ops_file(
source.path().to_path_buf(),
RuntimeConfig {
db_path: Some(PathBuf::from(&default_db_path)),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
},
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
ExecDbContext {
raw: Some(default_db_path.clone()),
anchor: Some(PathBuf::from(&default_db_path)),
config: Some(config_path.clone()),
},
false,
false,
None,
false,
)
.await
.expect_err("default mode must reject the typed-invalid second chunk");
assert!(
default_error.to_string().contains("$prev"),
"{default_error:#}"
);
let default_server = isolated_server(&default_db_path);
let default_response =
dispatch_json(&default_server, r#"list(kind="concept", limit=200)"#).await;
assert_eq!(
default_response["results"][0]["result"]["items"]
.as_array()
.expect("default concept rows")
.len(),
OPS_FILE_CHUNK_SIZE,
"backward-compatible default mode must commit its first logical chunk"
);
let serial_db = NamedTempFile::new().expect("serial db");
let serial_db_path = serial_db.path().to_str().expect("utf8").to_string();
let serial_error = run_exec_ops_file(
source.path().to_path_buf(),
RuntimeConfig {
db_path: Some(PathBuf::from(&serial_db_path)),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
},
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
ExecDbContext {
raw: Some(serial_db_path.clone()),
anchor: Some(PathBuf::from(&serial_db_path)),
config: Some(config_path),
},
true,
false,
None,
false,
)
.await
.expect_err("serial mode must reject the later invalid op before dispatch");
assert!(
serial_error.to_string().contains("$prev"),
"{serial_error:#}"
);
let serial_server = isolated_server(&serial_db_path);
let serial_response =
dispatch_json(&serial_server, r#"list(kind="concept", limit=200)"#).await;
assert_eq!(
serial_response["results"][0]["result"]["items"],
serde_json::json!([]),
"serial whole-snapshot preflight must reject before its first write"
);
}
#[tokio::test]
async fn ops_file_write_leaves_no_wal_sidecar_after_return() {
if crate::test_process::run_in_child() {
return;
}
use std::io::Write as _;
let db_dir = tempfile::tempdir().expect("db dir");
let db_path = db_dir.path().join("close-order.db");
let db_path_str = db_path.to_str().expect("utf8").to_string();
let wal = {
let mut name = db_path.file_name().unwrap().to_os_string();
name.push("-wal");
db_path.parent().unwrap().join(name)
};
let shm = {
let mut name = db_path.file_name().unwrap().to_os_string();
name.push("-shm");
db_path.parent().unwrap().join(name)
};
let config_dir = tempfile::tempdir().expect("config dir");
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").expect("write empty config");
let mut source = NamedTempFile::new().expect("ops source");
source
.write_all(
b"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"wal-sidecar-witness\"}}\n",
)
.expect("write create op");
let cfg = RuntimeConfig {
db_path: Some(db_path.clone()),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
run_exec_ops_file(
source.path().to_path_buf(),
cfg,
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
ExecDbContext {
raw: Some(db_path_str),
anchor: Some(db_path),
config: Some(config_path),
},
false,
false,
None,
false,
)
.await
.expect("ops-file write must succeed");
assert!(
!wal.exists(),
"a daemonless write through run_exec_ops_file must not leave a -wal sidecar behind after return"
);
assert!(
!shm.exists(),
"a daemonless write through run_exec_ops_file must not leave a -shm sidecar behind after return"
);
}
#[tokio::test]
async fn oversized_single_ops_file_reaches_handler_validation() {
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 ops = vec![OpsFileEntry {
tool: "stats".to_string(),
args: serde_json::json!({
"payload": "x".repeat(khive_request::MAX_OPS_INPUT_LEN + 1),
}),
}];
let mut observed_handler_error = false;
let error = apply_ops_file_with_response_transform(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
None,
false,
|_, raw| {
let response: serde_json::Value =
serde_json::from_str(&raw).expect("handler response must be JSON");
assert_eq!(response["results"][0]["tool"], "stats");
assert_eq!(response["results"][0]["ok"], false);
assert!(
response["results"][0]["error"]
.to_string()
.contains("payload"),
"the oversized typed op must reach stats argument validation: {response}"
);
observed_handler_error = true;
raw
},
)
.await
.expect_err("the only op is intentionally invalid at the handler boundary");
assert!(
observed_handler_error,
"ops-file dispatch must not reapply the public 1 MiB raw-DSL limit"
);
assert!(error.to_string().contains("every op failed"));
}
#[tokio::test]
async fn oversized_multi_op_chunk_preserves_order_save_and_strict() {
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 payload = "x".repeat(600 * 1024);
let ops = vec![
OpsFileEntry {
tool: "create".to_string(),
args: serde_json::json!({
"kind": "concept",
"name": "oversized ordered success",
"description": payload,
}),
},
OpsFileEntry {
tool: "stats".to_string(),
args: serde_json::json!({
"payload": "y".repeat(600 * 1024),
}),
},
];
let encoded_len = serde_json::to_vec(&ops).unwrap().len();
assert!(encoded_len > khive_request::MAX_OPS_INPUT_LEN);
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("oversized-ordered.jsonl");
let error = apply_ops_file(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
true,
)
.await
.expect_err("strict mode must report the handler-level stats failure");
assert!(error.to_string().contains("--strict"), "{error:#}");
let rows: Vec<serde_json::Value> = std::fs::read_to_string(&save_path)
.expect("strict failure still publishes the complete ordered result file")
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0]["tool"], "create");
assert_eq!(rows[0]["ok"], true);
assert_eq!(rows[1]["tool"], "stats");
assert_eq!(rows[1]["ok"], false);
assert_eq!(rows[1]["reason"], "strict-op-failure");
}
#[tokio::test]
async fn public_dispatch_still_rejects_oversized_raw_ops() {
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 params = RequestParams {
plan: None,
ops: serde_json::json!({
"tool": "stats",
"args": {"payload": "x".repeat(khive_request::MAX_OPS_INPUT_LEN + 1)},
})
.to_string(),
presentation: Some("verbose".to_string()),
presentation_per_op: None,
save_to: None,
format: Some("json".to_string()),
format_per_op: None,
request_id: None,
};
let error = server
.dispatch_request_local(params)
.await
.expect_err("the raw request surface must retain its 1 MiB safety bound");
assert!(
error.to_string().contains("ops input is")
&& error.to_string().contains(&format!(
"max is {} bytes",
khive_request::MAX_OPS_INPUT_LEN
)),
"unexpected public dispatch error: {error}"
);
}
#[tokio::test]
async fn multi_chunk_save_retains_order_rows_checksum_and_json_override() {
use sha2::Digest as _;
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)
.with_default_output_format(khive_runtime::OutputFormat::Table);
let ops: Vec<OpsFileEntry> = (0..=OPS_FILE_CHUNK_SIZE)
.map(|index| OpsFileEntry {
tool: "create".to_string(),
args: serde_json::json!({
"kind": "concept",
"name": format!("ordered-{index:03}"),
}),
})
.collect();
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("ordered.jsonl");
let manifest = apply_ops_file(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
true,
)
.await
.unwrap();
let manifest_keys: std::collections::BTreeSet<_> = manifest
.as_object()
.unwrap()
.keys()
.map(String::as_str)
.collect();
assert_eq!(
manifest_keys,
std::collections::BTreeSet::from([
"checksum",
"path",
"per_column_null_counts",
"rows",
"schema_fingerprint",
"summary",
]),
"the successful manifest shape must remain unchanged"
);
assert_eq!(manifest["rows"], OPS_FILE_CHUNK_SIZE + 1);
assert_eq!(manifest["summary"]["total"], OPS_FILE_CHUNK_SIZE + 1);
assert_eq!(manifest["summary"]["succeeded"], OPS_FILE_CHUNK_SIZE + 1);
assert_eq!(manifest["summary"]["failed"], 0);
assert_eq!(manifest["summary"]["aborted"], 0);
let bytes = std::fs::read(&save_path).unwrap();
let checksum = format!("{:x}", sha2::Sha256::digest(&bytes));
assert_eq!(manifest["checksum"], checksum);
let rows: Vec<serde_json::Value> = bytes
.split(|byte| *byte == b'\n')
.filter(|line| !line.is_empty())
.map(|line| serde_json::from_slice(line).unwrap())
.collect();
assert_eq!(rows.len(), OPS_FILE_CHUNK_SIZE + 1);
for (index, row) in rows.iter().enumerate() {
assert_eq!(row["tool"], "create");
assert_eq!(row["ok"], true);
assert_eq!(row["result"]["name"], format!("ordered-{index:03}"));
}
}
#[tokio::test]
async fn malformed_later_chunk_emits_aborted_manifest_for_prior_commits() {
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 ops: Vec<OpsFileEntry> = (0..=OPS_FILE_CHUNK_SIZE)
.map(|index| OpsFileEntry {
tool: "create".to_string(),
args: serde_json::json!({
"kind": "concept",
"name": format!("abort-manifest-{index:03}"),
}),
})
.collect();
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("must-not-publish.jsonl");
let error = apply_ops_file_with_response_transform(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
true,
|chunk_number, raw| {
if chunk_number == 2 {
"{malformed-response".to_string()
} else {
raw
}
},
)
.await
.unwrap_err();
let aborted = error
.downcast_ref::<AbortedOpsFileError>()
.expect("post-dispatch failure must carry its emitted manifest");
assert_eq!(aborted.manifest["status"], "aborted");
assert_eq!(aborted.manifest["committed_chunks"], serde_json::json!([1]));
assert_eq!(aborted.manifest["dispatched_chunk"], 2);
assert_eq!(aborted.manifest["file_published"], false);
assert_eq!(
aborted.manifest["summary"]["succeeded"],
OPS_FILE_CHUNK_SIZE
);
assert_eq!(aborted.manifest["summary"]["total"], OPS_FILE_CHUNK_SIZE);
assert_eq!(aborted.manifest["summary"]["aborted"], 0);
assert_eq!(aborted.manifest["unconfirmed_ops"], 1);
assert!(
!save_path.exists(),
"an aborted run must not publish partial JSONL"
);
let params = RequestParams {
plan: None,
ops: r#"list(kind="concept", limit=200)"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: Some("json".to_string()),
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let response: serde_json::Value = serde_json::from_str(&raw).unwrap();
let rows = response["results"][0]["result"]["items"]
.as_array()
.unwrap_or_else(|| {
panic!(
"read-back list result must contain an items array; got {}",
response["results"][0]["result"]
)
});
let mut observed: Vec<String> = rows
.iter()
.map(|row| {
row["name"]
.as_str()
.unwrap_or_else(|| panic!("read-back row carries no string name: {row}"))
.to_owned()
})
.collect();
assert!(
!observed.is_empty(),
"read-back yielded no rows; result was {}",
response["results"][0]["result"]
);
observed.sort_unstable();
let committed: Vec<String> = (0..OPS_FILE_CHUNK_SIZE)
.map(|index| format!("abort-manifest-{index:03}"))
.collect();
let mut with_unconfirmed = committed.clone();
with_unconfirmed.push(format!("abort-manifest-{OPS_FILE_CHUNK_SIZE:03}"));
assert!(
observed == committed || observed == with_unconfirmed,
"manifest reports chunk 1 committed and chunk 2 unconfirmed, so the database must \
hold exactly the {OPS_FILE_CHUNK_SIZE} confirmed rows, optionally plus the one \
unconfirmed row; found {} rows: {observed:?}",
observed.len()
);
}
#[tokio::test]
async fn invalid_save_directory_is_rejected_before_any_op_side_effect() {
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 output_dir = tempfile::tempdir().unwrap();
let ops = vec![OpsFileEntry {
tool: "create".to_string(),
args: serde_json::json!({"kind":"concept","name":"must-not-exist"}),
}];
let error = apply_ops_file(
&server,
ops,
Some("verbose".to_string()),
Some("json".to_string()),
Some(output_dir.path().to_string_lossy().into_owned()),
true,
)
.await
.unwrap_err();
assert!(error
.to_string()
.contains("absent or an existing regular file"));
let params = RequestParams {
plan: None,
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: Some("json".to_string()),
format_per_op: None,
request_id: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let response: serde_json::Value = serde_json::from_str(&raw).unwrap();
assert_eq!(
response["results"][0]["result"]["items"],
serde_json::json!([])
);
}
#[tokio::test]
async fn save_manifest_preserves_partial_summary_and_strict_writes_rows() {
fn partial_ops(success_name: &str) -> Vec<OpsFileEntry> {
vec![
OpsFileEntry {
tool: "create".to_string(),
args: serde_json::json!({"kind":"concept","name":success_name}),
},
OpsFileEntry {
tool: "search".to_string(),
args: serde_json::json!({"kind":"not_a_real_kind","query":"x"}),
},
]
}
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 output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("partial.jsonl");
let manifest = apply_ops_file(
&server,
partial_ops("partial-ok"),
Some("verbose".to_string()),
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
false,
)
.await
.unwrap();
assert_eq!(manifest["rows"], 2);
assert_eq!(manifest["summary"]["succeeded"], 1);
assert_eq!(manifest["summary"]["failed"], 1);
assert_eq!(manifest["summary"]["aborted"], 0);
let strict_path = output_dir.path().join("strict-partial.jsonl");
let error = apply_ops_file(
&server,
partial_ops("strict-partial-ok"),
Some("verbose".to_string()),
Some("json".to_string()),
Some(strict_path.to_string_lossy().into_owned()),
true,
)
.await
.unwrap_err();
assert!(error.to_string().contains("--strict"));
assert_eq!(
std::fs::read_to_string(strict_path)
.unwrap()
.lines()
.count(),
2
);
}
#[test]
fn prepare_exec_output_preserves_specific_reasons_and_fills_strict_failures() {
let raw = serde_json::json!({
"results": [
{"ok": true, "tool": "stats", "result": {}},
{"ok": false, "tool": "get", "error": "missing id"},
{
"ok": false,
"tool": "not_loaded",
"error": "unknown verb",
"reason": "verb-refused"
},
{"ok": false, "tool": "update", "aborted": true},
],
"summary": {"total": 4, "succeeded": 1, "failed": 2, "aborted": 1},
"status": "partial",
})
.to_string();
let parsed: serde_json::Value =
serde_json::from_str(&prepare_exec_output(&raw, true)).unwrap();
assert_eq!(parsed["results"][1]["reason"], "strict-op-failure");
assert_eq!(parsed["results"][2]["reason"], "verb-refused");
assert_eq!(parsed["results"][3]["reason"], "strict-op-failure");
assert!(parsed["results"][0].get("reason").is_none());
}
#[test]
fn enforce_strict_batch_result_ok_when_strict_off_and_partially_failed() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 2, "succeeded": 1, "failed": 1, "aborted": 0},
})
.to_string();
assert!(enforce_strict_batch_result(&raw, false).is_ok());
}
#[test]
fn enforce_strict_batch_result_errs_when_strict_off_and_every_op_failed() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 1, "succeeded": 0, "failed": 1, "aborted": 0},
})
.to_string();
let err = enforce_strict_batch_result(&raw, false).unwrap_err();
assert!(format!("{err}").contains("every op failed"));
}
#[test]
fn enforce_strict_batch_result_errs_when_strict_off_and_chain_fully_aborted() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 3, "succeeded": 0, "failed": 1, "aborted": 2},
})
.to_string();
assert!(enforce_strict_batch_result(&raw, false).is_err());
}
#[test]
fn enforce_strict_batch_result_ok_on_empty_batch_summary() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 0, "succeeded": 0, "failed": 0, "aborted": 0},
})
.to_string();
assert!(enforce_strict_batch_result(&raw, false).is_ok());
assert!(enforce_strict_batch_result(&raw, true).is_ok());
}
#[test]
fn enforce_strict_batch_result_ok_when_strict_on_and_nothing_failed() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 2, "succeeded": 2, "failed": 0, "aborted": 0},
})
.to_string();
assert!(enforce_strict_batch_result(&raw, true).is_ok());
}
#[test]
fn enforce_strict_batch_result_errs_when_strict_on_and_a_failure_present() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 2, "succeeded": 1, "failed": 1, "aborted": 0},
})
.to_string();
let err = enforce_strict_batch_result(&raw, true).unwrap_err();
assert!(format!("{err}").contains("1 op(s) failed"));
}
#[test]
fn enforce_strict_batch_result_errs_when_strict_on_and_chain_aborted() {
let raw = serde_json::json!({
"results": [],
"summary": {"total": 2, "succeeded": 0, "failed": 1, "aborted": 1},
})
.to_string();
assert!(enforce_strict_batch_result(&raw, true).is_err());
}
#[test]
fn enforce_strict_batch_result_errs_on_save_manifest_with_failures() {
let raw = r#"{"path":"/tmp/out.jsonl","rows":2,"checksum":"ab","summary":{"total":2,"succeeded":1,"failed":1,"aborted":0}}"#;
assert!(enforce_strict_batch_result(raw, true).is_err());
let clean = r#"{"path":"/tmp/out.jsonl","rows":2,"checksum":"ab","summary":{"total":2,"succeeded":2,"failed":0,"aborted":0}}"#;
assert!(enforce_strict_batch_result(clean, true).is_ok());
}
#[test]
fn enforce_strict_batch_result_ok_on_non_json_output() {
assert!(enforce_strict_batch_result("| a | b |\n", true).is_ok());
}
#[tokio::test]
async fn apply_ops_file_strict_errs_when_an_op_fails() {
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 _;
f.write_all(
b"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"StrictOne\"}}\n",
)
.unwrap();
f.write_all(
b"{\"tool\":\"search\",\"args\":{\"kind\":\"not_a_real_kind\",\"query\":\"x\"}}\n",
)
.unwrap();
let ops = parse_ops_file(f.path()).unwrap();
assert_eq!(ops.len(), 2);
let err = apply_ops_file(&server, ops, None, None, None, true)
.await
.expect_err("strict mode must surface the per-op failure as a process error");
assert!(format!("{err}").contains("1 op(s) failed"));
}
#[tokio::test]
async fn apply_ops_file_errs_without_strict_when_every_op_fails() {
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 _;
f.write_all(
b"{\"tool\":\"search\",\"args\":{\"kind\":\"not_a_real_kind\",\"query\":\"x\"}}\n",
)
.unwrap();
let ops = parse_ops_file(f.path()).unwrap();
let err = apply_ops_file(&server, ops, None, None, None, false)
.await
.expect_err("a fully-failed ops-file must exit non-zero even without --strict");
assert!(format!("{err}").contains("every op failed"));
}
#[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 {
plan: None,
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", "status"]),
"non-atomic envelope must carry exactly results+summary+status, no `atomic` block (#1220 added `status`): {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() {
if crate::test_process::run_in_child() {
return;
}
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,
None,
None,
true,
ExecDbContext::default(),
false,
false,
None,
false,
)
.await
.unwrap();
let server = isolated_server(&db_path);
let params = RequestParams {
plan: None,
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"]["items"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(count, 0, "dry-run must not write any entities");
}
#[tokio::test]
async fn atomic_dry_run_uses_the_real_read_only_admission() {
if crate::test_process::run_in_child() {
return;
}
let dir = tempfile::tempdir().expect("temp dir");
let db_path = dir.path().join("dry-run-target.db");
let config_path = dir.path().join("khive.toml");
std::fs::write(&config_path, "").expect("empty config");
let source_path = dir.path().join("ops.jsonl");
std::fs::write(
&source_path,
"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"MustNotLand\"}}\n",
)
.expect("ops file");
let db_text = db_path.to_str().expect("utf8 db path").to_string();
let context = || ExecDbContext {
raw: Some(db_text.clone()),
anchor: Some(db_path.clone()),
config: Some(config_path.clone()),
};
for dry_run in [true, false] {
let error = run_exec_ops_file(
source_path.clone(),
atomic_cfg(&db_text),
None,
None,
None,
dry_run,
context(),
false,
true,
None,
false,
)
.await
.expect_err("create is inadmissible in an atomic unit");
assert!(
error
.downcast_ref::<crate::atomic_apply::AtomicExecFailure>()
.is_some(),
"dry_run={dry_run}: {error:#}"
);
assert!(!db_path.exists(), "preflight must not open the target db");
}
let over_limit = run_exec_ops_file(
source_path,
atomic_cfg(&db_text),
None,
None,
None,
true,
context(),
false,
true,
Some(0),
false,
)
.await
.expect_err("atomic dry-run must apply the configured op ceiling");
assert!(over_limit.to_string().contains("op count 1"));
assert!(
!db_path.exists(),
"over-limit preview must not open storage"
);
}
#[derive(Debug)]
struct DenyPinnedActorGate {
observed: std::sync::Arc<std::sync::Mutex<Vec<String>>>,
}
impl khive_runtime::Gate for DenyPinnedActorGate {
fn check(
&self,
req: &khive_runtime::GateRequest,
) -> std::result::Result<khive_runtime::GateDecision, khive_runtime::GateError> {
self.observed.lock().unwrap().push(req.actor.id.clone());
if req.actor.id == "lambda:pinned" {
Ok(khive_runtime::GateDecision::deny(
"test actor is not granted",
))
} else {
Ok(khive_runtime::GateDecision::allow())
}
}
}
#[tokio::test]
#[serial]
async fn unauthorized_explicit_actor_is_not_retried_as_fallback() {
if crate::test_process::run_in_child() {
return;
}
let previous_no_daemon = std::env::var("KHIVE_NO_DAEMON").ok();
std::env::set_var("KHIVE_NO_DAEMON", "1");
let (previous_home, _home_dir) = isolate_home_for_test();
let db_file = NamedTempFile::new().expect("temp db");
let observed = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
let mut cfg = RuntimeConfig {
db_path: Some(db_file.path().to_path_buf()),
actor_id: Some("lambda:fallback".to_string()),
gate: std::sync::Arc::new(DenyPinnedActorGate {
observed: observed.clone(),
}),
packs: vec!["kg".to_string()],
..RuntimeConfig::default()
};
apply_actor_pin_and_expectation(&mut cfg, Some("lambda:pinned"), Some("lambda:pinned"))
.unwrap();
let result = run_exec_inline(
r#"create(kind="concept", name="MustNotExist")"#.to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
)
.await;
match previous_no_daemon {
Some(value) => std::env::set_var("KHIVE_NO_DAEMON", value),
None => std::env::remove_var("KHIVE_NO_DAEMON"),
}
restore_home(previous_home);
assert!(result.is_err(), "the gate refusal must be terminal");
let observed = observed.lock().unwrap();
assert!(
!observed.is_empty(),
"the configured gate must be consulted"
);
assert!(
observed.iter().all(|actor| actor == "lambda:pinned"),
"no gate check may retry as the displaced fallback actor: {observed:?}"
);
}
#[tokio::test]
#[serial]
async fn strict_mode_rejects_before_daemon_forward_when_comm_and_no_actor() {
if crate::test_process::run_in_child() {
return;
}
let prev_strict = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
let prev_no_daemon = std::env::var("KHIVE_NO_DAEMON").ok();
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(), "comm".to_string()],
actor_id: None, ..RuntimeConfig::default()
};
let result = run_exec_inline(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
)
.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"),
}
assert!(
result.is_err(),
"run_exec_inline must return Err under strict mode + comm + no actor; got Ok"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
"error must name the strict-mode env var; got: {msg}"
);
assert!(
msg.contains("KHIVE_ACTOR"),
"error must name the remedy (KHIVE_ACTOR); got: {msg}"
);
}
#[tokio::test]
#[serial]
async fn strict_mode_allows_exec_when_comm_and_actor_configured() {
if crate::test_process::run_in_child() {
return;
}
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(), "comm".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(),
false,
)
.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_comm_no_actor() {
if crate::test_process::run_in_child() {
return;
}
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(), "comm".to_string()],
actor_id: None,
..RuntimeConfig::default()
};
let result = run_exec_inline(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
)
.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)]
#[tokio::test]
#[serial]
async fn resolved_pack_list_reaches_real_exec_adapter_boundary() {
if crate::test_process::run_in_child() {
return;
}
let (prev_home, _home_dir) = isolate_home_for_test();
khive_mcp::daemon::test_forward_seam::arm();
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string(), "gtd".to_string()],
actor_id: None,
..RuntimeConfig::default()
};
let result = run_exec_inline(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
)
.await;
restore_home(prev_home);
assert!(
result.is_ok(),
"the intercepted dispatch through the real production entry point must \
succeed: {result:?}"
);
assert_eq!(
khive_mcp::daemon::test_forward_seam::take_captured(),
Some(Some(vec!["kg".to_string(), "gtd".to_string()])),
"the real forward_or_spawn_boxed adapter in this crate must convert cfg.packs \
into Some(&packs) at the forward_or_spawn_with_config_and_packs call boundary"
);
}
#[cfg(unix)]
std::thread_local! {
static SPY_WAS_CALLED: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
}
#[cfg(unix)]
fn spy_forward_records_call<'a>(
_frame: &'a DaemonRequestFrame,
_config: Option<PathBuf>,
_db: Option<&'a str>,
_packs: Vec<String>,
) -> super::ForwardFuture<'a> {
SPY_WAS_CALLED.with(|c| c.set(true));
Box::pin(async { None })
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn strict_mode_spy_confirms_enforce_fires_before_forward() {
if crate::test_process::run_in_child() {
return;
}
let prev_strict = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::remove_var("KHIVE_NO_DAEMON");
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
SPY_WAS_CALLED.with(|c| c.set(false));
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string(), "comm".to_string()],
actor_id: None, ..RuntimeConfig::default()
};
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None, None,
ExecDbContext::default(),
false,
spy_forward_records_call,
)
.await;
let spy_was_called = SPY_WAS_CALLED.with(|c| c.get());
SPY_WAS_CALLED.with(|c| c.set(false));
match prev_strict {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
assert!(
result.is_err(),
"strict mode + comm + no actor must return Err; got Ok"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
"error must name the strict-mode env var; got: {msg}"
);
assert!(
!spy_was_called,
"spy forward_fn was called — enforce_strict_actor_mode fired AFTER forwarding, not before"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn strict_mode_spy_forward_reached_when_actor_configured() {
if crate::test_process::run_in_child() {
return;
}
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");
SPY_WAS_CALLED.with(|c| c.set(false));
let cfg = RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string(), "comm".to_string()],
actor_id: Some("lambda:tenant-x".to_string()), ..RuntimeConfig::default()
};
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None, None,
ExecDbContext::default(),
false,
spy_forward_records_call,
)
.await;
let spy_was_called = SPY_WAS_CALLED.with(|c| c.get());
SPY_WAS_CALLED.with(|c| c.set(false));
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(),
"gate must pass when actor is configured; got: {result:?}"
);
assert!(
spy_was_called,
"spy forward_fn must be called when gate passes (KHIVE_NO_DAEMON=1 causes in-process fallback)"
);
}
#[cfg(unix)]
std::thread_local! {
static SPY_CAPTURED_CONFIG_ID: std::cell::RefCell<Option<String>> =
const { std::cell::RefCell::new(None) };
static SPY_CAPTURED_CONFIG_PATH: std::cell::RefCell<Option<PathBuf>> =
const { std::cell::RefCell::new(None) };
static SPY_CAPTURED_DB: std::cell::RefCell<Option<String>> =
const { std::cell::RefCell::new(None) };
static SPY_CAPTURED_PACKS: std::cell::RefCell<Option<Vec<String>>> =
const { std::cell::RefCell::new(None) };
}
#[cfg(unix)]
fn spy_capture_config_id<'a>(
frame: &'a DaemonRequestFrame,
config: Option<PathBuf>,
db: Option<&'a str>,
packs: Vec<String>,
) -> super::ForwardFuture<'a> {
SPY_CAPTURED_CONFIG_ID.with(|c| *c.borrow_mut() = Some(frame.config_id.clone()));
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = config);
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = db.map(str::to_string));
SPY_CAPTURED_PACKS.with(|c| *c.borrow_mut() = Some(packs));
Box::pin(async { None })
}
#[cfg(unix)]
fn spy_capture_config_and_succeed<'a>(
frame: &'a DaemonRequestFrame,
config: Option<PathBuf>,
db: Option<&'a str>,
packs: Vec<String>,
) -> super::ForwardFuture<'a> {
SPY_CAPTURED_CONFIG_ID.with(|c| *c.borrow_mut() = Some(frame.config_id.clone()));
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = config);
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = db.map(str::to_string));
SPY_CAPTURED_PACKS.with(|c| *c.borrow_mut() = Some(packs));
Box::pin(async {
Some(Ok(
r#"{"results":[{"ok":true,"tool":"stats","result":{}}],"summary":{"total":1,"succeeded":1,"failed":0}}"#
.to_string(),
))
})
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn explicit_config_reaches_daemon_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
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");
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = None);
let dir = tempfile::tempdir().expect("config tempdir");
let config_path = dir.path().join("selected.toml");
std::fs::write(&config_path, "[runtime]\npacks = [\"kg\"]\n")
.expect("write explicit config");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: None,
config: Some(&config_path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext {
raw: None,
anchor: None,
config: Some(config_path.clone()),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"forwarded dispatch must succeed: {result:?}"
);
assert_eq!(
SPY_CAPTURED_CONFIG_PATH.with(|captured| captured.borrow_mut().take()),
Some(config_path),
"the daemon spawn seam must receive the same explicit config path used to resolve the exec frame"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn memory_db_override_reaches_daemon_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
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");
std::env::remove_var("KHIVE_DB");
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = None);
let dir = tempfile::tempdir().expect("config tempdir");
let config_path = dir.path().join("selected.toml");
std::fs::write(&config_path, "[runtime]\npacks = [\"kg\"]\n")
.expect("write explicit config");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&config_path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext {
raw: Some(":memory:".to_string()),
anchor: None,
config: Some(config_path.clone()),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"forwarded dispatch must succeed: {result:?}"
);
assert_eq!(
SPY_CAPTURED_DB.with(|captured| captured.borrow_mut().take()),
Some(":memory:".to_string()),
"the daemon spawn seam must receive the raw --db override so a spawned daemon \
can be constructed with the same ephemeral in-memory storage"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn force_memory_exec_frame_matches_opened_read_only_topology_runtime() {
if crate::test_process::run_in_child() {
return;
}
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");
std::env::remove_var("KHIVE_DB");
let (prev_home, _home_dir) = isolate_home_for_test();
SPY_CAPTURED_CONFIG_ID.with(|captured| *captured.borrow_mut() = None);
let fixture = tempfile::tempdir().expect("force-memory config tempdir");
let config_path = fixture.path().join("read-only-topology.toml");
let declared_main = fixture.path().join("declared-main.db");
let declared_archive = fixture.path().join("declared-archive.db");
std::fs::write(
&config_path,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
read_only = true
[[backends]]
name = "archive"
kind = "sqlite"
path = "{}"
"#,
declared_main.display(),
declared_archive.display(),
),
)
.expect("write read-only topology config");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&config_path),
namespace: Namespace::local(),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve force-memory exec config");
assert_eq!(
cfg.db_path, None,
"the force-memory anchor must be in-memory"
);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg.clone(),
None,
None,
None,
ExecDbContext {
raw: Some(":memory:".to_string()),
anchor: None,
config: Some(config_path.clone()),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(result.is_ok(), "force-memory dispatch failed: {result:?}");
let frame_config_id = SPY_CAPTURED_CONFIG_ID
.with(|captured| captured.borrow_mut().take())
.expect("spy must capture the forwarded config id");
let khive_cfg = KhiveConfig::load_with_home_fallback(Some(&config_path), None)
.expect("load force-memory topology")
.expect("explicit config must exist");
let opened =
khive_mcp::serve::build_registry_for_multi_backend(cfg, &khive_cfg, Some(":memory:"))
.await
.expect("force-memory runtime must build");
restore_home(prev_home);
assert!(
!opened.default_runtime.is_read_only(),
"force-memory replaces the declared read-only SQLite main with writable memory"
);
assert_eq!(
frame_config_id, opened.config_id,
"the pre-open exec frame and opened force-memory runtime must have identical config ids"
);
assert!(
!declared_main.exists() && !declared_archive.exists(),
"force-memory parity setup must not materialize either declared SQLite path"
);
}
#[cfg(unix)]
fn write_writable_multi_backend_config(
config_path: &Path,
main_path: &Path,
secondary_path: &Path,
) {
std::fs::write(
config_path,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "archive"
kind = "sqlite"
path = "{}"
"#,
main_path.display(),
secondary_path.display(),
),
)
.unwrap();
}
#[cfg(unix)]
fn chmod_read_only(path: &Path) {
use std::os::unix::fs::PermissionsExt;
let mut permissions = std::fs::metadata(path).unwrap().permissions();
permissions.set_mode(0o444);
std::fs::set_permissions(path, permissions).unwrap();
for suffix in ["-wal", "-shm"] {
let mut name = path.file_name().unwrap().to_os_string();
name.push(suffix);
let sidecar = path.parent().unwrap().join(name);
if sidecar.exists() {
let mut sidecar_permissions = std::fs::metadata(&sidecar).unwrap().permissions();
sidecar_permissions.set_mode(0o444);
std::fs::set_permissions(&sidecar, sidecar_permissions).unwrap();
}
}
}
#[cfg(unix)]
fn runtime_config_for_explicit_multi_backend(
config_path: &Path,
db: Option<&str>,
) -> RuntimeConfig {
resolve_runtime_config(RuntimeConfigInputs {
db,
config: Some(config_path),
namespace: Namespace::local(),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.unwrap()
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn multi_backend_main_chmod_refuses_before_daemon_forward() {
if crate::test_process::run_in_child() {
return;
}
std::env::remove_var("KHIVE_DB");
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
SPY_CAPTURED_CONFIG_ID.with(|captured| *captured.borrow_mut() = None);
let fixture = tempfile::tempdir().unwrap();
let config_path = fixture.path().join("khive.toml");
let main_path = fixture.path().join("main.db");
let archive_path = fixture.path().join("archive.db");
std::fs::write(&main_path, b"main snapshot fixture").unwrap();
std::fs::write(&archive_path, b"archive fixture").unwrap();
chmod_read_only(&main_path);
write_writable_multi_backend_config(&config_path, &main_path, &archive_path);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
runtime_config_for_explicit_multi_backend(&config_path, None),
None,
None,
None,
ExecDbContext {
raw: None,
anchor: None,
config: Some(config_path),
},
false,
spy_capture_config_and_succeed,
)
.await;
let error = result.expect_err("an undeclared main snapshot mode must fail closed");
assert!(error.to_string().contains("read_only = true"), "{error}");
assert!(
SPY_CAPTURED_CONFIG_ID.with(|captured| captured.borrow().is_none()),
"the retained writable daemon must never receive the frame"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn multi_backend_secondary_chmod_refuses_before_daemon_forward() {
if crate::test_process::run_in_child() {
return;
}
std::env::remove_var("KHIVE_DB");
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
SPY_CAPTURED_CONFIG_ID.with(|captured| *captured.borrow_mut() = None);
let fixture = tempfile::tempdir().unwrap();
let config_path = fixture.path().join("khive.toml");
let main_path = fixture.path().join("main.db");
let archive_path = fixture.path().join("archive.db");
std::fs::write(&main_path, b"main fixture").unwrap();
std::fs::write(&archive_path, b"archive snapshot fixture").unwrap();
chmod_read_only(&archive_path);
write_writable_multi_backend_config(&config_path, &main_path, &archive_path);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
runtime_config_for_explicit_multi_backend(&config_path, None),
None,
None,
None,
ExecDbContext {
raw: None,
anchor: None,
config: Some(config_path),
},
false,
spy_capture_config_and_succeed,
)
.await;
let error = result.expect_err("an undeclared secondary snapshot mode must fail closed");
assert!(error.to_string().contains("archive"), "{error}");
assert!(
SPY_CAPTURED_CONFIG_ID.with(|captured| captured.borrow().is_none()),
"the retained writable daemon must never receive the frame"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn force_memory_skips_declared_chmod_preflight_and_forwards() {
if crate::test_process::run_in_child() {
return;
}
std::env::remove_var("KHIVE_DB");
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
SPY_CAPTURED_CONFIG_ID.with(|captured| *captured.borrow_mut() = None);
let fixture = tempfile::tempdir().unwrap();
let config_path = fixture.path().join("khive.toml");
let main_path = fixture.path().join("main.db");
let archive_path = fixture.path().join("archive.db");
std::fs::write(&main_path, b"unused main snapshot fixture").unwrap();
std::fs::write(&archive_path, b"unused archive snapshot fixture").unwrap();
chmod_read_only(&main_path);
chmod_read_only(&archive_path);
write_writable_multi_backend_config(&config_path, &main_path, &archive_path);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
runtime_config_for_explicit_multi_backend(&config_path, Some(":memory:")),
None,
None,
None,
ExecDbContext {
raw: Some(":memory:".to_string()),
anchor: None,
config: Some(config_path),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"force-memory forwarding must remain valid: {result:?}"
);
assert!(
SPY_CAPTURED_CONFIG_ID.with(|captured| captured.borrow().is_some()),
"the force-memory frame must reach the forwarding seam"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn single_backend_concrete_db_override_reaches_daemon_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
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");
std::env::remove_var("KHIVE_DB");
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = None);
let dir = tempfile::tempdir().expect("config tempdir");
let config_path = dir.path().join("selected.toml");
std::fs::write(&config_path, "[runtime]\npacks = [\"kg\"]\n")
.expect("write explicit config");
let override_path = dir.path().join("override.db");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(override_path.to_str().expect("utf8")),
config: Some(&config_path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext {
raw: Some(override_path.display().to_string()),
anchor: khive_runtime::resolve_db_anchor(override_path.to_str()),
config: Some(config_path),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"forwarded dispatch must succeed: {result:?}"
);
assert_eq!(
SPY_CAPTURED_DB.with(|captured| captured.borrow_mut().take()),
Some(override_path.display().to_string()),
"the single-backend concrete override must reach the daemon spawn seam so a \
spawned daemon binds the operator's file instead of the default database"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn env_khive_packs_reaches_daemon_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
let _guard = EnvAndCwdGuard::capture();
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");
std::env::remove_var("KHIVE_DB");
std::env::set_var("KHIVE_PACKS", "kg,gtd,memory");
SPY_CAPTURED_PACKS.with(|c| *c.borrow_mut() = None);
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 from KHIVE_PACKS env");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"forwarded dispatch must succeed: {result:?}"
);
assert_eq!(
SPY_CAPTURED_PACKS.with(|captured| captured.borrow_mut().take()),
Some(vec![
"kg".to_string(),
"gtd".to_string(),
"memory".to_string()
]),
"the KHIVE_PACKS-resolved pack list must reach the daemon spawn seam so a \
spawned daemon serves the same packs this client resolved"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn no_env_control_forwards_built_in_default_packs_to_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
let _guard = EnvAndCwdGuard::capture();
let home_dir = tempfile::tempdir().expect("tempdir for isolated HOME");
let empty_project_root = tempfile::tempdir().expect("empty project-root tempdir");
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");
std::env::remove_var("KHIVE_DB");
std::env::remove_var("KHIVE_PACKS");
std::env::set_var("HOME", home_dir.path());
std::env::set_current_dir(empty_project_root.path())
.expect("chdir into isolated project root with no discoverable config");
SPY_CAPTURED_PACKS.with(|c| *c.borrow_mut() = None);
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 with no pack-selection input");
assert_eq!(
cfg.packs,
RuntimeConfig::built_in_packs(),
"the isolated no-selection environment must resolve to the built-in default \
pack set before the forwarding seam is exercised"
);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"forwarded dispatch must succeed: {result:?}"
);
assert_eq!(
SPY_CAPTURED_PACKS.with(|captured| captured.borrow_mut().take()),
Some(RuntimeConfig::built_in_packs()),
"with no pack-selection input the client must forward exactly the built-in \
default pack set — identical to what an independently-spawned daemon would \
have defaulted to on its own"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn redundant_db_override_forwards_discovered_config_to_spawn_seam() {
if crate::test_process::run_in_child() {
return;
}
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");
std::env::remove_var("KHIVE_DB");
let (prev_home, _home_dir) = isolate_home_for_test();
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = None);
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = None);
let backend_dir = tempfile::tempdir().expect("backend tempdir");
let main_backend_path = backend_dir.path().join("main-backend.db");
let sessions_backend_path = backend_dir.path().join("sessions-backend.db");
let anchor_dir = backend_dir.path().join(".khive");
std::fs::create_dir_all(&anchor_dir).expect("mkdir db-dir anchor");
let discovered_config_path = anchor_dir.join("config.toml");
std::fs::write(
&discovered_config_path,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "sessions"
kind = "sqlite"
path = "{}"
"#,
main_backend_path.display(),
sessions_backend_path.display(),
),
)
.expect("write tier-3 multi-backend config");
let canonical_config_path =
std::fs::canonicalize(&discovered_config_path).expect("canonicalize config path");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(main_backend_path.to_str().expect("utf8")),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg.clone(),
None,
None,
None,
ExecDbContext {
raw: Some(main_backend_path.display().to_string()),
anchor: khive_runtime::resolve_db_anchor(main_backend_path.to_str()),
config: None,
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"redundant-override dispatch must reach daemon forwarding: {result:?}"
);
assert_eq!(
SPY_CAPTURED_DB.with(|captured| captured.borrow_mut().take()),
None,
"the redundant override stays withheld from the spawn seam"
);
assert_eq!(
SPY_CAPTURED_CONFIG_PATH.with(|captured| captured.borrow_mut().take()),
Some(canonical_config_path.clone()),
"the spawn seam must receive the retained resolved config path as the \
child's explicit --config when the redundant override is withheld"
);
let explicit_config_path = backend_dir.path().join("explicit.toml");
std::fs::copy(&discovered_config_path, &explicit_config_path)
.expect("copy topology as explicit config");
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = None);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg.clone(),
None,
None,
None,
ExecDbContext {
raw: Some(main_backend_path.display().to_string()),
anchor: khive_runtime::resolve_db_anchor(main_backend_path.to_str()),
config: Some(explicit_config_path.clone()),
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"explicit-config dispatch must reach daemon forwarding: {result:?}"
);
assert_eq!(
SPY_CAPTURED_CONFIG_PATH.with(|captured| captured.borrow_mut().take()),
Some(explicit_config_path),
"with an explicit config the seam receives the operator's path, never a discovered one"
);
let single_dir = tempfile::tempdir().expect("single-backend tempdir");
let override_path = single_dir.path().join("override.db");
let single_cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(override_path.to_str().expect("utf8")),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve single-backend exec-shaped config");
SPY_CAPTURED_CONFIG_PATH.with(|c| *c.borrow_mut() = None);
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = None);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
single_cfg,
None,
None,
None,
ExecDbContext {
raw: Some(override_path.display().to_string()),
anchor: khive_runtime::resolve_db_anchor(override_path.to_str()),
config: None,
},
false,
spy_capture_config_and_succeed,
)
.await;
assert!(
result.is_ok(),
"single-backend dispatch must reach daemon forwarding: {result:?}"
);
assert_eq!(
SPY_CAPTURED_CONFIG_PATH.with(|captured| captured.borrow_mut().take()),
None,
"the empty-backends case forwards no config — the concrete override supplies the database"
);
assert_eq!(
SPY_CAPTURED_DB.with(|captured| captured.borrow_mut().take()),
Some(override_path.display().to_string()),
"the single-backend concrete override still reaches the spawn seam"
);
restore_home(prev_home);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn explicit_config_is_loaded_for_exec_forward_frame() {
if crate::test_process::run_in_child() {
return;
}
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(|captured| *captured.borrow_mut() = None);
let fixture = tempfile::tempdir().expect("config fixture tempdir");
let config_path = fixture.path().join("code-map.toml");
let main_backend_path = fixture.path().join("code-map.db");
let sessions_backend_path = fixture.path().join("sessions.db");
std::fs::write(
&config_path,
format!(
r#"
[[backends]]
name = "main"
kind = "sqlite"
path = "{}"
[[backends]]
name = "sessions"
kind = "sqlite"
path = "{}"
"#,
main_backend_path.display(),
sessions_backend_path.display(),
),
)
.expect("write explicit exec config");
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: None,
config: Some(&config_path),
namespace: Namespace::parse("local").expect("namespace"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve explicit exec config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg.clone(),
None,
None,
None,
ExecDbContext {
raw: None,
anchor: None,
config: Some(config_path.clone()),
},
false,
spy_capture_config_id,
)
.await;
assert!(
result.is_ok(),
"explicit-config dispatch failed: {result:?}"
);
let captured = SPY_CAPTURED_CONFIG_ID
.with(|value| value.borrow_mut().take())
.expect("spy must capture the forwarded config id");
let khive_cfg = KhiveConfig::load_with_home_fallback(Some(&config_path), None)
.expect("load explicit config")
.expect("explicit config must exist");
let mut expected_config = cfg;
expected_config.db_path = Some(main_backend_path);
let expected = compute_config_id(&expected_config, Some(&khive_cfg));
restore_home(prev_home);
assert_eq!(
captured, expected,
"the exec forward frame must fold the explicitly selected backend topology"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn exec_frame_config_id_matches_daemon_config_id_for_multi_backend_project_toml() {
if crate::test_process::run_in_child() {
return;
}
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 backend_dir = tempfile::tempdir().expect("backend tempdir");
let main_backend_path = backend_dir.path().join("main-backend.db");
let sessions_backend_path = backend_dir.path().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: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext::default(),
false,
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: Some(vec!["kg".to_string()]),
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 mut declared_main_config = serve_cfg;
declared_main_config.db_path = Some(main_backend_path);
let daemon_config_id = compute_config_id(&declared_main_config, 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)"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn declared_topology_exec_fallback_updates_omitted_db_anchor() {
if crate::test_process::run_in_child() {
return;
}
let fixture = tempfile::tempdir().expect("config fixture");
let unused_home = fixture.path().join("unused-home");
let anchor = unused_home.join(".khive/khive.db");
let config_path = fixture.path().join("memory-topology.toml");
std::fs::write(
&config_path,
"[[backends]]\nname = \"main\"\nkind = \"memory\"\n",
)
.expect("write explicit memory topology");
let cfg = RuntimeConfig {
db_path: Some(anchor.clone()),
actor_id: Some("config-anchor-fallback-test".to_string()),
packs: vec!["kg".to_string()],
brain_profile: None,
..RuntimeConfig::no_embeddings()
};
let topology = KhiveConfig::load_with_home_fallback(Some(&config_path), None)
.unwrap()
.expect("explicit topology exists");
let mut expected_cfg = cfg.clone();
expected_cfg.db_path = None;
let expected_id = compute_config_id(&expected_cfg, Some(&topology));
SPY_CAPTURED_CONFIG_ID.with(|captured| *captured.borrow_mut() = None);
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
Some("json".to_string()),
None,
ExecDbContext {
raw: None,
anchor: Some(anchor),
config: Some(config_path),
},
true,
spy_capture_config_id,
)
.await;
assert!(result.is_ok(), "local fallback must dispatch: {result:?}");
assert_eq!(
SPY_CAPTURED_CONFIG_ID.with(|captured| captured.borrow_mut().take()),
Some(expected_id),
"the forward attempt must use declared memory before falling back"
);
assert!(
!unused_home.exists(),
"the superseded HOME-shaped database anchor must never be opened"
);
}
#[cfg(unix)]
#[tokio::test]
#[serial]
async fn inline_db_override_guard_normalizes_main_config_id_and_rejects_conflict() {
if crate::test_process::run_in_child() {
return;
}
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 backend_dir = tempfile::tempdir().expect("backend tempdir");
let main_backend_path = backend_dir.path().join("main-backend.db");
let sessions_backend_path = backend_dir.path().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 no_override_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: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config without override");
let no_override_result = run_exec_inline_with_forward(
"stats()".to_string(),
no_override_cfg,
None,
None,
None,
ExecDbContext::default(),
false,
spy_capture_config_id,
)
.await;
assert!(
no_override_result.is_ok(),
"no-override dispatch must succeed: {no_override_result:?}"
);
let no_override_config_id = SPY_CAPTURED_CONFIG_ID
.with(|captured| captured.borrow_mut().take())
.expect("no-override frame must be captured");
let matching_override = main_backend_path.display().to_string();
SPY_CAPTURED_DB.with(|c| *c.borrow_mut() = Some("sentinel".to_string()));
let cfg = resolve_runtime_config(RuntimeConfigInputs {
db: Some(&matching_override),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
actor_explicit: false,
no_embed: true,
packs: Some(vec!["kg".to_string()]),
brain_profile: None,
})
.expect("resolve exec-shaped config");
let matching_result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg.clone(),
None,
None,
None,
ExecDbContext {
raw: Some(matching_override.clone()),
anchor: khive_runtime::resolve_db_anchor(Some(&matching_override)),
config: None,
},
false,
spy_capture_config_id,
)
.await;
assert!(
matching_result.is_ok(),
"an override matching the declared main backend must reach daemon forwarding: {matching_result:?}"
);
let matching_config_id = SPY_CAPTURED_CONFIG_ID
.with(|captured| captured.borrow_mut().take())
.expect("matching-override frame must be captured");
assert_eq!(
SPY_CAPTURED_DB.with(|captured| captured.borrow_mut().take()),
None,
"the redundant multi-backend concrete override must be WITHHELD from the spawn \
seam: the frame's fingerprint is normalized to the no-override anchor, and the \
spawned daemon's config-declared main path IS the override's target"
);
let conflicting_override = backend_dir.path().join("override.db");
let result = run_exec_inline_with_forward(
"stats()".to_string(),
cfg,
None,
None,
None,
ExecDbContext {
raw: Some(conflicting_override.display().to_string()),
anchor: None,
config: None,
},
false,
spy_capture_config_id,
)
.await;
restore_home(prev_home);
assert!(
result.is_err(),
"a --db/KHIVE_DB override that conflicts with a declared [[backends]] topology \
must be rejected on the inline path too, not only on --ops-file; got: {result:?}"
);
assert!(
SPY_CAPTURED_CONFIG_ID.with(|c| c.borrow().is_none()),
"the conflict must be caught BEFORE any daemon-forward attempt — the spy must \
never have been called"
);
assert_eq!(
matching_config_id, no_override_config_id,
"a matching --db override must emit the same config_id as no override for the same multi-backend config"
);
}
#[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 {
plan: None,
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"]["items"]
.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 {
plan: None,
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![],
packs: vec!["kg".to_string()],
..Default::default()
}
}
async fn replace_fts_entities_with_incompatible_table(db_path: &str) {
let runtime = KhiveRuntime::new(atomic_cfg(db_path)).expect("runtime for FTS fault setup");
let sql = runtime.sql();
let mut writer = sql.writer().await.expect("writer for FTS fault setup");
writer
.execute_script(
"DROP TABLE fts_entities; \
CREATE TABLE fts_entities (broken_column TEXT);"
.to_string(),
)
.await
.expect("replace FTS table to inject post-commit reindex failure");
}
#[tokio::test]
async fn atomic_kg_only_config_keeps_gtd_hook_and_lifecycle_execution() {
if crate::test_process::run_in_child() {
return;
}
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 (hook_task_id, transition_task_id, complete_task_id) = {
let server = isolated_server(&db_path);
let response = dispatch_json(
&server,
r#"[gtd.assign(title="HookGuard", status="next"), gtd.assign(title="TransitionGuard", status="inbox"), gtd.assign(title="CompleteGuard", status="active")]"#,
)
.await;
let full_id = |index: usize| {
response["results"][index]["result"]["full_id"]
.as_str()
.unwrap_or_else(|| panic!("missing task full_id at index {index}: {response}"))
.to_string()
};
(full_id(0), full_id(1), full_id(2))
};
let hook_error = crate::atomic_apply::execute_atomic_ops_file(
vec![atomic_op(
"update",
serde_json::json!({
"id": hook_task_id.as_str(),
"properties": {"depends_on": [hook_task_id.as_str()]},
}),
)],
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("the GTD task hook must reject a self-dependency");
assert!(
format!("{hook_error:#}").contains("cannot depend on itself"),
"the kg-only atomic registry must enforce the GTD hook: {hook_error:#}"
);
let server = isolated_server(&db_path);
let response = dispatch_json(&server, &format!(r#"get(id="{hook_task_id}")"#)).await;
assert!(
response["results"][0]["result"]["properties"]
.get("depends_on")
.is_none(),
"the rejected dependency update must not mutate the task: {response}"
);
let envelope = crate::atomic_apply::execute_atomic_ops_file(
vec![
atomic_op(
"gtd.transition",
serde_json::json!({"id": transition_task_id, "status": "next"}),
),
atomic_op(
"gtd.complete",
serde_json::json!({"id": complete_task_id, "result": "verified"}),
),
],
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("GTD lifecycle adapters must execute with a kg-only config");
assert_eq!(envelope["atomic"]["committed"], true, "{envelope}");
assert_eq!(envelope["results"][0]["result"]["to"], "next");
assert_eq!(envelope["results"][1]["result"]["to"], "done");
}
#[tokio::test]
async fn atomic_ops_file_success_commits_all_ops() {
if crate::test_process::run_in_child() {
return;
}
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_post_commit_reindex_failure_returns_committed_non_retryable_envelope() {
if crate::test_process::run_in_child() {
return;
}
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 response = dispatch_json(
&server,
r#"create(kind="concept", name="PostCommitReindexBefore")"#,
)
.await;
response["results"][0]["result"]["id"]
.as_str()
.expect("entity id")
.to_string()
};
replace_fts_entities_with_incompatible_table(&db_path).await;
let envelope = crate::atomic_apply::execute_atomic_ops_file(
vec![atomic_op(
"update",
serde_json::json!({"id": entity_id, "name": "PostCommitReindexAfter"}),
)],
atomic_cfg(&db_path),
&KhiveConfig::default(),
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("post-commit failure must return a reconciliation envelope");
assert_eq!(envelope["atomic"]["committed"], true, "{envelope}");
assert_eq!(
envelope["atomic"]["status"], "committed_degraded",
"{envelope}"
);
assert_eq!(envelope["atomic"]["retryable"], false, "{envelope}");
assert_eq!(
envelope["atomic"]["degradations"][0]["stage"], "post_commit_reindex",
"{envelope}"
);
assert_eq!(envelope["summary"]["succeeded"], 1, "{envelope}");
let runtime =
KhiveRuntime::new(atomic_cfg(&db_path)).expect("runtime for committed-row check");
let token = runtime.authorize(Namespace::local()).expect("authorize");
let entity = runtime
.get_entity(&token, Uuid::parse_str(&entity_id).unwrap())
.await
.expect("committed entity row");
assert_eq!(entity.name, "PostCommitReindexAfter");
}
#[tokio::test]
async fn atomic_result_read_failure_keeps_later_committed_delete_non_retryable() {
if crate::test_process::run_in_child() {
return;
}
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 response = dispatch_json(
&server,
r#"create(kind="concept", name="RenderThenDelete")"#,
)
.await;
response["results"][0]["result"]["id"]
.as_str()
.expect("entity id")
.to_string()
};
let envelope = crate::atomic_apply::execute_atomic_ops_file(
vec![
atomic_op(
"update",
serde_json::json!({"id": entity_id, "name": "NeverRendered"}),
),
atomic_op("delete", serde_json::json!({"id": entity_id, "hard": true})),
],
atomic_cfg(&db_path),
&KhiveConfig::default(),
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("render failure after commit must be a reconciliation envelope");
assert_eq!(envelope["atomic"]["committed"], true, "{envelope}");
assert_eq!(
envelope["atomic"]["status"], "committed_degraded",
"{envelope}"
);
assert_eq!(envelope["atomic"]["retryable"], false, "{envelope}");
assert_eq!(
envelope["atomic"]["degradations"][0]["stage"], "result_rendering",
"{envelope}"
);
assert_eq!(envelope["results"][0]["ok"], true, "{envelope}");
assert_eq!(
envelope["results"][0]["status"], "committed_degraded",
"{envelope}"
);
assert_eq!(envelope["results"][0]["retryable"], false, "{envelope}");
assert_eq!(envelope["results"][1]["result"]["deleted"], true);
let runtime =
KhiveRuntime::new(atomic_cfg(&db_path)).expect("runtime for deleted-row check");
let token = runtime.authorize(Namespace::local()).expect("authorize");
let deleted = runtime
.get_entity(&token, Uuid::parse_str(&entity_id).unwrap())
.await;
assert!(deleted.is_err(), "hard delete must remain committed");
}
#[tokio::test]
async fn atomic_ops_file_mid_unit_failure_rolls_back_whole_unit() {
if crate::test_process::run_in_child() {
return;
}
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_rollback_exits_nonzero_with_or_without_save_and_strict() {
if crate::test_process::run_in_child() {
return;
}
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 response = dispatch_json(
&server,
r#"[create(kind="concept", name="AtomicExitX"), create(kind="concept", name="AtomicExitY")]"#,
)
.await;
(
response["results"][0]["result"]["id"]
.as_str()
.unwrap()
.to_string(),
response["results"][1]["result"]["id"]
.as_str()
.unwrap()
.to_string(),
)
};
let mut file = NamedTempFile::new().unwrap();
for op in [
serde_json::json!({"tool":"delete","args":{"id":x_id.clone(),"hard":true}}),
serde_json::json!({
"tool":"link",
"args":{"source_id":y_id,"target_id":x_id,"relation":"extends"}
}),
] {
serde_json::to_writer(&mut file, &op).unwrap();
file.write_all(b"\n").unwrap();
}
let config_dir = tempfile::tempdir().unwrap();
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").unwrap();
let db_context = || ExecDbContext {
raw: Some(db_path.clone()),
anchor: Some(PathBuf::from(&db_path)),
config: Some(config_path.clone()),
};
let error = run_exec_ops_file(
file.path().to_path_buf(),
atomic_cfg(&db_path),
None,
None,
None,
false,
db_context(),
false,
true,
None,
false,
)
.await
.expect_err("atomic rollback must exit non-zero without --strict");
assert!(error.to_string().contains("rolled back"), "{error:#}");
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("atomic-rollback.jsonl");
let error = run_exec_ops_file(
file.path().to_path_buf(),
atomic_cfg(&db_path),
None,
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
false,
db_context(),
false,
true,
None,
true,
)
.await
.expect_err("atomic rollback must exit non-zero with --strict and --save-file");
assert!(error.to_string().contains("rolled back"), "{error:#}");
assert_eq!(
std::fs::read_to_string(save_path).unwrap().lines().count(),
2
);
}
#[test]
fn atomic_save_persist_failure_returns_committed_reconciliation_stdout() {
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("publish-race.jsonl");
let sink = khive_mcp::save_sink::JsonlSaveSink::new(&save_path, false)
.expect("preflight save sink");
std::fs::create_dir(&save_path).expect("occupy destination with a directory");
let mut envelope = serde_json::json!({
"results": [{
"ok": true,
"tool": "update",
"op_index": 0,
"result": {"id": Uuid::new_v4()}
}],
"summary": {"total": 1, "succeeded": 1, "failed": 0},
"atomic": {
"committed": true,
"rolled_back": false,
"failed_op_index": null,
"error": null
}
});
let failure = render_atomic_output(&mut envelope, Some(sink))
.expect_err("persist failure must remain a non-zero CLI outcome");
let stdout: serde_json::Value =
serde_json::from_str(&failure.stdout).expect("structured stdout envelope");
assert_eq!(stdout["atomic"]["committed"], true, "{stdout}");
assert_eq!(stdout["atomic"]["status"], "committed_degraded", "{stdout}");
assert_eq!(stdout["atomic"]["retryable"], false, "{stdout}");
assert_eq!(
stdout["atomic"]["degradations"][0]["stage"], "save_file_publish",
"{stdout}"
);
assert!(
stdout["atomic"]["degradations"][0]["error"]
.as_str()
.is_some_and(|error| error.contains("persist temp file")),
"{stdout}"
);
assert!(
format!("{:#}", failure.error).contains("do not replay the mutation"),
"terminal error must point automation at reconciliation"
);
}
#[tokio::test]
async fn atomic_invalid_save_directory_is_rejected_before_commit() {
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut file = NamedTempFile::new().unwrap();
serde_json::to_writer(
&mut file,
&serde_json::json!({
"tool":"create",
"args":{"kind":"concept","name":"atomic-must-not-exist"}
}),
)
.unwrap();
file.write_all(b"\n").unwrap();
let config_dir = tempfile::tempdir().unwrap();
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").unwrap();
let save_directory = tempfile::tempdir().unwrap();
let error = run_exec_ops_file(
file.path().to_path_buf(),
atomic_cfg(&db_path),
None,
Some("json".to_string()),
Some(save_directory.path().to_string_lossy().into_owned()),
false,
ExecDbContext {
raw: Some(db_path.clone()),
anchor: Some(PathBuf::from(&db_path)),
config: Some(config_path),
},
false,
true,
None,
true,
)
.await
.unwrap_err();
assert!(error
.to_string()
.contains("absent or an existing regular file"));
let server = isolated_server(&db_path);
let response = dispatch_json(&server, r#"list(kind="concept")"#).await;
assert_eq!(
response["results"][0]["result"]["items"],
serde_json::json!([])
);
}
#[tokio::test]
async fn atomic_preflighted_save_keeps_prior_file_on_execution_error() {
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut file = NamedTempFile::new().unwrap();
serde_json::to_writer(
&mut file,
&serde_json::json!({
"tool":"create",
"args":{"kind":"concept","name":"atomic-rejected"}
}),
)
.unwrap();
file.write_all(b"\n").unwrap();
let config_dir = tempfile::tempdir().unwrap();
let config_path = config_dir.path().join("khive.toml");
std::fs::write(&config_path, "").unwrap();
let output_dir = tempfile::tempdir().unwrap();
let save_path = output_dir.path().join("prior.jsonl");
std::fs::write(&save_path, b"prior-complete-output\n").unwrap();
let error = run_exec_ops_file(
file.path().to_path_buf(),
atomic_cfg(&db_path),
None,
Some("json".to_string()),
Some(save_path.to_string_lossy().into_owned()),
false,
ExecDbContext {
raw: Some(db_path),
anchor: Some(PathBuf::from(db_file.path())),
config: Some(config_path),
},
false,
true,
Some(1),
true,
)
.await
.unwrap_err();
assert!(
error.to_string().contains("--atomic rejected"),
"unexpected atomic execution error: {error:#}"
);
assert_eq!(
std::fs::read(&save_path).unwrap(),
b"prior-complete-output\n"
);
}
#[tokio::test]
async fn atomic_ops_file_rejects_same_unit_gtd_dependency_cycles() {
if crate::test_process::run_in_child() {
return;
}
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 response = dispatch_json(
&server,
r#"[gtd.assign(title="AtomicCycleA", status="next"), gtd.assign(title="AtomicCycleB", status="next")]"#,
)
.await;
(
response["results"][0]["result"]["full_id"]
.as_str()
.expect("task A id")
.to_string(),
response["results"][1]["result"]["full_id"]
.as_str()
.expect("task B id")
.to_string(),
)
};
let compact_a_id = a_id.replace('-', "");
let compact_b_id = b_id.replace('-', "");
let alternate_spelling_error = crate::atomic_apply::execute_atomic_ops_file(
vec![
atomic_op(
"update",
serde_json::json!({
"id": a_id.clone(),
"properties": {"depends_on": [compact_b_id]}
}),
),
atomic_op(
"update",
serde_json::json!({
"id": b_id.clone(),
"properties": {"depends_on": [compact_a_id]}
}),
),
],
atomic_cfg(&db_path),
&KhiveConfig::default(),
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect_err("atomic preparation must reject an alternate dependency UUID spelling");
let alternate_spelling_message = format!("{alternate_spelling_error:#}");
assert!(
alternate_spelling_message.contains("canonical lowercase hyphenated UUID"),
"unexpected alternate-spelling error: {alternate_spelling_message}"
);
{
let server = isolated_server(&db_path);
let response = dispatch_json(&server, &format!(r#"get(id="{a_id}")"#)).await;
assert!(
response["results"][0]["result"]["properties"]
.get("depends_on")
.is_none(),
"alternate dependency spelling must not persist: {response}"
);
}
let property_envelope = crate::atomic_apply::execute_atomic_ops_file(
vec![
atomic_op(
"update",
serde_json::json!({
"id": a_id.clone(),
"properties": {"depends_on": [b_id.clone()]}
}),
),
atomic_op(
"update",
serde_json::json!({
"id": b_id.clone(),
"properties": {"depends_on": [a_id.clone()]}
}),
),
],
atomic_cfg(&db_path),
&KhiveConfig::default(),
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("cycle is a clean atomic rollback, not a seam failure");
assert_eq!(property_envelope["atomic"]["rolled_back"], true);
assert_eq!(property_envelope["atomic"]["failed_op_index"], 1);
assert!(
property_envelope["atomic"]["error"]
.as_str()
.is_some_and(|error| error.contains("dependency cycle")),
"envelope: {property_envelope}"
);
let server = isolated_server(&db_path);
for task_id in [&a_id, &b_id] {
let response = dispatch_json(&server, &format!(r#"get(id="{task_id}")"#)).await;
assert!(
response["results"][0]["result"]["properties"]
.get("depends_on")
.is_none(),
"the earlier update must roll back too: {response}"
);
}
let edge_envelope = crate::atomic_apply::execute_atomic_ops_file(
vec![
atomic_op(
"link",
serde_json::json!({
"source_id": a_id.clone(),
"target_id": b_id.clone(),
"relation": "depends_on"
}),
),
atomic_op(
"link",
serde_json::json!({
"source_id": b_id.clone(),
"target_id": a_id.clone(),
"relation": "depends_on"
}),
),
],
atomic_cfg(&db_path),
&KhiveConfig::default(),
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("edge cycle is a clean atomic rollback, not a seam failure");
assert_eq!(edge_envelope["atomic"]["rolled_back"], true);
assert_eq!(edge_envelope["atomic"]["failed_op_index"], 1);
assert!(
edge_envelope["atomic"]["error"]
.as_str()
.is_some_and(|error| error.contains("dependency cycle")),
"envelope: {edge_envelope}"
);
let server = isolated_server(&db_path);
let response = dispatch_json(
&server,
&format!(r#"neighbors(id="{a_id}", direction="out", relations=["depends_on"])"#),
)
.await;
assert_eq!(
response["results"][0]["result"],
serde_json::json!([]),
"the earlier link must roll back too: {response}"
);
}
#[tokio::test]
async fn atomic_symmetric_update_absorbs_into_same_unit_link_and_renders_correct_id() {
if crate::test_process::run_in_child() {
return;
}
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.6,
"ADR-039 DO NOTHING: the surviving row keeps its OWN pre-existing weight (0.6, \
set by the link above), not the discarded update's patched weight (0.9): {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.6,
"the committed row itself must keep its pre-existing weight, not the discarded \
update's patch: {surviving_resp}"
);
}
#[tokio::test]
async fn atomic_symmetric_update_absorbs_into_pre_existing_tombstoned_survivor_and_renders_it()
{
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let (canonical_id, x_id) = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"[create(kind="concept", name="TombA"), create(kind="concept", name="TombB")]"#,
)
.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="competes_with", weight=0.6)"#
),
)
.await;
let canonical_id = link_resp["results"][0]["result"]["id"]
.as_str()
.expect("canonical id")
.to_string();
dispatch_json(&server, &format!(r#"delete(id="{canonical_id}")"#)).await;
let x_resp = dispatch_json(
&server,
&format!(
r#"link(source_id="{a_id}", target_id="{b_id}", relation="extends", weight=0.2)"#
),
)
.await;
let x_id = x_resp["results"][0]["result"]["id"]
.as_str()
.expect("x id")
.to_string();
(canonical_id, x_id)
};
let ops = vec![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 by absorbing into the tombstoned survivor");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
let rendered_id = envelope["results"][0]["result"]["id"]
.as_str()
.expect("update result id")
.to_string();
assert_eq!(
rendered_id, canonical_id,
"must render the pre-existing tombstoned canonical survivor, not X's stale \
requested id: {envelope}"
);
assert!(
!envelope["results"][0]["result"]["deleted_at"].is_null(),
"the rendered survivor must show its OWN tombstoned state (non-null deleted_at) \
— absorbing a conflicting update must not resurrect it: {envelope}"
);
assert_eq!(
envelope["results"][0]["result"]["weight"], 0.6,
"ADR-039 DO NOTHING: the survivor keeps its own pre-existing weight, not X's \
discarded patched weight (0.9): {envelope}"
);
}
#[tokio::test]
async fn atomic_cli_boundary_rejections_happen_before_any_write() {
if crate::test_process::run_in_child() {
return;
}
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"]["items"]
.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"]["items"]
.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"]["items"]
.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() {
if crate::test_process::run_in_child() {
return;
}
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() {
if crate::test_process::run_in_child() {
return;
}
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_task_note_update_projects_status_like_canonical_handlers() {
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let task_id = {
let server = isolated_server(&db_path);
let response = dispatch_json(
&server,
r#"gtd.assign(title="AtomicProjectionTask", status="next")"#,
)
.await;
response["results"][0]["result"]["full_id"]
.as_str()
.expect("task id")
.to_string()
};
let khive_cfg = KhiveConfig::default();
let update = || {
vec![atomic_op(
"update",
serde_json::json!({"id": task_id.clone(), "content": "projected body"}),
)]
};
let changed = crate::atomic_apply::execute_atomic_ops_file(
update(),
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("atomic note update");
let changed_result = &changed["results"][0]["result"];
assert_eq!(changed_result["status"], "next", "{changed}");
assert_eq!(changed_result["lifecycle"], "active", "{changed}");
assert_eq!(
changed_result["display_name"], "AtomicProjectionTask",
"{changed}"
);
assert!(changed_result.get("unchanged").is_none(), "{changed}");
let server = isolated_server(&db_path);
let canonical_get = dispatch_json(&server, &format!(r#"get(id="{task_id}")"#)).await;
assert_eq!(changed_result, &canonical_get["results"][0]["result"]);
let unchanged = crate::atomic_apply::execute_atomic_ops_file(
update(),
atomic_cfg(&db_path),
&khive_cfg,
khive_types::pack::ATOMIC_MAX_OPS_DEFAULT,
)
.await
.expect("atomic note no-op update");
let unchanged_result = &unchanged["results"][0]["result"];
assert_eq!(unchanged_result["unchanged"], true, "{unchanged}");
let canonical_update = dispatch_json(
&server,
&format!(r#"update(id="{task_id}", content="projected body")"#),
)
.await;
assert_eq!(unchanged_result, &canonical_update["results"][0]["result"]);
}
#[tokio::test]
async fn atomic_delete_unknown_field_rejected_well_formed_succeeds() {
if crate::test_process::run_in_child() {
return;
}
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() {
if crate::test_process::run_in_child() {
return;
}
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_gtd_transition_unknown_field_rejected_well_formed_succeeds() {
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let task_id = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"gtd.assign(title="TransitionTypoGuard", status="inbox")"#,
)
.await;
resp["results"][0]["result"]["full_id"]
.as_str()
.expect("full_id")
.to_string()
};
let khive_cfg = KhiveConfig::default();
let ops = vec![atomic_op(
"gtd.transition",
serde_json::json!({"id": task_id, "status": "next", "notee": "typo"}),
)];
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 `notee` must be rejected");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let ops = vec![atomic_op(
"gtd.transition",
serde_json::json!({"id": task_id, "status": "next"}),
)];
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 gtd.transition must succeed");
assert_eq!(
envelope["atomic"]["committed"], true,
"envelope: {envelope}"
);
}
#[tokio::test]
async fn atomic_gtd_complete_unknown_field_rejected_well_formed_succeeds() {
if crate::test_process::run_in_child() {
return;
}
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let task_id = {
let server = isolated_server(&db_path);
let resp = dispatch_json(
&server,
r#"gtd.assign(title="CompleteTypoGuard", status="next")"#,
)
.await;
resp["results"][0]["result"]["full_id"]
.as_str()
.expect("full_id")
.to_string()
};
let khive_cfg = KhiveConfig::default();
let ops = vec![atomic_op(
"gtd.complete",
serde_json::json!({"id": task_id, "resutl": "typo"}),
)];
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 `resutl` must be rejected");
assert!(
format!("{err:#}").contains("unknown field"),
"error: {err:#}"
);
let ops = vec![atomic_op(
"gtd.complete",
serde_json::json!({"id": task_id, "result": "shipped"}),
)];
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 gtd.complete 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() {
if crate::test_process::run_in_child() {
return;
}
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() {
if crate::test_process::run_in_child() {
return;
}
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_gtd_transition_accept_8_hex_prefix_ids() {
if crate::test_process::run_in_child() {
return;
}
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, task_full_id) = {
let server = isolated_server(&db_path);
let resp =
dispatch_json(&server, r#"create(kind="concept", name="PrefixEntity")"#).await;
let entity_id = resp["results"][0]["result"]["id"]
.as_str()
.expect("entity id")
.to_string();
let resp =
dispatch_json(&server, r#"gtd.assign(title="PrefixTask", status="next")"#).await;
let task_id = resp["results"][0]["result"]["full_id"]
.as_str()
.expect("task full_id")
.to_string();
(entity_id, task_id)
};
let entity_prefix = &entity_full_id[..8];
let task_prefix = &task_full_id[..8];
let ops = vec![
atomic_op(
"update",
serde_json::json!({"id": entity_prefix, "name": "PrefixEntity-renamed"}),
),
atomic_op(
"gtd.transition",
serde_json::json!({"id": task_prefix, "status": "active"}),
),
];
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() {
if crate::test_process::run_in_child() {
return;
}
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, transition_task_id, complete_task_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()
};
let resp = dispatch_json(
&server,
r#"gtd.assign(title="ResultTransitionTask", status="next")"#,
)
.await;
let transition_task_id = resp["results"][0]["result"]["full_id"]
.as_str()
.expect("task full_id")
.to_string();
let resp = dispatch_json(
&server,
r#"gtd.assign(title="ResultCompleteTask", status="active")"#,
)
.await;
let complete_task_id = resp["results"][0]["result"]["full_id"]
.as_str()
.expect("task full_id")
.to_string();
(
id(0),
id(1),
id(2),
id(3),
transition_task_id,
complete_task_id,
)
};
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",
}),
),
atomic_op(
"gtd.transition",
serde_json::json!({"id": transition_task_id, "status": "active"}),
),
atomic_op(
"gtd.complete",
serde_json::json!({"id": complete_task_id, "result": "shipped"}),
),
];
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("all five 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(), 5, "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}"
);
assert_eq!(
results[3]["result"]["transitioned"], true,
"gtd.transition result: {envelope}"
);
assert_eq!(
results[3]["result"]["to"], "active",
"gtd.transition result must carry the new status: {envelope}"
);
assert_eq!(
results[4]["result"]["completed"], true,
"gtd.complete result: {envelope}"
);
assert_eq!(
results[4]["result"]["to"], "done",
"gtd.complete result must carry the terminal status: {envelope}"
);
}
}