use std::collections::BTreeSet;
use std::num::NonZeroUsize;
use std::sync::Arc;
use std::time::Duration;
use uuid::Uuid;
use khive_pack_kg::handlers::ValidatedSearchRequest;
use khive_runtime::Namespace as RuntimeNamespace;
use khive_runtime::{
BackendId, KhiveRuntime, NoteSearchHit, PackRegistry, SearchHit, SearchSource,
VerbRegistryBuilder,
};
use khive_score::DeterministicScore;
use khive_storage::types::Direction;
use khive_storage::EdgeRelation;
use khive_types::namespace::Namespace;
use super::dispatch::bounded_backend_cause_for_log;
use super::{BackendRegistry, LocatorCache, SubstrateCoordinator, SubstrateCoordinatorService};
fn memory_runtime() -> Arc<KhiveRuntime> {
Arc::new(KhiveRuntime::memory().expect("memory runtime"))
}
#[derive(Clone, Default)]
struct CoordinatorLogCapture(Arc<std::sync::Mutex<Vec<u8>>>);
impl std::io::Write for CoordinatorLogCapture {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
impl CoordinatorLogCapture {
fn contents(&self) -> String {
String::from_utf8(self.0.lock().unwrap().clone()).expect("captured logs are UTF-8")
}
}
struct MakeCoordinatorLogCapture(CoordinatorLogCapture);
impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for MakeCoordinatorLogCapture {
type Writer = CoordinatorLogCapture;
fn make_writer(&'a self) -> Self::Writer {
self.0.clone()
}
}
#[derive(Debug)]
struct DenyWithCauseGate {
cause: String,
}
impl khive_runtime::Gate for DenyWithCauseGate {
fn check(
&self,
_req: &khive_runtime::GateRequest,
) -> Result<khive_runtime::GateDecision, khive_runtime::GateError> {
Ok(khive_runtime::GateDecision::deny(self.cause.clone()))
}
}
fn memory_runtime_denied_with(cause: String) -> Arc<KhiveRuntime> {
Arc::new(
KhiveRuntime::new(khive_runtime::RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
gate: Arc::new(DenyWithCauseGate { cause }),
..khive_runtime::RuntimeConfig::no_embeddings()
})
.expect("memory runtime with denying gate"),
)
}
fn search_hit(entity_id: Uuid, source: SearchSource) -> SearchHit {
SearchHit {
entity_id,
score: DeterministicScore::from_f64(1.0),
source,
title: None,
snippet: None,
}
}
fn note_search_hit(note_id: Uuid, source: SearchSource) -> NoteSearchHit {
NoteSearchHit {
note_id,
score: DeterministicScore::from_f64(1.0),
source,
title: None,
snippet: None,
}
}
fn packs_registry(runtime: Arc<KhiveRuntime>, pack_names: &[&str]) -> khive_runtime::VerbRegistry {
let gate = runtime.config().gate.clone();
let default_ns = runtime.config().default_namespace.clone();
let actor_id = runtime.config().actor_id.clone();
let mut builder = VerbRegistryBuilder::new();
builder.with_gate(gate);
builder.with_default_namespace(default_ns.as_str());
builder.with_actor_id(actor_id);
let names: Vec<String> = pack_names.iter().map(|s| s.to_string()).collect();
PackRegistry::register_packs(&names, (*runtime).clone(), &mut builder)
.unwrap_or_else(|n| panic!("pack {n:?} declared in inventory but factory missing"));
let registry = builder.build().expect("build registry");
runtime.install_edge_rules(registry.all_edge_rules());
registry
}
fn validated_kg_search(params: serde_json::Value) -> ValidatedSearchRequest {
let registry = packs_registry(memory_runtime(), &["kg"]);
ValidatedSearchRequest::from_value(params, ®istry).expect("valid KG search request")
}
#[test]
fn validated_search_reconciles_compatible_granular_kind_fields() {
let entity = validated_kg_search(serde_json::json!({
"kind": "concept",
"query": "typed request",
"entity_kind": "concept",
"entity_type": "algorithm",
"source": "text",
}));
assert_eq!(entity.kind_filter(), Some("concept"));
assert_eq!(entity.entity_type(), Some("algorithm"));
assert_eq!(entity.source(), Some(SearchSource::Text));
let note = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "typed request",
"note_kind": "observation",
"include_superseded": true,
}));
assert_eq!(note.kind_filter(), Some("observation"));
assert!(note.include_superseded());
}
#[test]
fn validated_search_rejects_contradictory_compatibility_kind_fields() {
let registry = packs_registry(memory_runtime(), &["kg"]);
for params in [
serde_json::json!({
"kind": "concept",
"query": "typed request",
"entity_kind": "document",
}),
serde_json::json!({
"kind": "observation",
"query": "typed request",
"note_kind": "decision",
}),
] {
let error = ValidatedSearchRequest::from_value(params, ®istry)
.expect_err("contradictory compatibility kinds must reject");
assert!(
error.to_string().contains("contradicts"),
"validation error must identify the contradiction: {error}"
);
}
}
#[test]
fn single_coordinator_is_single_backend() {
let coord = SubstrateCoordinator::single(memory_runtime());
assert!(coord.is_single_backend());
assert_eq!(coord.backend_count(), 1);
assert_eq!(coord.backend_ids().len(), 1);
assert_eq!(coord.backend_ids()[0].as_str(), "main");
}
#[test]
fn registry_register_dedup() {
let mut reg = BackendRegistry::new();
let rt = memory_runtime();
assert!(reg.register(BackendId::new("main"), Arc::clone(&rt)));
assert!(!reg.register(BackendId::new("main"), Arc::clone(&rt)));
assert_eq!(reg.len(), 1);
}
#[test]
fn registry_primary_is_first_registered() {
let mut reg = BackendRegistry::new();
let rt1 = memory_runtime();
let rt2 = memory_runtime();
reg.register(BackendId::new("main"), rt1);
reg.register(BackendId::new("lore"), rt2);
assert_eq!(reg.primary().unwrap().id.as_str(), "main");
}
#[test]
fn multi_backend_coordinator_not_single() {
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), memory_runtime());
registry.register(BackendId::new("lore"), memory_runtime());
let coord = SubstrateCoordinator::new(registry);
assert!(!coord.is_single_backend());
assert_eq!(coord.backend_count(), 2);
}
#[test]
fn backend_id_display() {
let id = BackendId::new("archive");
assert_eq!(id.to_string(), "archive");
assert_eq!(id.as_str(), "archive");
}
#[test]
fn backend_id_main_constant() {
assert_eq!(BackendId::main().as_str(), BackendId::MAIN);
}
#[test]
fn locator_cache_miss_returns_none() {
let cache = LocatorCache::new();
let id = Uuid::new_v4();
assert!(cache.get(id).is_none());
}
#[test]
fn locator_cache_insert_then_get_returns_backend() {
let cache = LocatorCache::new();
let id = Uuid::new_v4();
cache.insert(id, BackendId::new("main"));
let result = cache.get(id);
assert!(result.is_some());
assert_eq!(result.unwrap().as_str(), "main");
}
#[test]
fn locator_cache_expired_entry_returns_none() {
let cache = LocatorCache::with_ttl(Duration::from_nanos(1));
let id = Uuid::new_v4();
cache.insert(id, BackendId::new("main"));
std::thread::sleep(Duration::from_micros(1));
assert!(cache.get(id).is_none());
}
#[test]
fn locator_cache_purge_removes_expired() {
let cache = LocatorCache::with_ttl(Duration::from_nanos(1));
for _ in 0..5 {
cache.insert(Uuid::new_v4(), BackendId::new("main"));
}
std::thread::sleep(Duration::from_micros(1));
cache.purge_expired();
assert_eq!(cache.len(), 0);
}
#[test]
fn locator_cache_insert_purges_expired_entry() {
let cache = LocatorCache::with_ttl(Duration::from_nanos(1));
let expired_id = Uuid::new_v4();
cache.insert(expired_id, BackendId::new("main"));
std::thread::sleep(Duration::from_micros(1));
let live_id = Uuid::new_v4();
cache.insert(live_id, BackendId::new("main"));
assert_eq!(cache.len(), 1);
assert!(cache.get(expired_id).is_none());
}
#[test]
fn locator_cache_evicts_least_recently_used_at_capacity() {
let cache =
LocatorCache::with_ttl_and_capacity(Duration::from_secs(60), NonZeroUsize::new(2).unwrap());
let first = Uuid::new_v4();
let second = Uuid::new_v4();
let third = Uuid::new_v4();
cache.insert(first, BackendId::new("main"));
cache.insert(second, BackendId::new("main"));
assert!(cache.get(first).is_some());
cache.insert(third, BackendId::new("main"));
assert_eq!(cache.len(), 2);
assert!(cache.get(second).is_none());
assert!(cache.get(first).is_some());
assert!(cache.get(third).is_some());
}
#[tokio::test]
async fn locator_cache_miss_then_hit() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let runtime = coord.primary_runtime().unwrap();
let token = runtime.authorize(ns.clone()).unwrap();
let entity = runtime
.create_entity(&token, "concept", None, "LoRA", None, None, vec![])
.await
.expect("create entity");
let first = coord.locate(entity.id, &ns).await;
assert!(
first.is_some(),
"locate should find the entity on first call"
);
assert_eq!(first.unwrap().as_str(), BackendId::MAIN);
assert_eq!(coord.locator_cache().len(), 1, "cache should be populated");
let second = coord.locate(entity.id, &ns).await;
assert!(second.is_some(), "second locate should hit cache");
}
#[tokio::test]
async fn locator_cache_returns_none_for_unknown_uuid() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let unknown = Uuid::new_v4();
let result = coord.locate(unknown, &ns).await;
assert!(result.is_none(), "unknown UUID should resolve to None");
}
#[tokio::test]
async fn fan_out_search_single_backend_returns_hits() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let runtime = coord.primary_runtime().unwrap();
let token = runtime.authorize(ns.clone()).unwrap();
runtime
.create_entity(
&token,
"concept",
None,
"FlashAttention",
Some("IO-aware exact attention"),
None,
vec![],
)
.await
.expect("create entity");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "FlashAttention",
"limit": 10,
}));
let (hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(!hits.is_empty(), "should find the entity");
assert_eq!(per_backend.len(), 1, "single backend report");
assert!(per_backend[0].error.is_none(), "no error");
}
#[tokio::test]
async fn fan_out_search_single_backend_applies_source_filter_before_limit() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let runtime = coord.primary_runtime().expect("single primary runtime");
let token = runtime.authorize(ns.clone()).expect("authorize local");
runtime
.create_entity(
&token,
"concept",
None,
"SingleBackendSourceEntity",
Some("single backend source filter entity probe"),
None,
vec![],
)
.await
.expect("create source-filter entity");
runtime
.create_note(
&token,
"observation",
Some("SingleBackendSourceNote"),
"single backend source filter note probe",
None,
None,
vec![],
)
.await
.expect("create source-filter note");
let entity_request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "SingleBackendSourceEntity",
"source": "vector",
"limit": 1,
}));
let (entity_hits, _, entity_backends) = coord.fan_out_search(&entity_request, &ns).await;
assert!(
entity_hits.is_empty(),
"single-backend entity results must apply source=vector before limit"
);
assert!(
entity_backends[0].error.is_none(),
"source filtering must not turn into a backend error"
);
let note_request = validated_kg_search(serde_json::json!({
"kind": "note",
"query": "SingleBackendSourceNote",
"source": "vector",
"limit": 1,
}));
let (_, note_hits, note_backends) = coord.fan_out_search(¬e_request, &ns).await;
assert!(
note_hits.is_empty(),
"single-backend note results must apply source=vector before limit"
);
assert!(
note_backends[0].error.is_none(),
"source filtering must not turn into a backend error"
);
}
#[tokio::test]
async fn fan_out_search_two_backends_merged() {
let mut registry = BackendRegistry::new();
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
rt_main
.create_entity(
&tok_main,
"concept",
None,
"LoRA",
Some("Low-rank adaptation"),
None,
vec![],
)
.await
.expect("create on main");
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"QLoRA",
Some("Quantised LoRA"),
None,
vec![],
)
.await
.expect("create on lore");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "LoRA",
"limit": 10,
}));
let (merged_hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "both backends in report");
assert!(
!merged_hits.is_empty(),
"merged results should not be empty"
);
}
#[tokio::test]
async fn fan_out_search_caps_merged_entity_hits_at_limit() {
let mut registry = BackendRegistry::new();
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
for (rt, prefix) in [(&rt_main, "Main"), (&rt_lore, "Lore")] {
let token = rt.authorize(ns.clone()).unwrap();
for i in 0..3 {
rt.create_entity(
&token,
"concept",
None,
&format!("{prefix}LimitProbe{i}"),
Some("shared limitprobe token"),
None,
vec![],
)
.await
.expect("create entity");
}
}
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "limitprobe",
"limit": 2,
}));
let (merged_hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "both backends in report");
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors"
);
assert!(
merged_hits.len() <= 2,
"merged entity hits must be capped at limit=2, got {}",
merged_hits.len()
);
}
#[tokio::test]
async fn fan_out_search_caps_merged_note_hits_at_limit() {
let mut registry = BackendRegistry::new();
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
for (rt, prefix) in [(&rt_main, "Main"), (&rt_lore, "Lore")] {
let token = rt.authorize(ns.clone()).unwrap();
for i in 0..3 {
rt.create_note(
&token,
"observation",
Some(&format!("{prefix}NoteLimitProbe{i}")),
"shared notelimitprobe token",
None,
None,
vec![],
)
.await
.expect("create note");
}
}
let request = validated_kg_search(serde_json::json!({
"kind": "note",
"query": "notelimitprobe",
"limit": 2,
}));
let (_entity_hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "both backends in report");
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors"
);
assert!(
note_hits.len() <= 2,
"merged note hits must be capped at limit=2, got {}",
note_hits.len()
);
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_hung_backend_times_out_sibling_still_returns() {
let mut registry = BackendRegistry::new();
let rt_main = memory_runtime();
let rt_hung = memory_runtime();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("hung"), Arc::clone(&rt_hung));
let coord = SubstrateCoordinator::new(registry).with_hanging_backend("hung");
let ns = Namespace::local();
let token = rt_main.authorize(ns.clone()).unwrap();
rt_main
.create_entity(
&token,
"concept",
None,
"TimeoutProbeHealthySibling",
Some("must still be returned despite the hung backend"),
None,
vec![],
)
.await
.expect("create entity on healthy backend");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "TimeoutProbeHealthySibling",
"limit": 10,
}));
let (hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "both backends must report");
let hung_report = per_backend
.iter()
.find(|r| r.backend_id.as_str() == "hung")
.expect("hung backend must have a report entry");
let err = hung_report
.error
.as_deref()
.expect("hung backend must carry an error");
assert!(
err.contains("timed out"),
"hung backend error must be timeout-specific, got: {err:?}"
);
let healthy_report = per_backend
.iter()
.find(|r| r.backend_id.as_str() == "main")
.expect("healthy backend must have a report entry");
assert!(
healthy_report.error.is_none(),
"healthy backend must not error"
);
assert!(
!hits.is_empty(),
"the healthy sibling's hit must still be present in the merged result \
despite the hung backend"
);
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_multiple_hung_backends_share_one_absolute_deadline() {
let mut registry = BackendRegistry::new();
for backend in ["hung-a", "hung-b", "hung-c"] {
registry.register(BackendId::new(backend), memory_runtime());
}
let coord =
SubstrateCoordinator::new(registry).with_hanging_backends(["hung-a", "hung-b", "hung-c"]);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "shared absolute timeout",
"limit": 10,
}));
let started = tokio::time::Instant::now();
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &Namespace::local()).await;
assert!(hits.is_empty());
assert!(note_hits.is_empty());
assert_eq!(per_backend.len(), 3);
assert!(per_backend.iter().all(|entry| entry
.error
.as_deref()
.is_some_and(|error| error.contains("timed out"))));
let elapsed = started.elapsed();
assert!(
elapsed >= Duration::from_secs(5) && elapsed < Duration::from_secs(6),
"three concurrently-started hung backends renewed the request budget; elapsed={elapsed:?}"
);
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_rejects_sibling_that_completed_during_interrupt_grace() {
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("a-hung"), memory_runtime());
registry.register(BackendId::new("b-late"), memory_runtime());
let coord = SubstrateCoordinator::new(registry)
.with_hanging_backend("a-hung")
.with_delayed_backend("b-late", Duration::from_millis(5_100));
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "late sibling",
"limit": 10,
}));
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &Namespace::local()).await;
assert!(hits.is_empty());
assert!(note_hits.is_empty());
assert_eq!(per_backend.len(), 2);
for backend in ["a-hung", "b-late"] {
let report = per_backend
.iter()
.find(|entry| entry.backend_id.as_str() == backend)
.expect("every backend is reported");
assert!(
report
.error
.as_deref()
.is_some_and(|error| error.contains("timed out")),
"{backend} completed outside the absolute deadline but was accepted: {:?}",
report.error
);
}
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_single_backend_hung_backend_times_out_entity_substrate() {
let mut registry = BackendRegistry::new();
let rt_hung = memory_runtime();
registry.register(BackendId::new("hung"), Arc::clone(&rt_hung));
let coord = SubstrateCoordinator::new(registry).with_hanging_backend("hung");
assert!(
coord.is_single_backend(),
"precondition: registry must hold exactly one backend so the \
single-entry early-return path (not the spawned fan-out) is taken"
);
let ns = Namespace::local();
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "SingleBackendTimeoutProbeEntity",
"limit": 10,
}));
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(hits.is_empty(), "no hits: the only backend timed out");
assert!(
note_hits.is_empty(),
"no note hits: the only backend timed out"
);
assert_eq!(per_backend.len(), 1, "the single backend must report");
let report = &per_backend[0];
assert_eq!(report.backend_id.as_str(), "hung");
let err = report
.error
.as_deref()
.expect("hung single backend must carry an error");
assert!(
err.contains("timed out"),
"single-backend timeout error must be timeout-specific, got: {err:?}"
);
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_timeout_masks_backend_credentials_in_coordinator_warning() {
let secret = format!("archive auth token sk_live_{}", "z".repeat(32));
let mut registry = BackendRegistry::new();
registry.register(BackendId::new(secret.clone()), memory_runtime());
let coord = SubstrateCoordinator::new(registry).with_hanging_backend(&secret);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "credential-safe timeout",
"limit": 10,
}));
let captured = CoordinatorLogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCoordinatorLogCapture(captured.clone()))
.with_ansi(false)
.without_time()
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let (_hits, _note_hits, per_backend) =
coord.fan_out_search(&request, &Namespace::local()).await;
drop(guard);
let logs = captured.contents();
assert_eq!(per_backend.len(), 1);
assert!(per_backend[0].error.is_some());
assert!(
!logs.contains("sk_live_"),
"coordinator WARN leaked backend credential: {logs}"
);
assert!(
logs.contains("***MASKED***"),
"coordinator WARN omitted masked backend identity: {logs}"
);
assert!(logs.contains("backend search task timed out"));
}
#[test]
fn coordinator_warning_cause_masker_is_bounded_and_fail_closed() {
let secret = format!("authorization token sk_live_{} denied", "q".repeat(32));
let masked = bounded_backend_cause_for_log(&secret);
assert!(masked.contains("***MASKED***"));
assert!(!masked.contains("sk_live_"));
let oversized = "x".repeat(5_000);
let bounded = bounded_backend_cause_for_log(&oversized);
assert_eq!(bounded.chars().count(), 1_025);
assert!(bounded.ends_with('…'));
assert_eq!(
bounded_backend_cause_for_log(" \t\n"),
"backend search failed without diagnostic detail"
);
}
#[tokio::test]
async fn fan_out_search_masks_real_authorization_cause_in_coordinator_warning() {
let secret = format!("authorization token sk_live_{} denied", "r".repeat(32));
let mut registry = BackendRegistry::new();
registry.register(
BackendId::new("archive"),
memory_runtime_denied_with(secret.clone()),
);
let coord = SubstrateCoordinator::new(registry);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "credential-safe authorization",
"limit": 10,
}));
let captured = CoordinatorLogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCoordinatorLogCapture(captured.clone()))
.with_ansi(false)
.without_time()
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let (_hits, _note_hits, per_backend) =
coord.fan_out_search(&request, &Namespace::local()).await;
drop(guard);
let logs = captured.contents();
assert_eq!(per_backend.len(), 1);
assert!(
per_backend[0]
.error
.as_deref()
.is_some_and(|error| error.contains("sk_live_")),
"internal result should retain the raw cause until the MCP sanitizer"
);
assert!(
!logs.contains("sk_live_"),
"coordinator WARN leaked authorization credential: {logs}"
);
assert!(logs.contains("***MASKED***"));
assert!(logs.contains("authorization denied for namespace"));
}
#[tokio::test(start_paused = true)]
async fn fan_out_search_single_backend_hung_backend_times_out_note_substrate() {
let mut registry = BackendRegistry::new();
let rt_hung = memory_runtime();
registry.register(BackendId::new("hung"), Arc::clone(&rt_hung));
let coord = SubstrateCoordinator::new(registry).with_hanging_backend("hung");
assert!(
coord.is_single_backend(),
"precondition: registry must hold exactly one backend so the \
single-entry early-return path (not the spawned fan-out) is taken"
);
let ns = Namespace::local();
let request = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "SingleBackendTimeoutProbeNote",
"limit": 10,
}));
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(
hits.is_empty(),
"no entity hits: the only backend timed out"
);
assert!(note_hits.is_empty(), "no hits: the only backend timed out");
assert_eq!(per_backend.len(), 1, "the single backend must report");
let report = &per_backend[0];
assert_eq!(report.backend_id.as_str(), "hung");
let err = report
.error
.as_deref()
.expect("hung single backend must carry an error");
assert!(
err.contains("timed out"),
"single-backend timeout error must be timeout-specific, got: {err:?}"
);
}
#[tokio::test]
async fn fan_out_search_with_visibility_single_backend_finds_extra_namespace_row() {
let coord = SubstrateCoordinator::single(memory_runtime());
let runtime = coord.primary_runtime().unwrap();
let tenant_ns = Namespace::parse("tenant-a").expect("valid namespace");
let token = runtime.authorize(tenant_ns.clone()).unwrap();
runtime
.create_entity(
&token,
"concept",
None,
"VisibilityProbeSingleBackend",
Some("only visible via extra_visible"),
None,
vec![],
)
.await
.expect("create entity in tenant-a");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "VisibilityProbeSingleBackend",
"limit": 10,
}));
let (widened_hits, _notes, per_backend) = coord
.fan_out_search_with_visibility(
&request,
&Namespace::local(),
std::slice::from_ref(&tenant_ns),
)
.await;
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors: {per_backend:?}"
);
assert!(
!widened_hits.is_empty(),
"widened visibility must find the tenant-a row on the single-backend path"
);
let (narrow_hits, _notes, _per_backend) = coord
.fan_out_search_with_visibility(&request, &Namespace::local(), &[])
.await;
assert!(
narrow_hits.is_empty(),
"primary-only visibility (no widening) must not see the tenant-a row"
);
}
#[tokio::test]
async fn fan_out_search_with_visibility_multi_backend_finds_extra_namespace_row() {
let mut registry = BackendRegistry::new();
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let tenant_ns = Namespace::parse("tenant-b").expect("valid namespace");
let token = rt_lore.authorize(tenant_ns.clone()).unwrap();
rt_lore
.create_entity(
&token,
"concept",
None,
"VisibilityProbeMultiBackend",
Some("only visible via extra_visible, on the second backend"),
None,
vec![],
)
.await
.expect("create entity in tenant-b on the lore backend");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "VisibilityProbeMultiBackend",
"limit": 10,
}));
let (widened_hits, _notes, per_backend) = coord
.fan_out_search_with_visibility(
&request,
&Namespace::local(),
std::slice::from_ref(&tenant_ns),
)
.await;
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors: {per_backend:?}"
);
assert!(
!widened_hits.is_empty(),
"widened visibility must find the tenant-b row on the spawned multi-backend path"
);
let (narrow_hits, _notes, _per_backend) = coord
.fan_out_search_with_visibility(&request, &Namespace::local(), &[])
.await;
assert!(
narrow_hits.is_empty(),
"primary-only visibility (no widening) must not see the tenant-b row"
);
}
#[tokio::test]
async fn fan_out_search_applies_source_filter_before_rrf_and_limit() {
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("alpha"), memory_runtime());
registry.register(BackendId::new("beta"), memory_runtime());
let vector_entity = Uuid::from_u128(1);
let cross_source_entity = Uuid::from_u128(2);
let text_entity = Uuid::from_u128(3);
let mut entity_overrides = std::collections::HashMap::new();
entity_overrides.insert(
"alpha".to_string(),
vec![
search_hit(vector_entity, SearchSource::Vector),
search_hit(cross_source_entity, SearchSource::Text),
search_hit(text_entity, SearchSource::Text),
],
);
entity_overrides.insert(
"beta".to_string(),
vec![
search_hit(vector_entity, SearchSource::Vector),
search_hit(cross_source_entity, SearchSource::Vector),
search_hit(Uuid::from_u128(4), SearchSource::Text),
],
);
let entity_coord =
SubstrateCoordinator::new(registry).with_entity_hits_override(entity_overrides);
let entity_request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "overridden",
"source": "text",
"limit": 1,
}));
let (entity_hits, _, _) = entity_coord
.fan_out_search(&entity_request, &Namespace::local())
.await;
assert_eq!(
entity_hits
.iter()
.map(|hit| hit.entity_id)
.collect::<Vec<_>>(),
vec![text_entity],
"filtering must drop the top vector hit and the cross-backend Both hit before limit"
);
let both_request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "overridden",
"source": "both",
"limit": 1,
}));
let (both_hits, _, _) = entity_coord
.fan_out_search(&both_request, &Namespace::local())
.await;
assert_eq!(
both_hits
.iter()
.map(|hit| hit.entity_id)
.collect::<Vec<_>>(),
vec![cross_source_entity],
"a text hit on one backend and vector hit on another has final source=both"
);
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("alpha"), memory_runtime());
registry.register(BackendId::new("beta"), memory_runtime());
let vector_note = Uuid::from_u128(5);
let cross_source_note = Uuid::from_u128(6);
let text_note = Uuid::from_u128(7);
let mut note_overrides = std::collections::HashMap::new();
note_overrides.insert(
"alpha".to_string(),
vec![
note_search_hit(vector_note, SearchSource::Vector),
note_search_hit(cross_source_note, SearchSource::Text),
note_search_hit(text_note, SearchSource::Text),
],
);
note_overrides.insert(
"beta".to_string(),
vec![
note_search_hit(vector_note, SearchSource::Vector),
note_search_hit(cross_source_note, SearchSource::Vector),
note_search_hit(Uuid::from_u128(8), SearchSource::Text),
],
);
let note_coord = SubstrateCoordinator::new(registry).with_note_hits_override(note_overrides);
let note_request = validated_kg_search(serde_json::json!({
"kind": "note",
"query": "overridden",
"source": "text",
"limit": 1,
}));
let (_, note_hits, _) = note_coord
.fan_out_search(¬e_request, &Namespace::local())
.await;
assert_eq!(
note_hits.iter().map(|hit| hit.note_id).collect::<Vec<_>>(),
vec![text_note],
"note filtering must preserve final-source semantics before limit"
);
}
#[tokio::test]
async fn fan_out_search_rrf_merge_uses_full_candidate_window_not_per_backend_limit() {
let mut registry = BackendRegistry::new();
let rt_a = memory_runtime();
let rt_b = memory_runtime();
registry.register(BackendId::new("alpha"), Arc::clone(&rt_a));
registry.register(BackendId::new("beta"), Arc::clone(&rt_b));
let y = Uuid::from_u128(1);
let x = Uuid::from_u128(2);
let w = Uuid::from_u128(3);
let z = Uuid::from_u128(4);
let v = Uuid::from_u128(5);
let alpha_list = vec![
search_hit(x, SearchSource::Text),
search_hit(w, SearchSource::Text),
search_hit(y, SearchSource::Text),
];
let beta_list = vec![
search_hit(z, SearchSource::Text),
search_hit(v, SearchSource::Text),
search_hit(y, SearchSource::Text),
];
let mut overrides = std::collections::HashMap::new();
overrides.insert("alpha".to_string(), alpha_list.clone());
overrides.insert("beta".to_string(), beta_list.clone());
let coord = SubstrateCoordinator::new(registry).with_entity_hits_override(overrides);
let ns = Namespace::local();
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "irrelevant, hits are overridden",
"limit": 2,
}));
let (hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors: {per_backend:?}"
);
let ids: Vec<Uuid> = hits.iter().map(|h| h.entity_id).collect();
assert_eq!(
ids,
vec![y, x],
"Y (rank-3 on both backends, fused score 2/63) must outrank the rank-1 \
singletons X and Z (score 1/61 each); tie between X and Z is broken by \
ascending entity_id, got {ids:?}"
);
let old_buggy_result = super::dispatch::rrf_merge_entity_hits(
vec![
alpha_list.into_iter().take(2).collect(),
beta_list.into_iter().take(2).collect(),
],
2,
);
let old_ids: Vec<Uuid> = old_buggy_result.iter().map(|h| h.entity_id).collect();
assert_eq!(
old_ids,
vec![x, z],
"sanity: the pre-fix per-backend-truncated merge should have produced \
[X, Z] (Y dropped before ever reaching RRF), got {old_ids:?}"
);
assert_ne!(
ids, old_ids,
"the fixed full-candidate-window result must differ from the pre-fix \
per-backend-truncated result — otherwise this test cannot distinguish them"
);
}
#[tokio::test]
async fn fan_out_search_rrf_merge_uses_full_candidate_window_not_per_backend_limit_notes() {
let mut registry = BackendRegistry::new();
let rt_a = memory_runtime();
let rt_b = memory_runtime();
registry.register(BackendId::new("alpha"), Arc::clone(&rt_a));
registry.register(BackendId::new("beta"), Arc::clone(&rt_b));
let y = Uuid::from_u128(1);
let x = Uuid::from_u128(2);
let w = Uuid::from_u128(3);
let z = Uuid::from_u128(4);
let v = Uuid::from_u128(5);
let alpha_list = vec![
note_search_hit(x, SearchSource::Text),
note_search_hit(w, SearchSource::Text),
note_search_hit(y, SearchSource::Text),
];
let beta_list = vec![
note_search_hit(z, SearchSource::Text),
note_search_hit(v, SearchSource::Text),
note_search_hit(y, SearchSource::Text),
];
let mut overrides = std::collections::HashMap::new();
overrides.insert("alpha".to_string(), alpha_list.clone());
overrides.insert("beta".to_string(), beta_list.clone());
let coord = SubstrateCoordinator::new(registry).with_note_hits_override(overrides);
let ns = Namespace::local();
let request = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "irrelevant, hits are overridden",
"limit": 2,
}));
let (_entity_hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"no backend errors: {per_backend:?}"
);
let ids: Vec<Uuid> = note_hits.iter().map(|h| h.note_id).collect();
assert_eq!(
ids,
vec![y, x],
"Y (rank-3 on both backends, fused score 2/63) must outrank the rank-1 \
singletons X and Z (score 1/61 each); tie between X and Z is broken by \
ascending note_id, got {ids:?}"
);
let old_buggy_result = super::dispatch::rrf_merge_note_hits(
vec![
alpha_list.into_iter().take(2).collect(),
beta_list.into_iter().take(2).collect(),
],
2,
);
let old_ids: Vec<Uuid> = old_buggy_result.iter().map(|h| h.note_id).collect();
assert_eq!(
old_ids,
vec![x, z],
"sanity: the pre-fix per-backend-truncated merge should have produced \
[X, Z] (Y dropped before ever reaching RRF), got {old_ids:?}"
);
assert_ne!(
ids, old_ids,
"the fixed full-candidate-window result must differ from the pre-fix \
per-backend-truncated result — otherwise this test cannot distinguish them"
);
}
#[test]
fn validated_search_rejects_unknown_fields() {
let registry = packs_registry(memory_runtime(), &["kg"]);
let error = ValidatedSearchRequest::from_value(
serde_json::json!({
"kind": "entity",
"query": "typed request",
"bogus_field": true,
}),
®istry,
)
.expect_err("unknown search field must reject");
assert!(
error.to_string().contains("unknown field"),
"error must name the unknown-field rejection: {error}"
);
}
#[test]
fn cross_backend_entity_merge_preserves_retrieval_leg_membership() {
let text_only = Uuid::new_v4();
let vector_only = Uuid::new_v4();
let both_only = Uuid::new_v4();
let text_and_vector = Uuid::new_v4();
let text_on_both = Uuid::new_v4();
let vector_on_both = Uuid::new_v4();
let both_and_vector = Uuid::new_v4();
let merged = super::dispatch::rrf_merge_entity_hits(
vec![
vec![
search_hit(text_only, SearchSource::Text),
search_hit(both_only, SearchSource::Both),
search_hit(text_and_vector, SearchSource::Text),
search_hit(text_on_both, SearchSource::Text),
search_hit(vector_on_both, SearchSource::Vector),
search_hit(both_and_vector, SearchSource::Both),
],
vec![
search_hit(vector_only, SearchSource::Vector),
search_hit(text_and_vector, SearchSource::Vector),
search_hit(text_on_both, SearchSource::Text),
search_hit(vector_on_both, SearchSource::Vector),
search_hit(both_and_vector, SearchSource::Vector),
],
],
10,
);
let sources: std::collections::HashMap<Uuid, SearchSource> = merged
.into_iter()
.map(|hit| (hit.entity_id, hit.source))
.collect();
assert_eq!(sources[&text_only], SearchSource::Text);
assert_eq!(sources[&vector_only], SearchSource::Vector);
assert_eq!(sources[&both_only], SearchSource::Both);
assert_eq!(sources[&text_and_vector], SearchSource::Both);
assert_eq!(sources[&text_on_both], SearchSource::Text);
assert_eq!(sources[&vector_on_both], SearchSource::Vector);
assert_eq!(sources[&both_and_vector], SearchSource::Both);
}
#[tokio::test]
async fn fan_out_search_empty_registry_returns_empty() {
let coord = SubstrateCoordinator::new(BackendRegistry::new());
let ns = Namespace::local();
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "anything",
"limit": 10,
}));
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(hits.is_empty());
assert!(note_hits.is_empty());
assert!(per_backend.is_empty());
}
#[tokio::test]
async fn fan_out_partial_failure_preserves_working_backend_hits() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"PartialFailureProbe",
Some("probe entity for partial-failure test"),
None,
vec![],
)
.await
.expect("create entity on lore");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry).with_failing_backend("main");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "PartialFailureProbe",
"limit": 10,
}));
let (merged_hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(
per_backend.len(),
2,
"both backends should appear in the report"
);
let main_result = per_backend
.iter()
.find(|r| r.backend_id.as_str() == "main")
.expect("main backend result must be present");
assert!(
main_result.error.is_some(),
"main backend should report an error"
);
assert!(
main_result.hits.is_empty(),
"main backend should have no hits"
);
let lore_result = per_backend
.iter()
.find(|r| r.backend_id.as_str() == "lore")
.expect("lore backend result must be present");
assert!(
lore_result.error.is_none(),
"lore backend should have no error"
);
assert!(
!merged_hits.is_empty(),
"merged hits must include results from the working backend"
);
}
#[tokio::test]
async fn fan_out_panicked_backend_is_explicit_in_per_backend() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"PanickedBackendProbe",
Some("healthy backend result retained when its sibling task panics"),
None,
vec![],
)
.await
.expect("create entity on healthy backend");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), rt_main);
registry.register(BackendId::new("lore"), rt_lore);
let coord = SubstrateCoordinator::new(registry).with_panicking_backend("main");
let request = validated_kg_search(serde_json::json!({
"kind": "concept",
"query": "PanickedBackendProbe",
"limit": 10,
}));
let captured = CoordinatorLogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCoordinatorLogCapture(captured.clone()))
.with_ansi(false)
.without_time()
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let (merged_hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
drop(guard);
let logs = captured.contents();
assert_eq!(per_backend.len(), 2, "every spawned backend is reported");
let panicked = per_backend
.iter()
.find(|result| result.backend_id.as_str() == "main")
.expect("panicked backend remains identified");
let error = panicked
.error
.as_deref()
.expect("panicked backend carries an explicit error");
assert!(
error.contains("join failed") && error.contains("panic"),
"join error should identify the task panic, got {error:?}"
);
assert!(logs.contains("backend search task failed"));
assert!(logs.contains("join failed") && logs.contains("panic"));
let healthy = per_backend
.iter()
.find(|result| result.backend_id.as_str() == "lore")
.expect("healthy backend is reported");
assert!(
healthy.error.is_none(),
"healthy backend must remain successful"
);
assert!(
!merged_hits.is_empty(),
"healthy backend hits survive a sibling task panic"
);
}
#[tokio::test]
async fn locate_finds_note_uuid() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let runtime = coord.primary_runtime().unwrap();
let token = runtime.authorize(ns.clone()).unwrap();
let note = runtime
.create_note(
&token,
"observation",
Some("locate-note-regression"),
"content for locate regression test",
None,
None,
vec![],
)
.await
.expect("create note");
let backend = coord.locate(note.id, &ns).await;
assert!(backend.is_some(), "locate should find the note's backend");
assert_eq!(backend.unwrap().as_str(), BackendId::MAIN);
assert_eq!(
coord.locator_cache().len(),
1,
"cache should be populated for the note"
);
}
#[test]
fn locator_cache_get_evicts_expired_entry() {
let cache = LocatorCache::with_ttl(Duration::from_nanos(1));
let id = Uuid::new_v4();
cache.insert(id, BackendId::new("main"));
assert_eq!(cache.len(), 1, "entry inserted");
std::thread::sleep(Duration::from_micros(1));
assert!(cache.get(id).is_none(), "expired entry returns None");
assert_eq!(cache.len(), 0, "expired entry must be evicted from the map");
}
#[test]
fn locator_cache_remove_evicts_live_entry() {
let cache = LocatorCache::new();
let id = Uuid::new_v4();
cache.insert(id, BackendId::new("main"));
assert!(cache.get(id).is_some(), "entry live before remove");
cache.remove(id);
assert!(cache.get(id).is_none(), "entry gone after remove");
assert_eq!(cache.len(), 0, "map must be empty after remove");
}
#[tokio::test]
async fn invalidate_clears_locate_cache() {
let coord = SubstrateCoordinator::single(memory_runtime());
let ns = Namespace::local();
let runtime = coord.primary_runtime().unwrap();
let token = runtime.authorize(ns.clone()).unwrap();
let entity = runtime
.create_entity(
&token,
"concept",
None,
"InvalidateTest",
None,
None,
vec![],
)
.await
.expect("create entity");
coord.locate(entity.id, &ns).await;
assert_eq!(coord.locator_cache().len(), 1, "cache populated");
coord.invalidate(entity.id);
assert_eq!(
coord.locator_cache().len(),
0,
"cache cleared after invalidate"
);
let found_again = coord.locate(entity.id, &ns).await;
assert!(found_again.is_some(), "locate re-finds after cache clear");
}
#[tokio::test]
async fn t1_single_backend_zero_change_invariant() {
let rt = memory_runtime();
let coord = SubstrateCoordinator::single(Arc::clone(&rt));
let ns = Namespace::local();
assert!(coord.is_single_backend(), "T1: must be single-backend");
let token = rt.authorize(ns.clone()).unwrap();
let entity = rt
.create_entity(&token, "concept", None, "T1Entity", None, None, vec![])
.await
.expect("T1: create entity");
let located = coord.locate(entity.id, &ns).await;
assert_eq!(
located.as_ref().map(|b| b.as_str()),
Some("main"),
"T1: single-backend locate must return main"
);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "T1Entity",
"limit": 10,
}));
let (hits, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert!(
!hits.is_empty(),
"T1: fan-out on single backend must return hits"
);
assert_eq!(per_backend.len(), 1, "T1: one backend in report");
assert!(
per_backend[0].error.is_none(),
"T1: no error on single backend"
);
}
#[tokio::test]
async fn t2_cross_backend_link_stamps_target_backend() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
let src = rt_main
.create_entity(
&tok_main,
"project",
None,
"SourceProject",
None,
None,
vec![],
)
.await
.expect("T2: create source on main");
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
let tgt = rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"TargetConcept",
None,
None,
vec![],
)
.await
.expect("T2: create target on lore");
let result = coord
.link_cross_backend(&ns, src.id, tgt.id, EdgeRelation::Implements, 1.0, None)
.await;
assert!(
result.is_ok(),
"T2: cross-backend link must succeed: {:?}",
result.err()
);
let edge = result.unwrap();
assert_eq!(
edge.target_backend.as_deref(),
Some("lore"),
"T2: edge must have target_backend stamped"
);
assert_eq!(edge.source_id, src.id, "T2: correct source_id");
assert_eq!(edge.target_id, tgt.id, "T2: correct target_id");
}
#[tokio::test]
async fn cross_backend_illegal_entity_pair_rejected_and_not_persisted() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), Arc::clone(&rt_main));
registry.register(BackendId::new("lore"), Arc::clone(&rt_lore));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
let src = rt_main
.create_entity(
&tok_main,
"concept",
None,
"SourceConcept",
None,
None,
vec![],
)
.await
.expect("T2b: create source on main");
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
let tgt = rt_lore
.create_entity(
&tok_lore,
"project",
None,
"TargetProject",
None,
None,
vec![],
)
.await
.expect("T2b: create target on lore");
let result = coord
.link_cross_backend(&ns, src.id, tgt.id, EdgeRelation::CompetesWith, 1.0, None)
.await;
let err = result.expect_err("T2b: illegal cross-backend link must be rejected");
assert!(
err.contains(
"currently legal relations for concept -> project under the loaded endpoint rules: none"
),
"T2b: rejection must expose the exact loaded legal set; got: {err}"
);
let main_neighbors = rt_main
.neighbors(&tok_main, src.id, Direction::Out, None, None)
.await
.expect("T2b: main neighbors query");
assert!(
main_neighbors.is_empty(),
"T2b: no edge must be written on the source backend after rejection"
);
let lore_neighbors = rt_lore
.neighbors(&tok_lore, tgt.id, Direction::In, None, None)
.await
.expect("T2b: lore neighbors query");
assert!(
lore_neighbors.is_empty(),
"T2b: no edge must be written on the target backend after rejection"
);
}
#[derive(Debug)]
struct ErrorAfterSetupCallsGate {
cause: String,
calls: std::sync::atomic::AtomicUsize,
}
impl khive_runtime::Gate for ErrorAfterSetupCallsGate {
fn check(
&self,
_req: &khive_runtime::GateRequest,
) -> Result<khive_runtime::GateDecision, khive_runtime::GateError> {
let call = self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
if call < 4 {
Ok(khive_runtime::GateDecision::allow())
} else {
Err(khive_runtime::GateError::Internal(self.cause.clone()))
}
}
}
fn memory_runtime_gate_erroring_after_first_call_with(cause: String) -> Arc<KhiveRuntime> {
Arc::new(
KhiveRuntime::new(khive_runtime::RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
gate: Arc::new(ErrorAfterSetupCallsGate {
cause,
calls: std::sync::atomic::AtomicUsize::new(0),
}),
..khive_runtime::RuntimeConfig::no_embeddings()
})
.expect("memory runtime with gate-erroring-after-setup-calls gate"),
)
}
#[tokio::test]
async fn t2c_cross_backend_link_authorize_gate_error_omits_backend_text_from_wire() {
const CANARY: &str = "postgres://svc:not-a-real-secret@internal-host";
let rt_main = memory_runtime_gate_erroring_after_first_call_with(CANARY.to_string());
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).expect("setup authorize");
let src = rt_main
.create_entity(
&tok_main,
"project",
None,
"SourceProject",
None,
None,
vec![],
)
.await
.expect("T2c: create source on main");
let tok_lore = rt_lore.authorize(ns.clone()).expect("setup authorize lore");
let tgt = rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"TargetConcept",
None,
None,
vec![],
)
.await
.expect("T2c: create target on lore");
let server = two_backend_server_with_packs(Arc::clone(&rt_main), Arc::clone(&rt_lore), &["kg"]);
let captured = CoordinatorLogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCoordinatorLogCapture(captured.clone()))
.with_ansi(false)
.without_time()
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let ops = format!(
r#"link(source_id="{}", target_id="{}", relation="implements")"#,
src.id, tgt.id
);
let raw = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops,
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T2c: dispatch must produce a response envelope");
drop(guard);
let logs = captured.contents();
let response: serde_json::Value =
serde_json::from_str(&raw).expect("T2c: response must be valid JSON");
let op = &response["results"][0];
assert_eq!(
op["ok"].as_bool(),
Some(false),
"T2c: link op must fail closed on the wire: {response}"
);
let wire_err = op["error"]
.as_str()
.unwrap_or_else(|| panic!("T2c: MCP-visible error must be a string: {response}"))
.to_string();
assert!(
!wire_err.contains(CANARY),
"T2c: MCP-visible link error must not embed backend error text: {wire_err:?}"
);
assert!(
!wire_err.contains("svc") && !wire_err.contains("internal-host"),
"T2c: MCP-visible link error must not embed backend error fragments: {wire_err:?}"
);
assert!(
wire_err.contains("gate backend unavailable"),
"T2c: MCP-visible link error must carry the stable classified reason: {wire_err:?}"
);
assert!(
!logs.contains(CANARY),
"T2c: backend error text must not reach the server-side log unmasked: {logs}"
);
assert!(
!logs.contains("not-a-real-secret"),
"T2c: the credential fragment must not reach the server-side log: {logs}"
);
assert!(
logs.contains("***MASKED***"),
"T2c: the gate failure log must still record the masked backend error: {logs}"
);
}
#[tokio::test]
async fn t2d_rego_gate_evaluator_failure_omits_canary_from_wire_and_logs() {
use khive_gate_rego::RegoGate;
const CANARY: &str = "AKIAFAKEKEY000000000";
let policy = r#"
package khive.gate
import rego.v1
decision := 1 + input.args.canary
"#;
let gate = Arc::new(RegoGate::from_policy_str(policy).expect("policy compiles"));
let rt = Arc::new(
KhiveRuntime::new(khive_runtime::RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
gate,
..khive_runtime::RuntimeConfig::no_embeddings()
})
.expect("T2d: memory runtime with rego gate"),
);
let server = khive_mcp::server::KhiveMcpServer::from_registry_with_meta(
packs_registry(Arc::clone(&rt), &["kg"]),
"local",
"test-rego-evaluator-failure",
);
let captured = CoordinatorLogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_writer(MakeCoordinatorLogCapture(captured.clone()))
.with_ansi(false)
.without_time()
.finish();
let guard = tracing::subscriber::set_default(subscriber);
let ops = format!(r#"list(kind="entity", canary="{CANARY}")"#);
let raw = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops,
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T2d: dispatch must produce a response envelope");
drop(guard);
let logs = captured.contents();
let response: serde_json::Value =
serde_json::from_str(&raw).expect("T2d: response must be valid JSON");
let op = &response["results"][0];
assert_eq!(
op["ok"].as_bool(),
Some(false),
"T2d: list op must fail closed on the wire: {response}"
);
let wire_err = op["error"]
.as_str()
.unwrap_or_else(|| panic!("T2d: MCP-visible error must be a string: {response}"))
.to_string();
assert!(
!wire_err.contains(CANARY),
"T2d: MCP-visible error must not embed the evaluator's raw error text: {wire_err:?}"
);
assert!(
wire_err.contains("policy evaluation failed"),
"T2d: MCP-visible error must carry the static classified reason: {wire_err:?}"
);
assert!(
wire_err.starts_with("permission denied for verb"),
"T2d: MCP-visible error must have the PermissionDenied shape, not GateUnavailable: {wire_err:?}"
);
assert!(
!logs.contains(CANARY),
"T2d: canary must not reach the server-side log: {logs}"
);
}
#[tokio::test]
async fn t3_fan_out_search_merged_from_two_backends() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("alpha"), Arc::clone(&rt_a));
registry.register(BackendId::new("beta"), Arc::clone(&rt_b));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_entity(
&tok_a,
"concept",
None,
"AlphaEntity",
Some("alpha side"),
None,
vec![],
)
.await
.expect("T3: create on alpha");
let tok_b = rt_b.authorize(ns.clone()).unwrap();
rt_b.create_entity(
&tok_b,
"concept",
None,
"BetaEntity",
Some("beta side"),
None,
vec![],
)
.await
.expect("T3: create on beta");
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "Entity",
"limit": 20,
}));
let (merged, _note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "T3: both backends in report");
assert!(
per_backend.iter().all(|r| r.error.is_none()),
"T3: no errors"
);
assert!(
merged.len() >= 2,
"T3: merged results must include hits from both backends, got {}",
merged.len()
);
}
#[tokio::test]
async fn t4_locate_namespace_agnostic() {
let rt = memory_runtime();
let coord = SubstrateCoordinator::single(Arc::clone(&rt));
let ns = Namespace::local();
let token = rt.authorize(ns.clone()).unwrap();
let entity = rt
.create_entity(&token, "concept", None, "T4NSAgnostic", None, None, vec![])
.await
.expect("T4: create entity");
let found = coord.locate(entity.id, &ns).await;
assert!(
found.is_some(),
"T4: locate must find the record with local namespace"
);
let other_ns = Namespace::parse("other").expect("T4: parse namespace");
let found_other = coord.locate(entity.id, &other_ns).await;
let _ = found_other; }
#[tokio::test]
async fn t5_record_created_prewarns_locator() {
let rt = memory_runtime();
let coord = SubstrateCoordinator::single(Arc::clone(&rt));
let ns = Namespace::local();
let token = rt.authorize(ns.clone()).unwrap();
let entity = rt
.create_entity(&token, "concept", None, "T5Prewarm", None, None, vec![])
.await
.expect("T5: create entity");
coord.record_created(entity.id, BackendId::main());
assert_eq!(
coord.locator_cache().len(),
1,
"T5: cache must be populated after record_created"
);
let backend = coord.locate(entity.id, &ns).await;
assert_eq!(
backend.as_ref().map(|b| b.as_str()),
Some("main"),
"T5: locate must return main from cache"
);
assert_eq!(coord.locator_cache().len(), 1, "T5: cache size stable");
}
#[tokio::test]
async fn fan_out_note_search_two_backends() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), Arc::clone(&rt_a));
registry.register(BackendId::new("lore"), Arc::clone(&rt_b));
let coord = SubstrateCoordinator::new(registry);
let ns = Namespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_note(
&tok_a,
"observation",
Some("AlphaObs"),
"alpha observation text",
None,
None,
vec![],
)
.await
.expect("create note on main");
let tok_b = rt_b.authorize(ns.clone()).unwrap();
rt_b.create_note(
&tok_b,
"observation",
Some("BetaObs"),
"beta observation text",
None,
None,
vec![],
)
.await
.expect("create note on lore");
let request = validated_kg_search(serde_json::json!({
"kind": "note",
"query": "observation",
"limit": 10,
}));
let (_entity_hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
assert_eq!(per_backend.len(), 2, "both backends in report");
assert!(per_backend.iter().all(|r| r.error.is_none()), "no errors");
assert!(
!note_hits.is_empty(),
"note fan-out must return hits, got 0"
);
}
#[tokio::test]
async fn fan_out_search_props_filter_drops_non_matching() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
rt_main
.create_entity(
&tok_main,
"concept",
None,
"PropsFanDecoy",
Some("propsfiltertest decoy entity without the matching property"),
None,
vec![],
)
.await
.expect("create decoy on main");
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
let target = rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"PropsFanTarget",
Some("propsfiltertest target entity with the matching property"),
Some(serde_json::json!({"status": "keep"})),
vec![],
)
.await
.expect("create target on lore");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), rt_main);
registry.register(BackendId::new("lore"), rt_lore);
let coord = SubstrateCoordinator::new(registry);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "propsfiltertest",
"limit": 10,
"properties": {"status": "keep"},
}));
let (hits, _note_hits, _per_backend) = coord.fan_out_search(&request, &ns).await;
let hit_ids: Vec<uuid::Uuid> = hits.iter().map(|h| h.entity_id).collect();
assert!(
hit_ids.contains(&target.id),
"entity with matching property must be in results; got {:?}",
hit_ids
);
assert!(
hit_ids.iter().all(|id| *id == target.id),
"only the matching entity should be returned; got {:?}",
hit_ids
);
}
#[tokio::test]
async fn fan_out_search_props_filter_before_truncation_semantics() {
let rt = memory_runtime();
let ns = Namespace::local();
let tok = rt.authorize(ns.clone()).unwrap();
rt.create_entity(
&tok,
"concept",
None,
"TruncSemAlpha",
Some("truncsemtest decoy entity without the filter property"),
None,
vec![],
)
.await
.expect("create decoy");
let target = rt
.create_entity(
&tok,
"concept",
None,
"TruncSemBeta",
Some("truncsemtest target entity with the filter property"),
Some(serde_json::json!({"keep": true})),
vec![],
)
.await
.expect("create target");
let coord = SubstrateCoordinator::single(rt);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "truncsemtest",
"limit": 1,
"properties": {"keep": true},
}));
let (hits, _note_hits, _per_backend) = coord.fan_out_search(&request, &ns).await;
let hit_ids: Vec<uuid::Uuid> = hits.iter().map(|h| h.entity_id).collect();
assert_eq!(
hits.len(),
1,
"exactly one hit expected with limit=1; got {:?}",
hit_ids
);
assert_eq!(
hits[0].entity_id, target.id,
"the matching entity must be returned even at limit=1; got {:?}",
hit_ids
);
}
#[tokio::test]
async fn fan_out_search_tags_filter_drops_non_matching() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
rt_main
.create_entity(
&tok_main,
"concept",
None,
"TagsFanDecoy",
Some("tagsfiltertest decoy entity without the target tag"),
None,
vec![],
)
.await
.expect("create untagged on main");
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
let tagged = rt_lore
.create_entity(
&tok_lore,
"concept",
None,
"TagsFanMarked",
Some("tagsfiltertest target entity with the target tag"),
None,
vec!["target-tag".to_string()],
)
.await
.expect("create tagged on lore");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), rt_main);
registry.register(BackendId::new("lore"), rt_lore);
let coord = SubstrateCoordinator::new(registry);
let request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "tagsfiltertest",
"limit": 10,
"tags": ["target-tag"],
}));
let (hits, _note_hits, _per_backend) = coord.fan_out_search(&request, &ns).await;
let hit_ids: Vec<uuid::Uuid> = hits.iter().map(|h| h.entity_id).collect();
assert!(
hit_ids.contains(&tagged.id),
"tagged entity must be in results; got {:?}",
hit_ids
);
assert!(
hit_ids.iter().all(|id| *id == tagged.id),
"only the tagged entity should be returned; got {:?}",
hit_ids
);
}
#[tokio::test]
async fn fan_out_search_preserves_full_entity_filter_contract() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
rt_main
.create_entity(
&tok_main,
"concept",
Some("algorithm"),
"FullEntityFilterWrongTag",
Some("fullentityfilter contract probe"),
Some(serde_json::json!({"scope": "keep"})),
vec!["other-tag".to_string()],
)
.await
.expect("create tag decoy");
rt_main
.create_entity(
&tok_main,
"concept",
Some("technique"),
"FullEntityFilterWrongType",
Some("fullentityfilter contract probe"),
Some(serde_json::json!({"scope": "keep"})),
vec!["entity-target".to_string()],
)
.await
.expect("create entity-type decoy");
rt_main
.create_entity(
&tok_main,
"document",
Some("paper"),
"FullEntityFilterWrongKind",
Some("fullentityfilter contract probe"),
Some(serde_json::json!({"scope": "keep"})),
vec!["entity-target".to_string()],
)
.await
.expect("create kind decoy");
rt_lore
.create_entity(
&tok_lore,
"concept",
Some("algorithm"),
"FullEntityFilterWrongProperties",
Some("fullentityfilter contract probe"),
Some(serde_json::json!({"scope": "drop"})),
vec!["entity-target".to_string()],
)
.await
.expect("create properties decoy");
let target = rt_lore
.create_entity(
&tok_lore,
"concept",
Some("algorithm"),
"FullEntityFilterTarget",
Some("fullentityfilter contract probe"),
Some(serde_json::json!({"scope": "keep", "extra": true})),
vec!["entity-target".to_string()],
)
.await
.expect("create matching entity");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), rt_main);
registry.register(BackendId::new("lore"), rt_lore);
let coord = SubstrateCoordinator::new(registry);
let request = validated_kg_search(serde_json::json!({
"kind": "concept",
"query": "fullentityfilter",
"limit": 20,
"entity_kind": "concept",
"entity_type": "algorithm",
"properties": {"scope": "keep"},
"tags": ["entity-target"],
}));
let (hits, note_hits, per_backend) = coord.fan_out_search(&request, &ns).await;
let hit_ids: Vec<Uuid> = hits.iter().map(|hit| hit.entity_id).collect();
assert!(
note_hits.is_empty(),
"entity request cannot return note hits"
);
assert!(
per_backend.iter().all(|result| result.error.is_none()),
"both backends should search successfully: {per_backend:?}"
);
assert_eq!(
hit_ids,
vec![target.id],
"all entity filters must be applied before merge"
);
}
#[tokio::test]
async fn fan_out_search_preserves_full_note_filter_contract() {
let rt_main = memory_runtime();
let rt_lore = memory_runtime();
let ns = Namespace::local();
let tok_main = rt_main.authorize(ns.clone()).unwrap();
let tok_lore = rt_lore.authorize(ns.clone()).unwrap();
rt_main
.create_note(
&tok_main,
"decision",
Some("wrong kind"),
"fullnotefilter contract probe",
Some(0.8),
Some(serde_json::json!({
"scope": "keep",
"tags": ["note-target"],
})),
vec![],
)
.await
.expect("create note-kind decoy");
let superseded = rt_lore
.create_note(
&tok_lore,
"observation",
Some("superseded target"),
"fullnotefilter contract probe",
Some(0.8),
Some(serde_json::json!({
"scope": "keep",
"tags": ["note-target"],
"extra": true,
})),
vec![],
)
.await
.expect("create superseded target note");
let replacement = rt_lore
.create_note(
&tok_lore,
"observation",
Some("replacement decoy"),
"fullnotefilter contract probe",
Some(0.8),
Some(serde_json::json!({
"scope": "drop",
"tags": ["other-tag"],
})),
vec![],
)
.await
.expect("create replacement note");
rt_lore
.link(
&tok_lore,
replacement.id,
superseded.id,
EdgeRelation::Supersedes,
1.0,
None,
)
.await
.expect("mark target note superseded");
let mut registry = BackendRegistry::new();
registry.register(BackendId::new("main"), rt_main);
registry.register(BackendId::new("lore"), rt_lore);
let coord = SubstrateCoordinator::new(registry);
let include_request = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "fullnotefilter",
"limit": 20,
"note_kind": "observation",
"include_superseded": true,
"properties": {"scope": "keep"},
"tags": ["note-target"],
}));
let (entity_hits, note_hits, per_backend) = coord.fan_out_search(&include_request, &ns).await;
let note_ids: Vec<Uuid> = note_hits.iter().map(|hit| hit.note_id).collect();
assert!(
entity_hits.is_empty(),
"note request cannot return entity hits"
);
assert!(
per_backend.iter().all(|result| result.error.is_none()),
"both backends should search successfully: {per_backend:?}"
);
assert_eq!(
note_ids,
vec![superseded.id],
"all note filters, including include_superseded, must reach fan-out"
);
let exclude_request = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "fullnotefilter",
"limit": 20,
"note_kind": "observation",
"include_superseded": false,
"properties": {"scope": "keep"},
"tags": ["note-target"],
}));
let (_entity_hits, note_hits, _per_backend) = coord.fan_out_search(&exclude_request, &ns).await;
assert!(
note_hits.is_empty(),
"the only matching note must be hidden when superseded notes are excluded"
);
}
fn two_backend_server(
rt_a: Arc<KhiveRuntime>,
rt_b: Arc<KhiveRuntime>,
) -> khive_mcp::server::KhiveMcpServer {
two_backend_server_with_packs(rt_a, rt_b, &["kg"])
}
fn two_backend_server_with_packs(
rt_a: Arc<KhiveRuntime>,
rt_b: Arc<KhiveRuntime>,
pack_names: &[&str],
) -> khive_mcp::server::KhiveMcpServer {
let registry = packs_registry(Arc::clone(&rt_a), pack_names);
let mut backend_reg = BackendRegistry::new();
backend_reg.register(BackendId::new("alpha"), Arc::clone(&rt_a));
backend_reg.register(BackendId::new("beta"), Arc::clone(&rt_b));
let coordinator = SubstrateCoordinatorService::new(SubstrateCoordinator::new(backend_reg));
khive_mcp::server::KhiveMcpServer::from_registry_with_meta(
registry,
"local",
"test-two-backend",
)
.with_coordinator(Arc::new(coordinator) as Arc<dyn khive_mcp::coordinator::CoordinatorService>)
}
#[tokio::test]
async fn t7a_multi_backend_search_populates_real_entity_kind() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let ns = RuntimeNamespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_entity(
&tok_a,
"concept",
None,
"T7aConceptAlpha",
Some("concept on alpha backend"),
None,
vec![],
)
.await
.expect("T7a: create concept on alpha");
let tok_b = rt_b.authorize(ns.clone()).unwrap();
rt_b.create_entity(
&tok_b,
"concept",
None,
"T7aConceptBeta",
Some("concept on beta backend"),
None,
vec![],
)
.await
.expect("T7a: create concept on beta");
let server = two_backend_server(Arc::clone(&rt_a), Arc::clone(&rt_b));
let result_str = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops: r#"search(kind="concept", query="T7aConcept")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T7a: dispatch");
let response: serde_json::Value =
serde_json::from_str(&result_str).expect("T7a: parse response JSON");
let results = response["results"].as_array().expect("T7a: results array");
assert!(
!results.is_empty(),
"T7a: should have at least one result op"
);
let op = &results[0];
assert!(
op["ok"].as_bool() == Some(true),
"T7a: search op must succeed, got: {op}"
);
let hits = op["result"].as_array().expect("T7a: result must be array");
assert!(!hits.is_empty(), "T7a: must find at least one concept hit");
for hit in hits {
let entity_kind = hit.get("entity_kind");
assert!(
entity_kind.is_some(),
"T7a: entity_kind field must be present in hit: {hit}"
);
assert!(
entity_kind.and_then(|v| v.as_str()).is_some(),
"T7a: entity_kind must be a non-null string, got: {hit}"
);
assert_eq!(
entity_kind.and_then(|v| v.as_str()),
Some("concept"),
"T7a: entity_kind must be 'concept', got: {hit}"
);
}
}
#[tokio::test]
async fn multi_backend_and_direct_search_rows_have_exact_key_set_parity() {
async fn first_hit_keys(
server: &khive_mcp::server::KhiveMcpServer,
ops: &str,
) -> BTreeSet<String> {
let raw = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops: ops.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("search dispatch must succeed");
let response: serde_json::Value =
serde_json::from_str(&raw).expect("search response must be valid JSON");
response["results"][0]["result"][0]
.as_object()
.unwrap_or_else(|| panic!("search must return an object row: {response}"))
.keys()
.cloned()
.collect()
}
let primary = memory_runtime();
let empty_secondary = memory_runtime();
let namespace = RuntimeNamespace::local();
let token = primary.authorize(namespace).expect("authorize primary");
primary
.create_entity(
&token,
"concept",
None,
"ExactShapeEntityProbe",
Some("entity used to compare direct and coordinator row keys"),
None,
vec![],
)
.await
.expect("create entity shape probe");
primary
.create_note(
&token,
"observation",
Some("Exact shape note probe"),
"exactshapenoteprobe content used to compare search row keys",
None,
None,
vec![],
)
.await
.expect("create note shape probe");
let direct = khive_mcp::server::KhiveMcpServer::from_registry_with_meta(
packs_registry(Arc::clone(&primary), &["kg"]),
"local",
"test-direct-shape",
);
let coordinated = two_backend_server(Arc::clone(&primary), empty_secondary);
for ops in [
r#"search(kind="concept", query="ExactShapeEntityProbe")"#,
r#"search(kind="observation", query="exactshapenoteprobe")"#,
] {
let direct_keys = first_hit_keys(&direct, ops).await;
let coordinated_keys = first_hit_keys(&coordinated, ops).await;
assert_eq!(
coordinated_keys, direct_keys,
"neither search route may add or omit a row field for {ops}"
);
}
}
#[tokio::test]
async fn t7b_multi_backend_search_kind_filter_excludes_off_kind() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let ns = RuntimeNamespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_entity(
&tok_a,
"concept",
None,
"T7bTargetConcept",
Some("the concept we want"),
None,
vec![],
)
.await
.expect("T7b: create concept on alpha");
rt_a.create_entity(
&tok_a,
"document",
None,
"T7bTargetDocument",
Some("a document that must be excluded"),
None,
vec![],
)
.await
.expect("T7b: create document on alpha");
let _ = rt_b.authorize(ns.clone()).unwrap();
let server = two_backend_server(Arc::clone(&rt_a), Arc::clone(&rt_b));
let result_str = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops: r#"search(kind="concept", query="T7bTarget")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T7b: dispatch");
let response: serde_json::Value = serde_json::from_str(&result_str).expect("T7b: parse");
let results = response["results"].as_array().expect("T7b: results array");
let op = &results[0];
assert!(
op["ok"].as_bool() == Some(true),
"T7b: search op must succeed"
);
let hits = op["result"].as_array().expect("T7b: result array");
for hit in hits {
let kind = hit["entity_kind"].as_str().unwrap_or("null");
assert_eq!(
kind, "concept",
"T7b: only concept hits expected, got entity_kind={kind:?} in: {hit}"
);
}
}
#[tokio::test]
async fn t7c_multi_backend_search_min_score_applied() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let ns = RuntimeNamespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_entity(
&tok_a,
"concept",
None,
"T7cMinScoreProbe",
Some("entity for min_score test"),
None,
vec![],
)
.await
.expect("T7c: create entity");
let _ = rt_b.authorize(ns.clone()).unwrap();
let server = two_backend_server(Arc::clone(&rt_a), Arc::clone(&rt_b));
let result_str = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops: r#"search(kind="concept", query="T7cMinScoreProbe", min_score=1.0)"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T7c: dispatch");
let response: serde_json::Value = serde_json::from_str(&result_str).expect("T7c: parse");
let results = response["results"].as_array().expect("T7c: results");
let op = &results[0];
assert!(
op["ok"].as_bool() == Some(true),
"T7c: search op must succeed"
);
let hits = op["result"].as_array().expect("T7c: result array");
assert!(
hits.is_empty(),
"T7c: min_score=1.0 must filter all hits, got {} hit(s)",
hits.len()
);
}
#[tokio::test]
async fn t7d_multi_backend_search_session_kind_routes_to_note_substrate() {
let rt_a = memory_runtime();
let rt_b = memory_runtime();
let ns = RuntimeNamespace::local();
let tok_a = rt_a.authorize(ns.clone()).unwrap();
rt_a.create_note(
&tok_a,
"session",
Some("Daily standup"),
"standup notes for the team",
None,
None,
vec![],
)
.await
.expect("T7d: create session note on alpha");
let _ = rt_b.authorize(ns.clone()).unwrap();
let server =
two_backend_server_with_packs(Arc::clone(&rt_a), Arc::clone(&rt_b), &["kg", "session"]);
let result_str = server
.dispatch_request_local(khive_mcp::tools::request::RequestParams {
ops: r#"search(kind="session", query="standup")"#.to_string(),
presentation: None,
presentation_per_op: None,
save_to: None,
format: None,
format_per_op: None,
request_id: None,
})
.await
.expect("T7d: dispatch");
let response: serde_json::Value = serde_json::from_str(&result_str).expect("T7d: parse");
let results = response["results"].as_array().expect("T7d: results");
let op = &results[0];
assert!(
op["ok"].as_bool() == Some(true),
"T7d: search op must succeed, got: {op}"
);
let hits = op["result"].as_array().expect("T7d: result array");
assert!(
!hits.is_empty(),
"T7d: session note must be found through the coordinator path"
);
for hit in hits {
assert_eq!(
hit.get("note_kind").and_then(|v| v.as_str()),
Some("session"),
"T7d: hit must be note-shaped with note_kind='session', got: {hit}"
);
assert!(
hit.get("entity_kind").map(|v| v.is_null()).unwrap_or(true),
"T7d: note-substrate hit must not carry an entity_kind, got: {hit}"
);
}
}
#[tokio::test]
async fn substrate_coordinator_service_hydrates_entity_and_note_metadata() {
use khive_mcp::coordinator::CoordinatorService;
let mut backend_reg = BackendRegistry::new();
let rt = memory_runtime();
backend_reg.register(BackendId::new("main"), Arc::clone(&rt));
let service = SubstrateCoordinatorService::new(SubstrateCoordinator::new(backend_reg));
let ns = Namespace::local();
let token = rt.authorize(ns.clone()).unwrap();
let entity = rt
.create_entity(
&token,
"concept",
None,
"Min1HydrationEntityProbe",
Some("entity for MIN-1 hydration coverage"),
None,
vec![],
)
.await
.expect("create entity");
let note = rt
.create_note(
&token,
"observation",
Some("Min1HydrationNoteProbe"),
"note content for min1hydrationnoteprobe coverage",
None,
None,
vec![],
)
.await
.expect("create note");
let entity_request = validated_kg_search(serde_json::json!({
"kind": "entity",
"query": "Min1HydrationEntityProbe",
"limit": 10,
}));
let entity_result = service.fan_out_search(&entity_request, &ns, &[]).await;
let entity_errors: Vec<&str> = entity_result
.per_backend
.iter()
.filter_map(|r| r.error.as_deref())
.collect();
assert!(
entity_errors.is_empty(),
"no backend errors: entity search: {entity_errors:?}"
);
assert!(
entity_result
.entity_hits
.iter()
.any(|h| h.entity_id == entity.id),
"the seeded entity must be found"
);
assert_eq!(
entity_result
.entity_kinds
.get(&entity.id)
.map(String::as_str),
Some("concept"),
"entity_kinds must be hydrated from the real stored entity, got: {:?}",
entity_result.entity_kinds
);
assert_eq!(
entity_result.entity_created_at.get(&entity.id),
Some(&entity.created_at),
"entity_created_at must be hydrated from the real stored entity, got: {:?}",
entity_result.entity_created_at
);
let note_request = validated_kg_search(serde_json::json!({
"kind": "observation",
"query": "min1hydrationnoteprobe",
"limit": 10,
}));
let note_result = service.fan_out_search(¬e_request, &ns, &[]).await;
let note_errors: Vec<&str> = note_result
.per_backend
.iter()
.filter_map(|r| r.error.as_deref())
.collect();
assert!(
note_errors.is_empty(),
"no backend errors: note search: {note_errors:?}"
);
assert!(
note_result.note_hits.iter().any(|h| h.note_id == note.id),
"the seeded note must be found"
);
assert_eq!(
note_result.note_kinds.get(¬e.id).map(String::as_str),
Some("observation"),
"note_kinds must be hydrated from the real stored note, got: {:?}",
note_result.note_kinds
);
assert_eq!(
note_result.note_created_at.get(¬e.id),
Some(¬e.created_at),
"note_created_at must be hydrated from the real stored note, got: {:?}",
note_result.note_created_at
);
assert_eq!(
note_result.note_names.get(¬e.id),
Some(¬e.name),
"note_names must be hydrated from the real stored note, got: {:?}",
note_result.note_names
);
}