use std::collections::BTreeMap;
use std::io::BufRead as _;
use std::path::PathBuf;
use anyhow::{Context, Result};
use clap::Parser;
use khive_mcp::serve::{apply_env_output_format, enforce_strict_actor_mode};
#[cfg(unix)]
use khive_mcp::server::compute_config_id;
use khive_mcp::server::KhiveMcpServer;
use khive_mcp::tools::request::RequestParams;
#[cfg(unix)]
use khive_runtime::{daemon::PROTOCOL_VERSION, DaemonRequestFrame};
use khive_runtime::{KhiveConfig, KhiveRuntime, Namespace, RuntimeConfig};
#[cfg(unix)]
type ForwardFuture<'a> = std::pin::Pin<
Box<dyn std::future::Future<Output = Option<Result<String, rmcp::ErrorData>>> + Send + 'a>,
>;
#[cfg(unix)]
type ForwardFnPtr = for<'a> fn(&'a DaemonRequestFrame) -> ForwardFuture<'a>;
#[cfg(unix)]
fn forward_or_spawn_boxed(frame: &DaemonRequestFrame) -> ForwardFuture<'_> {
Box::pin(khive_mcp::daemon::forward_or_spawn(frame))
}
use crate::dbpath::resolve_db_override;
use crate::pending_events;
const OPS_FILE_CHUNK_SIZE: usize = 100;
#[derive(Parser, Debug)]
pub struct ExecArgs {
pub ops: Option<String>,
#[arg(long, conflicts_with = "ops", conflicts_with = "ops_file")]
pub pending_events: bool,
#[arg(long, env = "KHIVE_DB")]
pub db: Option<String>,
#[arg(long, default_value = "local")]
pub namespace: String,
#[arg(long, default_value = "verbose")]
pub presentation: Option<String>,
#[arg(long, value_name = "FORMAT")]
pub output_format: Option<String>,
#[arg(long, short = 'v')]
pub verbose: bool,
#[arg(long)]
pub save_file: Option<String>,
#[arg(long, value_name = "PATH")]
pub ops_file: Option<PathBuf>,
#[arg(long, requires = "ops_file")]
pub dry_run: bool,
}
#[derive(Debug)]
struct OpsFileEntry {
tool: String,
args: serde_json::Value,
}
fn parse_ops_file(path: &PathBuf) -> Result<Vec<OpsFileEntry>> {
let file =
std::fs::File::open(path).with_context(|| format!("open ops-file {}", path.display()))?;
let reader = std::io::BufReader::new(file);
let mut ops: Vec<OpsFileEntry> = Vec::new();
for (line_idx, result) in reader.lines().enumerate() {
let line_num = line_idx + 1;
let raw = result.with_context(|| format!("read ops-file line {line_num}"))?;
let trimmed = raw.trim();
if trimmed.is_empty() {
continue;
}
let obj: serde_json::Value = serde_json::from_str(trimmed)
.map_err(|e| anyhow::anyhow!("ops-file line {line_num}: invalid JSON: {e}"))?;
let obj = obj.as_object().ok_or_else(|| {
anyhow::anyhow!(
"ops-file line {line_num}: expected a JSON object {{\"tool\":...,\"args\":...}}, \
got a non-object value"
)
})?;
let tool = obj
.get("tool")
.and_then(|v| v.as_str())
.ok_or_else(|| {
anyhow::anyhow!("ops-file line {line_num}: missing or non-string \"tool\" field")
})?
.to_owned();
let args = match obj.get("args") {
None => serde_json::Value::Object(serde_json::Map::new()),
Some(v) => {
if !v.is_object() {
anyhow::bail!(
"ops-file line {line_num}: \"args\" must be a JSON object, got {v}"
);
}
v.clone()
}
};
ops.push(OpsFileEntry { tool, args });
}
Ok(ops)
}
async fn apply_ops_file(
server: &KhiveMcpServer,
ops: Vec<OpsFileEntry>,
presentation: Option<String>,
) -> Result<()> {
let total = ops.len();
let mut total_succeeded: usize = 0;
let mut total_failed: usize = 0;
for (chunk_idx, chunk) in ops.chunks(OPS_FILE_CHUNK_SIZE).enumerate() {
let applied_before = (chunk_idx * OPS_FILE_CHUNK_SIZE).min(total);
let batch_arr: Vec<serde_json::Value> = chunk
.iter()
.map(|e| {
serde_json::json!({
"tool": e.tool,
"args": e.args,
})
})
.collect();
let batch_json = serde_json::to_string(&batch_arr).context("serialize chunk to JSON")?;
let params = RequestParams {
ops: batch_json,
presentation: presentation.clone(),
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
};
let raw = server
.dispatch_request_local(params)
.await
.map_err(|e| anyhow::anyhow!("dispatch chunk {}: {}", chunk_idx + 1, e))?;
let parsed: serde_json::Value =
serde_json::from_str(&raw).context("parse dispatch result")?;
let chunk_succeeded = parsed["summary"]["succeeded"].as_u64().unwrap_or(0) as usize;
let chunk_failed = parsed["summary"]["failed"].as_u64().unwrap_or(0) as usize;
total_succeeded += chunk_succeeded;
total_failed += chunk_failed;
let applied_now = applied_before + chunk.len();
eprintln!("applied {applied_now}/{total} (ok={total_succeeded}, failed={total_failed})");
}
let summary = serde_json::json!({
"total": total,
"succeeded": total_succeeded,
"failed": total_failed,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize summary")
);
Ok(())
}
pub async fn run_exec(args: ExecArgs) -> Result<()> {
if args.pending_events {
let summary =
pending_events::run_pending_events(args.db.as_deref(), &args.namespace, args.verbose)
.await?;
pending_events::print_summary(&summary);
return Ok(());
}
let mode = match (&args.ops, &args.ops_file) {
(Some(_), Some(_)) => {
anyhow::bail!(
"cannot use both a positional ops string and --ops-file; supply exactly one"
);
}
(None, None) => {
anyhow::bail!(
"no ops provided; supply a DSL expression as a positional argument or use \
--ops-file <PATH>"
);
}
(Some(ops), None) => ExecMode::Inline(ops.clone()),
(None, Some(path)) => ExecMode::OpsFile(path.clone()),
};
let mut cfg = RuntimeConfig::default();
if let Some(db_path) = resolve_db_override(args.db.as_deref()) {
cfg.db_path = db_path;
}
cfg.default_namespace =
Namespace::parse(&args.namespace).map_err(|e| anyhow::anyhow!("{e}"))?;
match mode {
ExecMode::Inline(ops) => {
run_exec_inline(
ops,
cfg,
args.presentation,
args.output_format,
args.save_file,
)
.await
}
ExecMode::OpsFile(path) => {
run_exec_ops_file(path, cfg, args.presentation, args.dry_run).await
}
}
}
enum ExecMode {
Inline(String),
OpsFile(PathBuf),
}
async fn run_exec_inline(
ops: String,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
) -> Result<()> {
#[cfg(unix)]
return run_exec_inline_with_forward(
ops,
cfg,
presentation,
output_format,
save_file,
forward_or_spawn_boxed,
)
.await;
#[cfg(not(unix))]
return run_exec_inline_with_forward(ops, cfg, presentation, output_format, save_file).await;
}
#[cfg_attr(not(unix), allow(unused_variables))]
async fn run_exec_inline_with_forward(
ops: String,
cfg: RuntimeConfig,
presentation: Option<String>,
output_format: Option<String>,
save_file: Option<String>,
#[cfg(unix)] forward_fn: ForwardFnPtr,
) -> Result<()> {
enforce_strict_actor_mode(cfg.actor_id.as_deref(), &cfg.packs)?;
#[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(),
config_id: compute_config_id(&cfg, None),
protocol_version: PROTOCOL_VERSION,
probe_only: false,
format: output_format.clone(),
format_per_op: None,
from_wire: false,
};
if let Some(res) = forward_fn(&frame).await {
let output = res.map_err(|e| anyhow::anyhow!("{}", e.message))?;
println!("{output}");
return Ok(());
}
}
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
let toml_default = KhiveConfig::load_with_home_fallback(None)
.ok()
.flatten()
.and_then(|c| c.runtime.default_output_format);
let env_fmt = apply_env_output_format(toml_default);
let server = KhiveMcpServer::new(rt)
.map_err(|e| anyhow::anyhow!("{e}"))?
.with_default_output_format(env_fmt);
let params = RequestParams {
ops,
presentation,
presentation_per_op: None,
save_to: save_file,
format: output_format,
format_per_op: None,
};
let output = server
.dispatch_request_local(params)
.await
.map_err(|e| anyhow::anyhow!("{e}"))?;
println!("{output}");
Ok(())
}
async fn run_exec_ops_file(
path: PathBuf,
cfg: RuntimeConfig,
presentation: Option<String>,
dry_run: bool,
) -> Result<()> {
let ops = parse_ops_file(&path)?;
if ops.is_empty() {
anyhow::bail!("ops-file is empty (no non-blank lines): {}", path.display());
}
if dry_run {
let mut per_verb: BTreeMap<String, usize> = BTreeMap::new();
for op in &ops {
*per_verb.entry(op.tool.clone()).or_insert(0) += 1;
}
let summary = serde_json::json!({
"dry_run": true,
"total": ops.len(),
"per_verb": per_verb,
});
println!(
"{}",
serde_json::to_string_pretty(&summary).expect("serialize dry-run summary")
);
return Ok(());
}
let rt = KhiveRuntime::new(cfg).map_err(|e| anyhow::anyhow!("{e}"))?;
enforce_strict_actor_mode(rt.config().actor_id.as_deref(), &rt.config().packs)?;
let server = KhiveMcpServer::new(rt).map_err(|e| anyhow::anyhow!("{e}"))?;
apply_ops_file(&server, ops, presentation).await
}
#[cfg(test)]
mod tests {
use super::*;
use clap::Parser;
use serial_test::serial;
use tempfile::NamedTempFile;
#[test]
#[serial]
fn khive_db_env_binds_to_db_arg() {
std::env::set_var("KHIVE_DB", "/tmp/kkernel-exec-env.db");
let args = ExecArgs::parse_from(["exec", "stats()"]);
std::env::remove_var("KHIVE_DB");
assert_eq!(args.db.as_deref(), Some("/tmp/kkernel-exec-env.db"));
}
#[test]
fn pending_events_flag_sets_mode() {
let args = ExecArgs::parse_from(["exec", "--pending-events"]);
assert!(args.pending_events);
assert!(args.ops.is_none());
}
#[test]
fn pending_events_conflicts_with_ops() {
let result = ExecArgs::try_parse_from(["exec", "--pending-events", "stats()"]);
assert!(
result.is_err(),
"--pending-events and positional ops must conflict"
);
}
#[test]
fn pending_events_conflicts_with_ops_file() {
let result =
ExecArgs::try_parse_from(["exec", "--pending-events", "--ops-file", "/tmp/x.jsonl"]);
assert!(
result.is_err(),
"--pending-events and --ops-file must conflict"
);
}
#[test]
fn ops_positional_is_optional() {
let args = ExecArgs::parse_from(["exec", "--ops-file", "/tmp/batch.jsonl"]);
assert!(args.ops.is_none());
assert_eq!(
args.ops_file.as_deref(),
Some(std::path::Path::new("/tmp/batch.jsonl"))
);
}
#[test]
fn ops_positional_works_without_pending_events() {
let args = ExecArgs::parse_from(["exec", "stats()"]);
assert_eq!(args.ops.as_deref(), Some("stats()"));
assert!(!args.pending_events);
}
#[test]
fn presentation_defaults_to_verbose_when_flag_omitted() {
let args = ExecArgs::parse_from(["exec", "stats()"]);
assert_eq!(args.presentation.as_deref(), Some("verbose"));
}
#[test]
fn presentation_agent_flag_still_selects_agent() {
let args = ExecArgs::parse_from(["exec", "stats()", "--presentation", "agent"]);
assert_eq!(args.presentation.as_deref(), Some("agent"));
}
#[test]
fn presentation_human_flag_still_selects_human() {
let args = ExecArgs::parse_from(["exec", "stats()", "--presentation", "human"]);
assert_eq!(args.presentation.as_deref(), Some("human"));
}
#[test]
fn dry_run_requires_ops_file() {
let result = ExecArgs::try_parse_from(["exec", "stats()", "--dry-run"]);
assert!(
result.is_err(),
"dry-run without --ops-file should be rejected by clap"
);
}
fn isolated_server(db_path: &str) -> KhiveMcpServer {
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(db_path)),
embedding_model: None,
additional_embedding_models: vec![],
..Default::default()
};
let rt = KhiveRuntime::new(cfg).expect("runtime on temp db");
KhiveMcpServer::new(rt).expect("server on temp db")
}
#[test]
fn parse_ops_file_skips_blank_lines() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
f.write_all(b"\n").unwrap(); f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
let ops = parse_ops_file(&f.path().to_path_buf()).unwrap();
assert_eq!(ops.len(), 2);
}
#[test]
fn parse_ops_file_reports_line_number_on_malformed() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"tool\":\"stats\",\"args\":{}}\n").unwrap();
f.write_all(b"not-json\n").unwrap(); let err = parse_ops_file(&f.path().to_path_buf()).unwrap_err();
let msg = format!("{err:#}");
assert!(
msg.contains("line 2"),
"error should name the bad line number, got: {msg}"
);
}
#[test]
fn parse_ops_file_missing_tool_field() {
use std::io::Write as _;
let mut f = NamedTempFile::new().unwrap();
f.write_all(b"{\"notool\":\"x\",\"args\":{}}\n").unwrap();
let err = parse_ops_file(&f.path().to_path_buf()).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("line 1"), "should report line number: {msg}");
}
#[tokio::test]
async fn ops_file_applies_ops_and_summary_matches() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let server = isolated_server(&db_path);
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
for name in ["Alpha", "Beta", "Gamma"] {
let line = format!(
"{{\"tool\":\"create\",\"args\":{{\"kind\":\"concept\",\"name\":\"{name}\"}}}}\n"
);
f.write_all(line.as_bytes()).unwrap();
}
let ops = parse_ops_file(&f.path().to_path_buf()).unwrap();
assert_eq!(ops.len(), 3);
apply_ops_file(&server, ops, None).await.unwrap();
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(
count, 3,
"all 3 entities should be present after apply\nraw: {resp}"
);
}
#[tokio::test]
async fn ops_file_dry_run_writes_nothing() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
for name in ["DryA", "DryB"] {
let line = format!(
"{{\"tool\":\"create\",\"args\":{{\"kind\":\"concept\",\"name\":\"{name}\"}}}}\n"
);
f.write_all(line.as_bytes()).unwrap();
}
let path = f.path().to_path_buf();
let cfg = RuntimeConfig {
db_path: Some(PathBuf::from(&db_path)),
..Default::default()
};
run_exec_ops_file(path.clone(), cfg.clone(), None, true)
.await
.unwrap();
let server = isolated_server(&db_path);
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(count, 0, "dry-run must not write any entities");
}
#[tokio::test]
#[serial]
async fn strict_mode_rejects_before_daemon_forward_when_comm_and_no_actor() {
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).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() {
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: Some("lambda:tenant-x".to_string()), ..RuntimeConfig::default()
};
let result = run_exec_inline("stats()".to_string(), cfg, None, None, None).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_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() {
let prev_strict = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
let prev_no_daemon = std::env::var("KHIVE_NO_DAEMON").ok();
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).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_ok(),
"run_exec_inline must NOT reject when strict mode is OFF (OSS default); got: {result:?}"
);
}
#[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(_frame: &DaemonRequestFrame) -> super::ForwardFuture<'_> {
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() {
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,
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() {
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");
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,
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"),
}
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)"
);
}
#[tokio::test]
async fn ops_file_malformed_line_aborts_before_writes() {
let db_file = NamedTempFile::new().expect("temp db");
let db_path = db_file.path().to_str().expect("utf8").to_string();
let mut f = NamedTempFile::new().unwrap();
use std::io::Write as _;
f.write_all(
b"{\"tool\":\"create\",\"args\":{\"kind\":\"concept\",\"name\":\"ShouldNotExist\"}}\n",
)
.unwrap();
f.write_all(b"INVALID JSON LINE\n").unwrap();
let path = f.path().to_path_buf();
let err = parse_ops_file(&path).unwrap_err();
let msg = format!("{err:#}");
assert!(
msg.contains("line 2"),
"should report line 2 as malformed: {msg}"
);
let server = isolated_server(&db_path);
let params = RequestParams {
ops: r#"list(kind="concept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
};
let raw = server.dispatch_request_local(params).await.unwrap();
let resp: serde_json::Value = serde_json::from_str(&raw).unwrap();
let count = resp["results"][0]["result"]
.as_array()
.map(|a| a.len())
.unwrap_or(0);
assert_eq!(
count, 0,
"nothing should be written when any line fails to parse"
);
}
}