use std::sync::Arc;
use std::time::Duration;
use uuid::Uuid;
use khive_runtime::Namespace as RuntimeNamespace;
use khive_runtime::{BackendId, KhiveRuntime, PackRegistry, VerbRegistryBuilder};
use khive_storage::EdgeRelation;
use khive_types::namespace::Namespace;
use super::{BackendRegistry, LocatorCache, SubstrateCoordinator, SubstrateCoordinatorService};
fn memory_runtime() -> Arc<KhiveRuntime> {
Arc::new(KhiveRuntime::memory().expect("memory runtime"))
}
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
}
#[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);
}
#[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 (hits, _note_hits, per_backend) = coord
.fan_out_search("FlashAttention", &ns, 10, false, None, None, &[])
.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_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 (merged_hits, _note_hits, per_backend) = coord
.fan_out_search("LoRA", &ns, 10, false, None, None, &[])
.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_empty_registry_returns_empty() {
let coord = SubstrateCoordinator::new(BackendRegistry::new());
let ns = Namespace::local();
let (hits, note_hits, per_backend) = coord
.fan_out_search("anything", &ns, 10, false, None, None, &[])
.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 (merged_hits, _note_hits, per_backend) = coord
.fan_out_search("PartialFailureProbe", &ns, 10, false, None, None, &[])
.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 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 (hits, _note_hits, per_backend) = coord
.fan_out_search("T1Entity", &ns, 10, false, None, None, &[])
.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 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 (merged, _note_hits, per_backend) = coord
.fan_out_search("Entity", &ns, 20, false, None, None, &[])
.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 (_entity_hits, note_hits, per_backend) = coord
.fan_out_search("observation", &ns, 10, true, None, None, &[])
.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 props = serde_json::json!({"status": "keep"});
let (hits, _note_hits, _per_backend) = coord
.fan_out_search("propsfiltertest", &ns, 10, false, None, Some(&props), &[])
.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 props = serde_json::json!({"keep": true});
let (hits, _note_hits, _per_backend) = coord
.fan_out_search("truncsemtest", &ns, 1, false, None, Some(&props), &[])
.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 (hits, _note_hits, _per_backend) = coord
.fan_out_search(
"tagsfiltertest",
&ns,
10,
false,
None,
None,
&["target-tag".to_string()],
)
.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
);
}
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 note_kinds: std::collections::HashSet<String> = registry
.all_note_kinds()
.into_iter()
.map(str::to_string)
.collect();
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), note_kinds);
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 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}"
);
}
}