use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use khive_runtime::{
config_from_env, run_migrations, runtime_config_from_khive_config, BackendConfig, BackendId,
BackendKind, ConnectionPool, KhiveConfig, KhiveRuntime, OutputFormat, PackRegistry,
RuntimeConfig, StorageBackend, VerbRegistryBuilder,
};
use crate::args::{resolve_cli_namespace, Args};
use crate::server::KhiveMcpServer;
use crate::transport::{ServeOptions, TransportRegistry};
pub struct MultiBackendRegistry {
pub registry: khive_runtime::VerbRegistry,
pub default_namespace: String,
pub config_id: String,
pub per_pack_runtimes: HashMap<String, Arc<KhiveRuntime>>,
pub main_backend: Arc<StorageBackend>,
}
pub async fn run(args: Args, registry: &TransportRegistry) -> anyhow::Result<()> {
let server = build_server(&args)?;
#[cfg(feature = "channel-email")]
{
use khive_channel::ChannelRegistry;
use khive_channel_email::EmailChannel;
use std::sync::Arc;
match EmailChannel::from_env() {
Ok(email_ch) => {
let email_ch = Arc::new(email_ch);
let mut ch_registry = ChannelRegistry::new();
let dyn_ch: Arc<dyn khive_channel::Channel> = email_ch.clone();
ch_registry.register(dyn_ch);
let ch_registry = Arc::new(ch_registry);
let verb_reg = server.verb_registry_clone();
let ingest_ns = ingest_namespace_from_env();
let default_actor = default_inbound_actor_from_env();
let mut allowlist = allowed_recipients_from_env();
if allowlist.is_empty() {
allowlist.push(email_ch.maintainer_address().to_string());
}
let mailbox = email_ch.mailbox().to_string();
let ingest_ns_clone = ingest_ns.clone();
let default_actor_clone = default_actor.clone();
let verb_reg_poll = verb_reg.clone();
let verb_reg_outbox = verb_reg.clone();
let ingest_ns_outbox = ingest_ns.clone();
let allowlist_clone = allowlist.clone();
let mailbox_clone = mailbox.clone();
let email_ch_clone = Arc::clone(&email_ch);
let spawned = run_if_authorized(&ingest_ns, &verb_reg, || {
tokio::task::spawn(channel_poll_loop(
ch_registry,
verb_reg_poll,
ingest_ns_clone,
default_actor_clone,
));
tokio::task::spawn(channel_outbox_loop(
email_ch_clone,
verb_reg_outbox,
ingest_ns_outbox,
mailbox_clone,
allowlist_clone,
));
tracing::info!("email channel polling and outbox loops started");
});
if !spawned {
tracing::error!(
namespace = %ingest_ns,
"email channel loops NOT started: ingest namespace authorization failed (fail-closed)"
);
}
}
Err(e) => {
tracing::warn!(
"channel-email feature is enabled but configuration is incomplete: {e}; \
email polling is disabled"
);
}
}
}
#[cfg(unix)]
if args.daemon {
khive_runtime::daemon::run_daemon(server).await?;
return Ok(());
}
#[cfg(not(unix))]
if args.daemon {
anyhow::bail!(
"--daemon mode requires Unix (macOS/Linux). On Windows, use the stdio transport."
);
}
let transport_name = args.transport.as_deref().unwrap_or("stdio");
let transport = registry.get(transport_name).ok_or_else(|| {
anyhow::anyhow!(
"unknown transport {transport_name:?}; registered: {}",
registry.names().join(", ")
)
})?;
let opts = ServeOptions {
bind: args.bind.clone(),
};
transport.serve(server, &opts).await
}
#[cfg(feature = "channel-email")]
fn ingest_namespace_from_env() -> String {
std::env::var("KHIVE_EMAIL_INGEST_NAMESPACE")
.ok()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| "local".to_string())
}
#[cfg(feature = "channel-email")]
fn default_inbound_actor_from_env() -> String {
std::env::var("KHIVE_EMAIL_DEFAULT_ACTOR")
.ok()
.filter(|s| !s.trim().is_empty())
.unwrap_or_else(|| "lambda:leo".to_string())
}
#[cfg(feature = "channel-email")]
fn allowed_recipients_from_env() -> Vec<String> {
std::env::var("KHIVE_EMAIL_SEND_ALLOWED_RECIPIENTS")
.ok()
.map(|s| {
s.split(',')
.map(|r| r.trim().to_string())
.filter(|r| !r.is_empty())
.collect()
})
.unwrap_or_default()
}
#[cfg(feature = "channel-email")]
fn run_if_authorized(
ns_str: &str,
registry: &khive_runtime::VerbRegistry,
on_authorized: impl FnOnce(),
) -> bool {
if preflight_ingest_namespace(ns_str, registry) {
on_authorized();
true
} else {
false
}
}
#[cfg(feature = "channel-email")]
fn preflight_ingest_namespace(ns_str: &str, registry: &khive_runtime::VerbRegistry) -> bool {
match khive_runtime::Namespace::parse(ns_str) {
Ok(ns) => match registry.authorize_namespace(ns) {
Ok(()) => true,
Err(e) => {
tracing::error!(
namespace = %ns_str,
error = %e,
"ingest namespace authorization denied; email polling will not start"
);
false
}
},
Err(e) => {
tracing::error!(
namespace = %ns_str,
error = %e,
"invalid ingest namespace string; email polling will not start"
);
false
}
}
}
#[cfg(feature = "channel-email")]
async fn channel_poll_loop(
channels: std::sync::Arc<khive_channel::ChannelRegistry>,
registry: khive_runtime::VerbRegistry,
ingest_namespace: String,
default_inbound_actor: String,
) {
use chrono::Utc;
use serde_json::json;
let mut last_poll = Utc::now();
loop {
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let since = last_poll;
last_poll = Utc::now();
for (kind, channel) in channels.iter() {
match channel.poll(since).await {
Ok(envelopes) => {
for env in envelopes {
let params = json!({
"namespace": ingest_namespace,
"from": env.from,
"to": env.to,
"content": env.content,
"subject": env.subject,
"channel_kind": kind,
"external_id": env.external_id,
"sent_at": env.sent_at.map(|ts| ts.to_rfc3339()),
"correlation_external_id": env.correlation_external_id,
"default_inbound_actor": default_inbound_actor,
});
if let Err(e) = registry.dispatch("comm.ingest", params).await {
tracing::warn!(
channel = kind,
"comm.ingest failed for inbound message: {e}"
);
}
}
}
Err(e) => {
tracing::warn!(channel = kind, "channel poll failed: {e}");
}
}
}
}
}
#[cfg(feature = "channel-email")]
fn note_already_delivered(props: &serde_json::Map<String, serde_json::Value>) -> bool {
props
.get("delivered_at")
.map(|v| !v.is_null())
.unwrap_or(false)
}
#[cfg(feature = "channel-email")]
async fn channel_outbox_loop(
email_channel: std::sync::Arc<khive_channel_email::EmailChannel>,
registry: khive_runtime::VerbRegistry,
ingest_namespace: String,
mailbox: String,
allowlist: Vec<String>,
) {
use chrono::Utc;
use khive_channel::{Channel, ChannelEnvelope};
use serde_json::json;
let domain = mailbox.split('@').nth(1).unwrap_or("localhost").to_string();
loop {
tokio::time::sleep(std::time::Duration::from_secs(5)).await;
let list_params = json!({
"namespace": ingest_namespace,
"kind": "message",
"direction": "outbound",
"delivered": false,
"limit": 200,
});
let list_result = match registry.dispatch("list", list_params).await {
Ok(r) => r,
Err(e) => {
tracing::warn!(error = %e, "outbox loop: list failed");
continue;
}
};
let notes = match list_result.as_array() {
Some(arr) => arr.clone(),
None => continue,
};
for note_val in notes {
let props = match note_val.get("properties") {
Some(serde_json::Value::Object(m)) => m.clone(),
_ => continue,
};
if props.get("direction").and_then(|v| v.as_str()) != Some("outbound") {
continue;
}
let to_actor = match props.get("to_actor").and_then(|v| v.as_str()) {
Some(a) if a.starts_with("email:") => a.to_string(),
_ => continue,
};
if note_already_delivered(&props) {
continue;
}
let note_id = match note_val.get("id").and_then(|v| v.as_str()) {
Some(id) => id.to_string(),
None => continue,
};
let recipient = to_actor
.strip_prefix("email:")
.unwrap_or(to_actor.as_str())
.to_string();
if !allowlist.is_empty() && !allowlist.contains(&recipient) {
tracing::warn!(
note_id = %note_id,
recipient = %recipient,
"outbox loop: recipient not in allowlist; skipping"
);
continue;
}
let subject = props
.get("subject")
.and_then(|v| v.as_str())
.unwrap_or("(no subject)")
.to_string();
let content = match note_val.get("content").and_then(|v| v.as_str()) {
Some(c) => c.to_string(),
None => continue,
};
let thread_id = props
.get("thread_id")
.and_then(|v| v.as_str())
.map(|s| s.to_string());
let message_id = match props.get("external_id").and_then(|v| v.as_str()) {
Some(eid) if !eid.is_empty() => eid.to_string(),
_ => {
let mid = format!("<{note_id}@{domain}>");
let claim_result = registry
.dispatch(
"update",
json!({
"namespace": ingest_namespace,
"id": note_id,
"properties": { "external_id": mid.clone() },
}),
)
.await;
if let Err(e) = claim_result {
tracing::warn!(
note_id = %note_id,
error = %e,
"outbox loop: failed to claim external_id; skipping"
);
continue;
}
mid
}
};
let mut env = ChannelEnvelope::new(
format!("email:{mailbox}"),
format!("email:{recipient}"),
content,
)
.with_subject(subject)
.with_message_id(message_id.clone());
if let Some(tid) = thread_id {
env = env.with_correlation(tid);
}
match email_channel.send(env).await {
Ok(()) => {
let delivered_at = Utc::now().to_rfc3339();
let mark_result = registry
.dispatch(
"update",
json!({
"namespace": ingest_namespace,
"id": note_id,
"properties": { "delivered_at": delivered_at },
}),
)
.await;
match mark_result {
Ok(_) => {
tracing::info!(
note_id = %note_id,
recipient = %recipient,
message_id = %message_id,
"outbox loop: delivered"
);
}
Err(e) => {
tracing::warn!(
note_id = %note_id,
error = %e,
"outbox loop: failed to set delivered_at (AT-LEAST-ONCE: will retry)"
);
}
}
}
Err(e) => {
tracing::warn!(
note_id = %note_id,
recipient = %recipient,
error = %e,
"outbox loop: send failed; will retry next cycle"
);
}
}
}
}
}
pub async fn serve_server(
server: KhiveMcpServer,
args: &Args,
registry: &TransportRegistry,
) -> anyhow::Result<()> {
#[cfg(unix)]
if args.daemon {
khive_runtime::daemon::run_daemon(server).await?;
return Ok(());
}
#[cfg(not(unix))]
if args.daemon {
anyhow::bail!(
"--daemon mode requires Unix (macOS/Linux). On Windows, use the stdio transport."
);
}
let transport_name = args.transport.as_deref().unwrap_or("stdio");
let transport = registry.get(transport_name).ok_or_else(|| {
anyhow::anyhow!(
"unknown transport {transport_name:?}; registered: {}",
registry.names().join(", ")
)
})?;
let opts = ServeOptions {
bind: args.bind.clone(),
};
transport.serve(server, &opts).await
}
pub fn build_registry_for_multi_backend(
base_config: RuntimeConfig,
khive_cfg: &KhiveConfig,
) -> anyhow::Result<MultiBackendRegistry> {
let mut backends: HashMap<String, Arc<StorageBackend>> = HashMap::new();
let mut path_to_backend: HashMap<std::path::PathBuf, Arc<StorageBackend>> = HashMap::new();
for backend_cfg in &khive_cfg.backends {
let canonical = canonical_backend_path(backend_cfg)?;
if let Some(ref canon) = canonical {
if let Some(existing) = path_to_backend.get(canon) {
backends.insert(backend_cfg.name.clone(), existing.clone());
continue;
}
}
let backend = open_backend(backend_cfg)?;
{
let mut writer = backend.pool().try_writer().map_err(|e| {
anyhow::anyhow!("backend {}: migration writer: {e}", backend_cfg.name)
})?;
run_migrations(writer.conn_mut())
.map_err(|e| anyhow::anyhow!("backend {}: migration: {e}", backend_cfg.name))?;
}
let arc = Arc::new(backend);
if let Some(canon) = canonical {
path_to_backend.insert(canon, arc.clone());
}
backends.insert(backend_cfg.name.clone(), arc);
}
let main_backend = backends
.get(BackendId::MAIN)
.ok_or_else(|| {
anyhow::anyhow!(
"[[backends]] is declared but no backend named \"main\" was found; \
add a [[backends]] entry with name = \"main\""
)
})?
.clone();
let pack_names = &base_config.packs;
let mut per_pack_runtimes_local: HashMap<String, KhiveRuntime> = HashMap::new();
for pack_name in pack_names {
let backend_name = khive_cfg
.packs
.get(pack_name.as_str())
.map(|pc| pc.backend.as_str())
.unwrap_or(BackendId::MAIN);
let backend = backends
.get(backend_name)
.cloned()
.unwrap_or_else(|| main_backend.clone());
let mut rt_config = base_config.clone();
rt_config.backend_id = BackendId::new(backend_name);
per_pack_runtimes_local.insert(
pack_name.clone(),
build_pack_runtime(backend, backend_name, rt_config, &main_backend),
);
}
let default_runtime = KhiveRuntime::from_backend(main_backend.clone(), {
let mut cfg = base_config.clone();
cfg.backend_id = BackendId::main();
cfg
});
#[cfg(feature = "bench-embedder")]
{
for rt in per_pack_runtimes_local.values() {
for name in rt.registered_embedding_model_names() {
rt.register_embedder(crate::bench_embedder::FeatureHashProvider::new(name));
}
}
for name in default_runtime.registered_embedding_model_names() {
default_runtime
.register_embedder(crate::bench_embedder::FeatureHashProvider::new(name));
}
}
enforce_strict_actor_mode(
default_runtime.config().actor_id.as_deref(),
&default_runtime.config().packs,
)?;
if should_warn_unattributed(
default_runtime.config().actor_id.as_deref(),
&default_runtime.config().packs,
) {
tracing::warn!(
"actor identity resolved to \"local\": comm sends will be stamped from \
\"local\" (unattributed) and comm.inbox will be unscoped (party-line). \
Set KHIVE_ACTOR or --actor to this lambda's id."
);
}
let gate = default_runtime.config().gate.clone();
let default_namespace = default_runtime.config().default_namespace.clone();
let config_id = crate::server::compute_config_id(default_runtime.config(), Some(khive_cfg));
let visible_namespaces = default_runtime.config().visible_namespaces.clone();
let mut builder = khive_runtime::VerbRegistryBuilder::new();
builder.with_gate(gate);
builder.with_default_namespace(default_namespace.as_str());
builder.with_visible_namespaces(visible_namespaces);
builder.with_actor_id(default_runtime.config().actor_id.clone());
if let Ok(tok) = default_runtime.authorize(khive_runtime::Namespace::local()) {
if let Ok(event_store) = default_runtime.events(&tok) {
builder.with_event_store(event_store);
}
}
khive_runtime::PackRegistry::register_packs_with_runtimes(
pack_names,
&per_pack_runtimes_local,
&default_runtime,
&mut builder,
)
.map_err(|e| anyhow::anyhow!("pack registration: {e}"))?;
let registry = builder
.build()
.map_err(|e| anyhow::anyhow!("registry build: {e}"))?;
default_runtime.install_edge_rules(registry.all_edge_rules());
for rt in per_pack_runtimes_local.values() {
rt.install_edge_rules(registry.all_edge_rules());
}
registry.call_register_embedders(&default_runtime);
registry.call_register_entity_type_validators(&default_runtime);
let backend_for_pack: HashMap<&str, &StorageBackend> = per_pack_runtimes_local
.iter()
.map(|(name, rt)| (name.as_str(), rt.backend()))
.collect();
let main_ref: &StorageBackend = main_backend.as_ref();
registry
.apply_schema_plans_with_map(&backend_for_pack, main_ref)
.map_err(|e| anyhow::anyhow!("pack schema boot failure: {e}"))?;
let per_pack_runtimes_arc: HashMap<String, Arc<KhiveRuntime>> = per_pack_runtimes_local
.into_iter()
.map(|(k, v)| (k, Arc::new(v)))
.collect();
Ok(MultiBackendRegistry {
registry,
default_namespace: default_namespace.as_str().to_string(),
config_id,
per_pack_runtimes: per_pack_runtimes_arc,
main_backend,
})
}
pub(crate) fn should_warn_unattributed(actor_id: Option<&str>, loaded_packs: &[String]) -> bool {
let is_local = actor_id.map(|id| id == "local").unwrap_or(true);
is_local && loaded_packs.iter().any(|p| p == "comm")
}
pub(crate) fn is_strict_actor_mode() -> bool {
std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR")
.map(|v| v.trim() == "1")
.unwrap_or(false)
}
pub fn enforce_strict_actor_mode(
actor_id: Option<&str>,
loaded_packs: &[String],
) -> anyhow::Result<()> {
if is_strict_actor_mode() && should_warn_unattributed(actor_id, loaded_packs) {
anyhow::bail!(
"KHIVE_REQUIRE_ATTRIBUTED_ACTOR=1 is set but no actor identity is \
configured. Set KHIVE_ACTOR or --actor to this lambda's id before \
starting in strict mode (comm pack requires an attributed actor to \
prevent party-line inbox exposure)."
);
}
Ok(())
}
pub fn build_server(args: &Args) -> anyhow::Result<KhiveMcpServer> {
let (cli_namespace_explicit, cli_namespace) =
resolve_cli_namespace(args).map_err(|e| anyhow::anyhow!("{e}"))?;
let config = resolve_runtime_config(RuntimeConfigInputs {
db: args.db.as_deref(),
config: args.config.as_deref(),
namespace: cli_namespace,
namespace_explicit: cli_namespace_explicit,
no_embed: args.no_embed,
packs: if args.pack.is_empty() {
None
} else {
Some(args.pack.clone())
},
brain_profile: args.brain_profile.clone(),
})?;
let khive_cfg = KhiveConfig::load_with_home_fallback(args.config.as_deref())
.map_err(|e| anyhow::anyhow!("config error: {e}"))?
.unwrap_or_default();
if khive_cfg.backends.is_empty() {
let runtime = KhiveRuntime::new(config)?;
#[cfg(feature = "bench-embedder")]
{
for name in runtime.registered_embedding_model_names() {
runtime.register_embedder(crate::bench_embedder::FeatureHashProvider::new(name));
}
}
enforce_strict_actor_mode(
runtime.config().actor_id.as_deref(),
&runtime.config().packs,
)?;
if should_warn_unattributed(
runtime.config().actor_id.as_deref(),
&runtime.config().packs,
) {
tracing::warn!(
"actor identity resolved to \"local\": comm sends will be stamped from \
\"local\" (unattributed) and comm.inbox will be unscoped (party-line). \
Set KHIVE_ACTOR or --actor to this lambda's id."
);
}
let fmt = apply_env_output_format(khive_cfg.runtime.default_output_format);
return KhiveMcpServer::new(runtime)
.map(|s| s.with_default_output_format(fmt))
.map_err(|e| anyhow::anyhow!("{e}"));
}
build_server_multi_backend(config, &khive_cfg)
}
fn canonical_backend_path(cfg: &BackendConfig) -> anyhow::Result<Option<PathBuf>> {
if cfg.kind == BackendKind::Memory {
return Ok(None);
}
let path = match cfg.path.as_ref() {
Some(p) => expand_tilde(p),
None => return Ok(None),
};
let parent = path
.parent()
.ok_or_else(|| anyhow::anyhow!("backend {}: path has no parent directory", cfg.name))?;
let file_name = path
.file_name()
.ok_or_else(|| anyhow::anyhow!("backend {}: path has no file name", cfg.name))?;
std::fs::create_dir_all(parent).map_err(|e| {
anyhow::anyhow!(
"backend {}: cannot create parent dir {}: {e}",
cfg.name,
parent.display()
)
})?;
let canon_parent = parent.canonicalize().map_err(|e| {
anyhow::anyhow!(
"backend {}: cannot canonicalize parent dir {}: {e}",
cfg.name,
parent.display()
)
})?;
Ok(Some(canon_parent.join(file_name)))
}
fn build_server_multi_backend(
base_config: RuntimeConfig,
khive_cfg: &KhiveConfig,
) -> anyhow::Result<KhiveMcpServer> {
let mut backends: HashMap<String, Arc<StorageBackend>> = HashMap::new();
let mut path_to_backend: HashMap<PathBuf, Arc<StorageBackend>> = HashMap::new();
for backend_cfg in &khive_cfg.backends {
let canonical = canonical_backend_path(backend_cfg)?;
if let Some(ref canon) = canonical {
if let Some(existing) = path_to_backend.get(canon) {
backends.insert(backend_cfg.name.clone(), existing.clone());
continue;
}
}
let backend = open_backend(backend_cfg)?;
{
let mut writer = backend.pool().try_writer().map_err(|e| {
anyhow::anyhow!("backend {}: migration writer: {e}", backend_cfg.name)
})?;
run_migrations(writer.conn_mut())
.map_err(|e| anyhow::anyhow!("backend {}: migration: {e}", backend_cfg.name))?;
}
let arc = Arc::new(backend);
if let Some(canon) = canonical {
path_to_backend.insert(canon, arc.clone());
}
backends.insert(backend_cfg.name.clone(), arc);
}
let main_backend = backends
.get(BackendId::MAIN)
.ok_or_else(|| {
anyhow::anyhow!(
"[[backends]] is declared but no backend named \"main\" was found; \
add a [[backends]] entry with name = \"main\""
)
})?
.clone();
let pack_names = &base_config.packs;
let mut per_pack_runtimes: HashMap<String, KhiveRuntime> = HashMap::new();
for pack_name in pack_names {
let backend_name = khive_cfg
.packs
.get(pack_name.as_str())
.map(|pc| pc.backend.as_str())
.unwrap_or(BackendId::MAIN);
let backend = backends
.get(backend_name)
.cloned()
.unwrap_or_else(|| main_backend.clone());
let mut rt_config = base_config.clone();
rt_config.backend_id = BackendId::new(backend_name);
per_pack_runtimes.insert(
pack_name.clone(),
build_pack_runtime(backend, backend_name, rt_config, &main_backend),
);
}
let default_runtime = KhiveRuntime::from_backend(main_backend.clone(), {
let mut cfg = base_config.clone();
cfg.backend_id = BackendId::main();
cfg
});
#[cfg(feature = "bench-embedder")]
{
for rt in per_pack_runtimes.values() {
for name in rt.registered_embedding_model_names() {
rt.register_embedder(crate::bench_embedder::FeatureHashProvider::new(name));
}
}
for name in default_runtime.registered_embedding_model_names() {
default_runtime
.register_embedder(crate::bench_embedder::FeatureHashProvider::new(name));
}
}
enforce_strict_actor_mode(
default_runtime.config().actor_id.as_deref(),
&default_runtime.config().packs,
)?;
if should_warn_unattributed(
default_runtime.config().actor_id.as_deref(),
&default_runtime.config().packs,
) {
tracing::warn!(
"actor identity resolved to \"local\": comm sends will be stamped from \
\"local\" (unattributed) and comm.inbox will be unscoped (party-line). \
Set KHIVE_ACTOR or --actor to this lambda's id."
);
}
let gate = default_runtime.config().gate.clone();
let default_namespace = default_runtime.config().default_namespace.clone();
let config_id = crate::server::compute_config_id(default_runtime.config(), Some(khive_cfg));
let visible_namespaces = default_runtime.config().visible_namespaces.clone();
let mut builder = VerbRegistryBuilder::new();
builder.with_gate(gate);
builder.with_default_namespace(default_namespace.as_str());
builder.with_visible_namespaces(visible_namespaces);
builder.with_actor_id(default_runtime.config().actor_id.clone());
if let Ok(tok) = default_runtime.authorize(khive_runtime::Namespace::local()) {
if let Ok(event_store) = default_runtime.events(&tok) {
builder.with_event_store(event_store);
}
}
PackRegistry::register_packs_with_runtimes(
pack_names,
&per_pack_runtimes,
&default_runtime,
&mut builder,
)
.map_err(|e| anyhow::anyhow!("pack registration: {e}"))?;
let registry = builder
.build()
.map_err(|e| anyhow::anyhow!("registry build: {e}"))?;
default_runtime.install_edge_rules(registry.all_edge_rules());
for rt in per_pack_runtimes.values() {
rt.install_edge_rules(registry.all_edge_rules());
}
registry.call_register_embedders(&default_runtime);
registry.call_register_entity_type_validators(&default_runtime);
let backend_for_pack: HashMap<&str, &StorageBackend> = per_pack_runtimes
.iter()
.map(|(name, rt)| (name.as_str(), rt.backend()))
.collect();
let main_ref: &StorageBackend = main_backend.as_ref();
registry
.apply_schema_plans_with_map(&backend_for_pack, main_ref)
.map_err(|e| anyhow::anyhow!("pack schema boot failure: {e}"))?;
let pool: Option<Arc<ConnectionPool>> = if main_backend.is_file_backed() {
Some(main_backend.pool_arc())
} else {
None
};
let fmt = apply_env_output_format(khive_cfg.runtime.default_output_format);
let server =
KhiveMcpServer::from_registry_with_meta(registry, default_namespace.as_str(), &config_id)
.with_default_output_format(fmt);
Ok(if let Some(p) = pool {
server.with_pool(p)
} else {
server
})
}
fn build_pack_runtime(
backend: Arc<StorageBackend>,
backend_name: &str,
rt_config: RuntimeConfig,
main_backend: &Arc<StorageBackend>,
) -> KhiveRuntime {
let rt = KhiveRuntime::from_backend(backend, rt_config);
if backend_name != BackendId::MAIN {
rt.with_core_backend(main_backend.clone())
} else {
rt
}
}
fn open_backend(cfg: &BackendConfig) -> anyhow::Result<StorageBackend> {
match cfg.kind {
BackendKind::Memory => StorageBackend::memory()
.map_err(|e| anyhow::anyhow!("backend {}: memory open: {e}", cfg.name)),
BackendKind::Sqlite => {
let path = cfg.path.as_ref().ok_or_else(|| {
anyhow::anyhow!(
"backend {}: sqlite backend requires a `path` field",
cfg.name
)
})?;
let expanded = expand_tilde(path);
if let Some(parent) = expanded.parent() {
std::fs::create_dir_all(parent).map_err(|e| {
anyhow::anyhow!(
"backend {}: cannot create parent dir {}: {e}",
cfg.name,
parent.display()
)
})?;
}
if cfg.read_only {
StorageBackend::sqlite_read_only(&expanded).map_err(|e| {
anyhow::anyhow!("backend {}: sqlite read-only open: {e}", cfg.name)
})
} else {
StorageBackend::sqlite(&expanded)
.map_err(|e| anyhow::anyhow!("backend {}: sqlite open: {e}", cfg.name))
}
}
}
}
fn expand_tilde(path: &std::path::Path) -> PathBuf {
let s = path.to_string_lossy();
if let Some(rest) = s.strip_prefix("~/") {
let home = std::env::var("HOME").unwrap_or_else(|_| ".".into());
PathBuf::from(format!("{home}/{rest}"))
} else if s == "~" {
let home = std::env::var("HOME").unwrap_or_else(|_| ".".into());
PathBuf::from(home)
} else {
path.to_path_buf()
}
}
pub struct RuntimeConfigInputs<'a> {
pub db: Option<&'a str>,
pub config: Option<&'a std::path::Path>,
pub namespace: khive_runtime::Namespace,
pub namespace_explicit: bool,
pub no_embed: bool,
pub packs: Option<Vec<String>>,
pub brain_profile: Option<String>,
}
pub fn resolve_runtime_config(inputs: RuntimeConfigInputs<'_>) -> anyhow::Result<RuntimeConfig> {
let db_path = match inputs.db {
Some(":memory:") => None,
Some(path) => Some(PathBuf::from(path)),
None => {
let home = std::env::var("HOME").unwrap_or_else(|_| ".".into());
Some(PathBuf::from(format!("{home}/.khive/khive.db")))
}
};
let packs = inputs
.packs
.unwrap_or_else(|| RuntimeConfig::default().packs);
let cli_brain_profile = inputs.brain_profile.filter(|s| !s.trim().is_empty());
let base_config = RuntimeConfig {
db_path,
default_namespace: inputs.namespace,
packs,
brain_profile: cli_brain_profile,
..RuntimeConfig::default()
};
let resolved = if inputs.no_embed {
let no_embed_base = RuntimeConfig {
embedding_model: None,
additional_embedding_models: vec![],
..base_config
};
resolve_actor_from_config(inputs.config, no_embed_base, inputs.namespace_explicit)?
} else {
resolve_config(inputs.config, base_config)?
};
let resolved = {
let mut resolved = resolved;
let ns = resolved.default_namespace.as_str().to_string();
if resolved.actor_id.is_none() && inputs.namespace_explicit && ns != "local" {
resolved.actor_id = Some(ns);
}
resolved
};
Ok(apply_env_brain_profile(resolved))
}
fn apply_env_brain_profile(mut cfg: RuntimeConfig) -> RuntimeConfig {
if cfg.brain_profile.is_none() {
cfg.brain_profile = std::env::var("KHIVE_BRAIN_PROFILE")
.ok()
.filter(|s| !s.trim().is_empty());
}
cfg
}
pub fn apply_env_output_format(toml_default: Option<OutputFormat>) -> OutputFormat {
if let Ok(val) = std::env::var("KHIVE_OUTPUT_FORMAT") {
match val.trim() {
"json" => return OutputFormat::Json,
"auto" => return OutputFormat::Auto,
"table" => return OutputFormat::Table,
_ => {
tracing::warn!(
value = %val,
"KHIVE_OUTPUT_FORMAT has unknown value; falling back to TOML / builtin default"
);
}
}
}
toml_default.unwrap_or(OutputFormat::Json)
}
fn resolve_config(
config_path: Option<&std::path::Path>,
base: RuntimeConfig,
) -> anyhow::Result<RuntimeConfig> {
match KhiveConfig::load_with_home_fallback(config_path)
.map_err(|e| anyhow::anyhow!("config error: {e}"))?
{
Some(khive_cfg) => {
let env_primary = std::env::var("KHIVE_EMBEDDING_MODEL").ok();
let env_additional = std::env::var("KHIVE_ADDITIONAL_EMBEDDING_MODELS").ok();
if !khive_cfg.engines.is_empty() && (env_primary.is_some() || env_additional.is_some())
{
tracing::warn!(
"khive config [[engines]] present; KHIVE_EMBEDDING_MODEL / \
KHIVE_ADDITIONAL_EMBEDDING_MODELS env vars are overridden"
);
}
Ok(runtime_config_from_khive_config(&khive_cfg, base))
}
None => {
let env_cfg = config_from_env();
if env_cfg.engines.is_empty() {
Ok(base)
} else {
Ok(runtime_config_from_khive_config(&env_cfg, base))
}
}
}
}
fn resolve_actor_from_config(
config_path: Option<&std::path::Path>,
base: RuntimeConfig,
cli_namespace_explicit: bool,
) -> anyhow::Result<RuntimeConfig> {
if cli_namespace_explicit {
return Ok(base);
}
match KhiveConfig::load_with_home_fallback(config_path)
.map_err(|e| anyhow::anyhow!("config error: {e}"))?
{
Some(khive_cfg) => {
let resolved = runtime_config_from_khive_config(&khive_cfg, base);
Ok(RuntimeConfig {
embedding_model: None,
additional_embedding_models: vec![],
..resolved
})
}
None => Ok(base),
}
}
#[cfg(test)]
mod tests {
use super::*;
use khive_runtime::Namespace;
use serial_test::serial;
use std::io::Write;
fn write_config(dir: &std::path::Path, body: &str) -> PathBuf {
let path = dir.join("khive.toml");
let mut f = std::fs::File::create(&path).expect("create config file");
f.write_all(body.as_bytes()).expect("write config");
path
}
#[test]
#[serial]
fn resolver_uses_config_file_engines_over_defaults() {
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
let default_cfg = RuntimeConfig::default();
let default_primary = format!("{:?}", default_cfg.embedding_model);
assert!(
!default_cfg.additional_embedding_models.is_empty(),
"precondition: default config has additional engines"
);
let dir = tempfile::tempdir().expect("temp dir");
let path = write_config(
dir.path(),
r#"
[[engines]]
name = "primary"
model = "bge-small-en-v1.5"
default = true
"#,
);
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config");
let resolved_primary = format!("{:?}", resolved.embedding_model);
assert_ne!(
resolved_primary, default_primary,
"resolved primary engine must come from the config file, not the default"
);
assert!(
resolved.embedding_model.is_some(),
"config-file engine must resolve to a primary embedding model"
);
assert!(
resolved.additional_embedding_models.is_empty(),
"config file declares one engine; additional list must be empty (not the default's)"
);
assert_eq!(resolved.db_path, None, ":memory: must map to in-memory db");
}
#[test]
#[serial]
fn resolver_falls_back_to_env_when_config_has_no_engines() {
std::env::remove_var("KHIVE_ADDITIONAL_EMBEDDING_MODELS");
std::env::set_var("KHIVE_EMBEDDING_MODEL", "bge-small-en-v1.5");
let dir = tempfile::tempdir().expect("temp dir");
let path = write_config(
dir.path(),
r#"
[runtime]
brain_profile = "unrelated"
"#,
);
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config");
std::env::remove_var("KHIVE_EMBEDDING_MODEL");
assert_eq!(
format!("{:?}", resolved.embedding_model),
"Some(BgeSmallEnV15)",
"KHIVE_EMBEDDING_MODEL must be applied as the fallback when the \
config file has no [[engines]] block, not treated as ignored"
);
}
#[test]
#[serial]
fn brain_profile_config_beats_env() {
std::env::set_var("KHIVE_BRAIN_PROFILE", "env-profile");
let dir = tempfile::tempdir().expect("temp dir");
let path = write_config(
dir.path(),
r#"
[runtime]
brain_profile = "project-profile"
"#,
);
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: None, })
.expect("resolve config");
std::env::remove_var("KHIVE_BRAIN_PROFILE");
assert_eq!(
resolved.brain_profile.as_deref(),
Some("project-profile"),
"project TOML brain_profile must win over KHIVE_BRAIN_PROFILE env var"
);
}
#[test]
#[serial]
fn brain_profile_env_fallback_when_no_toml() {
std::env::set_var("KHIVE_BRAIN_PROFILE", "env-profile");
let dir = tempfile::tempdir().expect("temp dir");
let path = write_config(
dir.path(),
r#"
[[engines]]
name = "primary"
model = "bge-small-en-v1.5"
default = true
"#,
);
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: None,
})
.expect("resolve config");
std::env::remove_var("KHIVE_BRAIN_PROFILE");
assert_eq!(
resolved.brain_profile.as_deref(),
Some("env-profile"),
"env var must be used when no CLI flag and no TOML brain_profile is set"
);
}
#[test]
#[serial]
fn brain_profile_cli_wins_over_all() {
std::env::set_var("KHIVE_BRAIN_PROFILE", "env-profile");
let dir = tempfile::tempdir().expect("temp dir");
let path = write_config(
dir.path(),
r#"
[runtime]
brain_profile = "project-profile"
"#,
);
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: Some(&path),
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: false,
no_embed: false,
packs: None,
brain_profile: Some("cli-profile".to_string()), })
.expect("resolve config");
std::env::remove_var("KHIVE_BRAIN_PROFILE");
assert_eq!(
resolved.brain_profile.as_deref(),
Some("cli-profile"),
"CLI --brain-profile must win over both TOML and KHIVE_BRAIN_PROFILE env var"
);
}
#[test]
#[serial]
fn cli_actor_flag_populates_actor_id() {
std::env::remove_var("KHIVE_ACTOR");
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: None,
namespace: Namespace::parse("lambda:agent-x").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config");
assert_eq!(
resolved.actor_id.as_deref(),
Some("lambda:agent-x"),
"--actor flag must populate actor_id (flag==env parity), not just default_namespace"
);
assert_eq!(
resolved.default_namespace.as_str(),
"lambda:agent-x",
"the flag still sets the write namespace"
);
}
#[test]
#[serial]
fn cli_actor_flag_local_stays_anonymous() {
std::env::remove_var("KHIVE_ACTOR");
let resolved = resolve_runtime_config(RuntimeConfigInputs {
db: Some(":memory:"),
config: None,
namespace: Namespace::parse("local").expect("ns"),
namespace_explicit: true,
no_embed: true,
packs: None,
brain_profile: None,
})
.expect("resolve config");
assert_eq!(
resolved.actor_id, None,
"explicit --actor local must remain anonymous (no actor_id) so the \
unattributed-comm warning still fires"
);
}
fn base_runtime_config_for_multi_backend() -> RuntimeConfig {
use khive_runtime::{AllowAllGate, BackendId, Namespace};
RuntimeConfig {
db_path: None,
gate: std::sync::Arc::new(AllowAllGate),
default_namespace: Namespace::parse("local").expect("ns"),
embedding_model: None,
additional_embedding_models: vec![],
packs: vec!["kg".to_string(), "comm".to_string()],
backend_id: BackendId::main(),
..RuntimeConfig::default()
}
}
#[tokio::test]
#[serial]
async fn multi_backend_boots_ok_with_two_memory_backends() {
use crate::tools::request::RequestParams;
use khive_runtime::PackConfig;
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "secondary".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let base_cfg = base_runtime_config_for_multi_backend();
let server = build_server_multi_backend(base_cfg, &khive_cfg)
.expect("multi-backend boot must succeed");
let kg_resp = server
.dispatch_request_local(RequestParams {
ops: r#"create(kind="concept", name="MultiBackendTestEntity")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.expect("kg dispatch must not error");
let kg_json: serde_json::Value =
serde_json::from_str(&kg_resp).expect("kg response is valid JSON");
let first_ok = kg_json["results"][0]["ok"].as_bool();
assert_eq!(
first_ok,
Some(true),
"kg create must succeed; response: {kg_resp}"
);
let comm_resp = server
.dispatch_request_local(RequestParams {
ops: r#"comm.send(to="local", content="multi-backend-test")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.expect("comm dispatch must not error");
let comm_json: serde_json::Value =
serde_json::from_str(&comm_resp).expect("comm response is valid JSON");
let first_comm_ok = comm_json["results"][0]["ok"].as_bool();
assert_eq!(
first_comm_ok,
Some(true),
"comm.send must succeed; response: {comm_resp}"
);
}
#[test]
fn secondary_pack_runtime_core_resolves_to_main_after_build_registry() {
use khive_runtime::PackConfig;
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "secondary".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let base_cfg = base_runtime_config_for_multi_backend();
let result = build_registry_for_multi_backend(base_cfg, &khive_cfg)
.expect("multi-backend registry must boot");
let comm_rt = result
.per_pack_runtimes
.get("comm")
.expect("comm pack runtime must be present in per_pack_runtimes");
assert_eq!(
comm_rt.backend_id().as_str(),
"secondary",
"comm pack runtime's own backend_id must be \"secondary\""
);
assert_eq!(
comm_rt.core().backend_id().as_str(),
BackendId::MAIN,
"secondary-backend pack must have core_backend wired to main (ADR-073); \
core().backend_id() returned {:?} — build_pack_runtime wiring missing",
comm_rt.core().backend_id().as_str()
);
}
#[tokio::test]
#[serial]
async fn multi_backend_preserves_actor_filtering() {
use crate::tools::request::RequestParams;
use khive_runtime::PackConfig;
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "secondary".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let base_cfg = RuntimeConfig {
actor_id: Some("actor-b".to_string()),
..base_runtime_config_for_multi_backend()
};
let server = build_server_multi_backend(base_cfg, &khive_cfg)
.expect("multi-backend boot must succeed");
let dispatch = |ops: String| {
let server = &server;
async move {
let resp = server
.dispatch_request_local(RequestParams {
ops,
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.expect("dispatch must not error");
serde_json::from_str::<serde_json::Value>(&resp).expect("valid JSON")
}
};
let to_a = dispatch(r#"comm.send(to="actor-a", content="for-a")"#.to_string()).await;
assert_eq!(to_a["results"][0]["ok"].as_bool(), Some(true), "{to_a}");
let to_b = dispatch(r#"comm.send(to="actor-b", content="for-b")"#.to_string()).await;
assert_eq!(to_b["results"][0]["ok"].as_bool(), Some(true), "{to_b}");
let inbox = dispatch(r#"comm.inbox()"#.to_string()).await;
let result = &inbox["results"][0]["result"];
let messages = result["messages"]
.as_array()
.expect("inbox returns a messages array");
let contents: Vec<&str> = messages
.iter()
.filter_map(|m| m["content"].as_str())
.collect();
assert!(
contents.contains(&"for-b"),
"actor-b must see the message addressed to it; got {contents:?}"
);
assert!(
!contents.contains(&"for-a"),
"actor-b must NOT see the message addressed to actor-a (leak #75 / B-BLOCKER-1); \
got {contents:?} — actor identity was not threaded into the multi-backend registry"
);
}
#[test]
fn multi_backend_missing_main_returns_error_mentioning_main() {
let khive_cfg = KhiveConfig {
backends: vec![BackendConfig {
name: "secondary".to_string(), kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
}],
packs: std::collections::HashMap::new(),
..KhiveConfig::default()
};
let base_cfg = base_runtime_config_for_multi_backend();
let result = build_server_multi_backend(base_cfg, &khive_cfg);
assert!(
result.is_err(),
"missing main backend must produce an error"
);
if let Err(err) = result {
assert!(
err.to_string().contains("main"),
"error message must mention \"main\"; got: {err}"
);
}
}
#[test]
fn read_only_backend_rejects_writes() {
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("ro_test.db");
let rw = StorageBackend::sqlite(&db_path).expect("rw backend");
rw.apply_pack_ddl_statements(&[
"CREATE TABLE IF NOT EXISTS ro_check (id INTEGER PRIMARY KEY)",
])
.expect("DDL on rw backend");
drop(rw);
let ro = StorageBackend::sqlite_read_only(&db_path).expect("ro backend");
let result = ro.apply_pack_ddl_statements(&["INSERT INTO ro_check (id) VALUES (1)"]);
assert!(
result.is_err(),
"write to a read-only backend must fail; got Ok(())"
);
}
#[test]
fn duplicate_sqlite_paths_deduplicated_to_single_backend() {
use khive_runtime::PackConfig;
let dir = tempfile::tempdir().unwrap();
let db_path = dir.path().join("shared.db");
let db_path_str = db_path.to_str().unwrap();
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Sqlite,
path: Some(db_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "alias".to_string(),
kind: BackendKind::Sqlite,
path: Some(db_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "alias".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let _ = db_path_str;
let base_cfg = base_runtime_config_for_multi_backend();
let result = build_server_multi_backend(base_cfg, &khive_cfg);
if let Err(ref e) = result {
panic!(
"two backends with the same canonical path must share one Arc and boot ok; got: {e}"
);
}
}
#[test]
fn config_id_folds_backend_topology_when_non_empty() {
use khive_runtime::{BackendId, KhiveConfig, Namespace, PackConfig, RuntimeConfig};
let base_rt = RuntimeConfig {
db_path: None,
default_namespace: Namespace::parse("local").unwrap(),
embedding_model: None,
packs: vec!["kg".to_string(), "comm".to_string()],
backend_id: BackendId::main(),
..RuntimeConfig::default()
};
let id_no_backends = crate::server::compute_config_id(&base_rt, None);
let id_empty_backends =
crate::server::compute_config_id(&base_rt, Some(&KhiveConfig::default()));
assert_eq!(
id_no_backends, id_empty_backends,
"empty-backends config_id must be byte-identical to None-config config_id"
);
let mut packs_a = std::collections::HashMap::new();
packs_a.insert(
"comm".to_string(),
PackConfig {
backend: "secondary".to_string(),
},
);
let cfg_a = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "secondary".to_string(),
kind: BackendKind::Memory,
path: None,
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: packs_a,
..KhiveConfig::default()
};
let cfg_b = KhiveConfig {
backends: cfg_a.backends.clone(),
packs: std::collections::HashMap::new(),
..KhiveConfig::default()
};
let id_a = crate::server::compute_config_id(&base_rt, Some(&cfg_a));
let id_b = crate::server::compute_config_id(&base_rt, Some(&cfg_b));
assert_ne!(
id_a, id_b,
"configs differing only in pack→backend routing must produce different config_ids; \
both produced: {id_a}"
);
}
#[tokio::test]
#[serial]
async fn multi_backend_isolates_pack_data_to_separate_files() {
use crate::tools::request::RequestParams;
use khive_runtime::PackConfig;
use rusqlite::Connection;
let dir = tempfile::tempdir().expect("temp dir");
let main_path = dir.path().join("main.db");
let second_path = dir.path().join("second.db");
let khive_cfg = KhiveConfig {
backends: vec![
BackendConfig {
name: "main".to_string(),
kind: BackendKind::Sqlite,
path: Some(main_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
BackendConfig {
name: "second".to_string(),
kind: BackendKind::Sqlite,
path: Some(second_path.clone()),
cache_mb: None,
journal_mode: None,
read_only: false,
},
],
packs: {
let mut m = std::collections::HashMap::new();
m.insert(
"comm".to_string(),
PackConfig {
backend: "second".to_string(),
},
);
m
},
..KhiveConfig::default()
};
let base_cfg = base_runtime_config_for_multi_backend();
let server = build_server_multi_backend(base_cfg, &khive_cfg)
.expect("multi-backend boot must succeed");
let dispatch = |ops: String| {
let server = &server;
async move {
server
.dispatch_request_local(RequestParams {
ops,
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
})
.await
.expect("dispatch must not error")
}
};
let kg_resp =
dispatch(r#"create(kind="concept", name="MainOnlyEntity")"#.to_string()).await;
let kg_json: serde_json::Value =
serde_json::from_str(&kg_resp).expect("kg response is valid JSON");
assert_eq!(
kg_json["results"][0]["ok"].as_bool(),
Some(true),
"kg create must succeed; response: {kg_resp}"
);
let comm_resp =
dispatch(r#"comm.send(to="local", content="SecondOnlyMsg")"#.to_string()).await;
let comm_json: serde_json::Value =
serde_json::from_str(&comm_resp).expect("comm response is valid JSON");
assert_eq!(
comm_json["results"][0]["ok"].as_bool(),
Some(true),
"comm.send must succeed; response: {comm_resp}"
);
drop(server);
let main_conn = Connection::open(&main_path).expect("open main.db");
let main_entity_count: i64 = main_conn
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'MainOnlyEntity' AND deleted_at IS NULL",
[],
|row| row.get(0),
)
.expect("query entities in main.db");
assert_eq!(
main_entity_count, 1,
"main.db MUST contain MainOnlyEntity (written via kg pack); got count={main_entity_count}"
);
let main_msg_count: i64 = main_conn
.query_row(
"SELECT COUNT(*) FROM notes WHERE kind = 'message'",
[],
|row| row.get(0),
)
.expect("query notes in main.db");
assert_eq!(
main_msg_count, 0,
"main.db MUST NOT contain any message notes (comm is pinned to second.db); \
got count={main_msg_count}"
);
let second_conn = Connection::open(&second_path).expect("open second.db");
let second_msg_count: i64 = second_conn
.query_row(
"SELECT COUNT(*) FROM notes WHERE kind = 'message' AND content = 'SecondOnlyMsg'",
[],
|row| row.get(0),
)
.expect("query notes in second.db");
assert_eq!(
second_msg_count, 2,
"second.db MUST contain SecondOnlyMsg (dual-write: 1 outbound + 1 inbound copy); \
got count={second_msg_count}"
);
let second_entity_count: i64 = second_conn
.query_row(
"SELECT COUNT(*) FROM entities WHERE name = 'MainOnlyEntity'",
[],
|row| row.get(0),
)
.expect("query entities in second.db");
assert_eq!(
second_entity_count, 0,
"second.db MUST NOT contain MainOnlyEntity (kg is pinned to main.db); \
got count={second_entity_count}"
);
}
#[cfg(feature = "channel-email")]
mod ingest_ns_tests {
use super::*;
#[test]
#[serial]
fn ingest_namespace_defaults_to_local() {
std::env::remove_var("KHIVE_EMAIL_INGEST_NAMESPACE");
assert_eq!(ingest_namespace_from_env(), "local");
}
#[test]
#[serial]
fn ingest_namespace_reads_env_var() {
std::env::set_var("KHIVE_EMAIL_INGEST_NAMESPACE", "lambda:mybot");
let ns = ingest_namespace_from_env();
std::env::remove_var("KHIVE_EMAIL_INGEST_NAMESPACE");
assert_eq!(ns, "lambda:mybot");
}
#[test]
#[serial]
fn ingest_namespace_ignores_blank_env_var() {
std::env::set_var("KHIVE_EMAIL_INGEST_NAMESPACE", " ");
let ns = ingest_namespace_from_env();
std::env::remove_var("KHIVE_EMAIL_INGEST_NAMESPACE");
assert_eq!(ns, "local", "blank env var must fall back to default");
}
#[test]
fn preflight_fails_on_invalid_namespace_string() {
let registry = khive_runtime::VerbRegistryBuilder::new()
.build()
.expect("build empty registry");
assert!(
!preflight_ingest_namespace("", ®istry),
"preflight must return false for an invalid namespace string"
);
}
#[test]
fn preflight_fails_when_gate_denies_namespace() {
use khive_runtime::{Gate, GateDecision, GateError, GateRequest};
use std::fmt;
#[derive(Debug)]
struct AlwaysDenyGate;
impl fmt::Display for AlwaysDenyGate {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "AlwaysDenyGate")
}
}
impl Gate for AlwaysDenyGate {
fn check(&self, _req: &GateRequest) -> Result<GateDecision, GateError> {
Ok(GateDecision::deny("test: always deny"))
}
}
let mut builder = khive_runtime::VerbRegistryBuilder::new();
builder.with_gate(std::sync::Arc::new(AlwaysDenyGate));
let registry = builder.build().expect("build registry with deny gate");
assert!(
!preflight_ingest_namespace("local", ®istry),
"preflight must return false when the gate denies the namespace"
);
}
#[test]
fn preflight_succeeds_with_allow_gate_and_valid_namespace() {
let registry = khive_runtime::VerbRegistryBuilder::new()
.build()
.expect("build registry with default allow-all gate");
assert!(
preflight_ingest_namespace("local", ®istry),
"preflight must return true for a valid namespace when the gate allows"
);
}
#[test]
fn spawn_not_called_when_gate_denies() {
use khive_runtime::{Gate, GateDecision, GateError, GateRequest};
use std::fmt;
#[derive(Debug)]
struct AlwaysDenyGate2;
impl fmt::Display for AlwaysDenyGate2 {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "AlwaysDenyGate2")
}
}
impl Gate for AlwaysDenyGate2 {
fn check(&self, _req: &GateRequest) -> Result<GateDecision, GateError> {
Ok(GateDecision::deny("spawn seam test: always deny"))
}
}
let mut builder = khive_runtime::VerbRegistryBuilder::new();
builder.with_gate(std::sync::Arc::new(AlwaysDenyGate2));
let registry = builder.build().expect("build registry with deny gate");
let mut spawn_count = 0usize;
let authorized = run_if_authorized("local", ®istry, || {
spawn_count += 1;
});
assert!(
!authorized,
"run_if_authorized must return false when gate denies"
);
assert_eq!(
spawn_count, 0,
"spawn must not be called when preflight fails"
);
}
#[test]
fn spawn_not_called_when_namespace_invalid() {
let registry = khive_runtime::VerbRegistryBuilder::new()
.build()
.expect("build empty registry");
let mut spawn_count = 0usize;
let authorized = run_if_authorized("", ®istry, || {
spawn_count += 1;
});
assert!(
!authorized,
"run_if_authorized must return false for invalid namespace"
);
assert_eq!(
spawn_count, 0,
"spawn must not be called when namespace is invalid"
);
}
}
fn packs(names: &[&str]) -> Vec<String> {
names.iter().map(|s| s.to_string()).collect()
}
#[test]
fn warn_when_actor_is_none_and_comm_loaded() {
assert!(should_warn_unattributed(None, &packs(&["kg", "comm"])));
}
#[test]
fn warn_when_actor_is_local_and_comm_loaded() {
assert!(should_warn_unattributed(
Some("local"),
&packs(&["kg", "comm"])
));
}
#[test]
fn no_warn_when_actor_is_configured() {
assert!(!should_warn_unattributed(
Some("lambda:khive"),
&packs(&["kg", "comm"])
));
}
#[test]
fn no_warn_when_comm_not_loaded() {
assert!(!should_warn_unattributed(Some("local"), &packs(&["kg"])));
}
#[test]
fn no_warn_when_actor_none_and_no_comm() {
assert!(!should_warn_unattributed(None, &packs(&["kg", "memory"])));
}
#[test]
#[serial]
fn strict_mode_off_by_default() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
assert!(
!is_strict_actor_mode(),
"strict mode must be OFF when KHIVE_REQUIRE_ATTRIBUTED_ACTOR is unset"
);
if let Some(v) = prev {
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v);
}
}
#[test]
#[serial]
fn strict_mode_on_when_env_var_is_1() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
assert!(
is_strict_actor_mode(),
"strict mode must be ON when KHIVE_REQUIRE_ATTRIBUTED_ACTOR=1"
);
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
}
#[test]
#[serial]
fn strict_mode_off_when_env_var_is_not_1() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "0");
assert!(
!is_strict_actor_mode(),
"strict mode must be OFF when KHIVE_REQUIRE_ATTRIBUTED_ACTOR=0"
);
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
}
#[test]
#[serial]
fn enforce_strict_actor_mode_returns_err_when_strict_and_no_actor() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
let result = enforce_strict_actor_mode(None, &packs(&["kg", "comm", "memory"]));
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
assert!(
result.is_err(),
"enforce_strict_actor_mode must return Err when strict mode is ON \
and no actor is configured (comm pack loaded)"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
"error message must name the env var; got: {msg}"
);
assert!(
msg.contains("KHIVE_ACTOR"),
"error message must name the remedy; got: {msg}"
);
}
#[test]
#[serial]
fn enforce_strict_actor_mode_ok_when_strict_and_actor_configured() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
let result = enforce_strict_actor_mode(Some("lambda:tenant-x"), &packs(&["kg", "comm"]));
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
assert!(
result.is_ok(),
"enforce_strict_actor_mode must return Ok when actor is properly configured"
);
}
#[test]
#[serial]
fn enforce_strict_actor_mode_ok_when_strict_off_and_no_actor() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR");
let result = enforce_strict_actor_mode(None, &packs(&["kg", "comm", "memory"]));
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
assert!(
result.is_ok(),
"enforce_strict_actor_mode must return Ok when strict mode is OFF \
(default OSS path must be completely unchanged)"
);
}
#[test]
#[serial]
fn enforce_strict_actor_mode_ok_when_strict_on_but_no_comm_pack() {
let prev = std::env::var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR").ok();
std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", "1");
let result = enforce_strict_actor_mode(None, &packs(&["kg", "memory"]));
match prev {
Some(v) => std::env::set_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR", v),
None => std::env::remove_var("KHIVE_REQUIRE_ATTRIBUTED_ACTOR"),
}
assert!(
result.is_ok(),
"enforce_strict_actor_mode must return Ok when comm pack is not loaded \
(no party-line risk even without actor)"
);
}
#[cfg(feature = "channel-email")]
mod outbox_delivered_guard_tests {
use super::*;
use serde_json::json;
#[test]
fn missing_delivered_at_is_undelivered() {
let props = json!({}).as_object().unwrap().clone();
assert!(!note_already_delivered(&props));
}
#[test]
fn explicit_null_delivered_at_is_undelivered() {
let props = json!({ "delivered_at": null }).as_object().unwrap().clone();
assert!(!note_already_delivered(&props));
}
#[test]
fn present_non_null_delivered_at_is_delivered() {
let props = json!({ "delivered_at": "2026-06-30T12:00:00Z" })
.as_object()
.unwrap()
.clone();
assert!(note_already_delivered(&props));
}
}
}