#[path = "../src/test_process.rs"]
mod test_process;
use khive_runtime::{KhiveRuntime, Namespace, RuntimeConfig};
use khive_storage::blob::ContentRef;
use khive_storage::types::{Direction, PageRequest, TraversalOptions, TraversalRequest};
use khive_storage::{BlobStore, EdgeRelation, Event, EventFilter, NewAttachment};
use khive_types::{EventKind, SubstrateKind};
use std::collections::BTreeSet;
use uuid::Uuid;
fn rt() -> KhiveRuntime {
KhiveRuntime::memory().expect("in-memory runtime")
}
#[tokio::test]
async fn note_supports_link_endpoints_are_observed_as_target() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let source = rt
.create_note(&tok, "observation", None, "source", None, None, vec![])
.await
.unwrap();
let target = rt
.create_note(&tok, "insight", None, "target", None, None, vec![])
.await
.unwrap();
rt.link(
&tok,
source.id,
target.id,
EdgeRelation::Supports,
1.0,
None,
)
.await
.unwrap();
let link_events = rt
.list_events(
&tok,
EventFilter {
kinds: vec![EventKind::LinkCreated],
..EventFilter::default()
},
PageRequest::default(),
)
.await
.unwrap();
assert_eq!(link_events.items.len(), 1);
let link_event = &link_events.items[0];
assert_eq!(link_event.payload["source_kind"], "note");
assert_eq!(link_event.payload["target_kind"], "note");
let query = format!(
"MATCH (ev)-[:observed_as_target]->(t) WHERE ev.id = '{}' RETURN t.id",
link_event.id
);
let rows = rt.query(&tok, &query).await.unwrap();
let observed_ids: BTreeSet<_> = rows
.iter()
.flat_map(|row| row.columns.iter())
.filter_map(|column| match &column.value {
khive_storage::types::SqlValue::Text(value) => Some(value.clone()),
_ => None,
})
.collect();
assert_eq!(rows.len(), 2);
assert_eq!(
observed_ids,
BTreeSet::from([source.id.to_string(), target.id.to_string()])
);
}
#[tokio::test]
async fn annotates_event_target_has_no_phantom_target_observation() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let note = rt
.create_note(&tok, "observation", None, "annotation", None, None, vec![])
.await
.unwrap();
let note_events = rt
.list_events(
&tok,
EventFilter {
kinds: vec![EventKind::NoteCreated],
..EventFilter::default()
},
PageRequest::default(),
)
.await
.unwrap();
assert_eq!(note_events.items.len(), 1);
let event_target = note_events.items[0].id;
rt.link(
&tok,
note.id,
event_target,
EdgeRelation::Annotates,
1.0,
None,
)
.await
.unwrap();
let link_events = rt
.list_events(
&tok,
EventFilter {
kinds: vec![EventKind::LinkCreated],
..EventFilter::default()
},
PageRequest::default(),
)
.await
.unwrap();
assert_eq!(link_events.items.len(), 1);
let link_event = &link_events.items[0];
assert_eq!(link_event.payload["source_kind"], "note");
assert_eq!(link_event.payload["target_kind"], "event");
let query = format!(
"MATCH (ev)-[:observed_as_target]->(t) WHERE ev.id = '{}' RETURN t.id",
link_event.id
);
let rows = rt.query(&tok, &query).await.unwrap();
assert_eq!(rows.len(), 1);
assert!(rows[0].columns.iter().any(|column| {
matches!(&column.value, khive_storage::types::SqlValue::Text(value) if value == ¬e.id.to_string())
}));
}
#[tokio::test]
async fn entity_create_and_get_roundtrip() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"LoRA",
Some("Low-Rank Adaptation"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let fetched = rt.get_entity(&tok, entity.id).await.unwrap();
assert_eq!(fetched.id, entity.id);
assert_eq!(fetched.name, "LoRA");
assert_eq!(fetched.kind, "concept");
assert_eq!(fetched.description.as_deref(), Some("Low-Rank Adaptation"));
}
#[tokio::test]
async fn entity_create_with_properties_and_tags() {
let rt = rt();
let research_tok = rt.authorize(Namespace::parse("research").unwrap()).unwrap();
let props = serde_json::json!({"domain": "fine-tuning", "type": "technique"});
let entity = rt
.create_entity_with_embedding_report(
&research_tok,
"concept",
None,
"QLoRA",
Some("Quantized LoRA"),
Some(props.clone()),
vec!["fine-tuning".to_string(), "quantization".to_string()],
)
.await
.map(|(record, _report)| record)
.unwrap();
let fetched = rt.get_entity(&research_tok, entity.id).await.unwrap();
assert_eq!(fetched.properties, Some(props));
assert_eq!(fetched.tags, vec!["fine-tuning", "quantization"]);
}
#[tokio::test]
async fn entity_create_with_content_attachment_roundtrip() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let temp = tempfile::tempdir().unwrap();
let blob_store = std::sync::Arc::new(
khive_db::stores::blob::FsBlobStore::new(temp.path().to_path_buf(), 0).unwrap(),
);
rt.install_blob_store(blob_store.clone())
.expect("install blob store");
let content_ref = blob_store.put(b"asset bytes".to_vec()).await.unwrap();
let entity = rt
.create_entity_with_attachments(
&tok,
"artifact",
Some("visual_asset"),
"asset",
None,
None,
vec![],
vec![NewAttachment {
role: "content".to_string(),
content_ref: content_ref.clone(),
media_type: None,
size_bytes: Some(11),
}],
)
.await
.unwrap();
let fetched = rt.get_entity(&tok, entity.id).await.unwrap();
assert_eq!(fetched.content_ref.as_deref(), Some(content_ref.as_str()));
}
#[tokio::test]
async fn entity_create_with_content_attachment_rejects_unpublished_blob() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let temp = tempfile::tempdir().unwrap();
rt.install_blob_store(std::sync::Arc::new(
khive_db::stores::blob::FsBlobStore::new(temp.path().to_path_buf(), 0).unwrap(),
))
.expect("install blob store");
let missing = ContentRef::from_digest_bytes(&[7; 32]);
let error = rt
.create_entity_with_attachments(
&tok,
"artifact",
Some("visual_asset"),
"asset",
None,
None,
vec![],
vec![NewAttachment {
role: "content".to_string(),
content_ref: missing,
media_type: None,
size_bytes: None,
}],
)
.await
.expect_err("unpublished ref must fail");
assert!(error.to_string().contains("requires a published blob"));
assert!(rt
.list_entities(&tok, None, None, 10, 0)
.await
.unwrap()
.is_empty());
}
#[tokio::test]
async fn entity_list_by_kind() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.create_entity_with_embedding_report(
&tok,
"concept",
None,
"FlashAttention",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
rt.create_entity_with_embedding_report(&tok, "concept", None, "GQA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.create_entity_with_embedding_report(
&tok,
"document",
None,
"Attention Is All You Need",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let concepts = rt
.list_entities(&tok, Some("concept"), None, 50, 0)
.await
.unwrap();
assert_eq!(concepts.len(), 2);
assert!(concepts.iter().any(|e| e.name == "FlashAttention"));
assert!(concepts.iter().any(|e| e.name == "GQA"));
let docs = rt
.list_entities(&tok, Some("document"), None, 50, 0)
.await
.unwrap();
assert_eq!(docs.len(), 1);
assert_eq!(docs[0].name, "Attention Is All You Need");
let all = rt.list_entities(&tok, None, None, 50, 0).await.unwrap();
assert_eq!(all.len(), 3);
}
#[tokio::test]
async fn entity_delete_soft() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(&tok, "concept", None, "to-delete", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let deleted = rt.delete_entity(&tok, entity.id, false).await.unwrap();
assert!(deleted);
let fetched = rt.get_entity(&tok, entity.id).await;
assert!(fetched.is_err());
}
#[tokio::test]
async fn entity_count_by_kind() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
for _ in 0..3 {
rt.create_entity_with_embedding_report(
&tok,
"concept",
None,
"concept-X",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
}
for _ in 0..2 {
rt.create_entity_with_embedding_report(&tok, "document", None, "doc-Y", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
}
let concept_count = rt.count_entities(&tok, Some("concept")).await.unwrap();
let doc_count = rt.count_entities(&tok, Some("document")).await.unwrap();
let total = rt.count_entities(&tok, None).await.unwrap();
assert_eq!(concept_count, 3);
assert_eq!(doc_count, 2);
assert_eq!(total, 5);
}
#[tokio::test]
async fn link_and_neighbors() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let lora = rt
.create_entity_with_embedding_report(&tok, "concept", None, "LoRA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let qlora = rt
.create_entity_with_embedding_report(&tok, "concept", None, "QLoRA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, qlora.id, lora.id, EdgeRelation::VariantOf, 1.0, None)
.await
.unwrap();
let hits = rt
.neighbors(&tok, qlora.id, Direction::Out, None, None)
.await
.unwrap();
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].node_id, lora.id);
assert_eq!(hits[0].relation, EdgeRelation::VariantOf);
}
#[tokio::test]
async fn traverse_multi_hop() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "A", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "B", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let c = rt
.create_entity_with_embedding_report(&tok, "concept", None, "C", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, a.id, b.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
rt.link(&tok, b.id, c.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
let request = TraversalRequest {
roots: vec![a.id],
options: TraversalOptions {
max_depth: 2,
direction: Direction::Out,
relations: Some(vec![EdgeRelation::Extends]),
..Default::default()
},
include_roots: false,
include_properties: false,
execution_budget: Default::default(),
};
let paths = rt.traverse(&tok, request).await.unwrap();
assert!(!paths.is_empty());
let reachable_ids: Vec<Uuid> = paths
.iter()
.flat_map(|p| p.nodes.iter().map(|n| n.node_id))
.collect();
assert!(reachable_ids.contains(&b.id));
assert!(reachable_ids.contains(&c.id));
}
#[tokio::test]
async fn traverse_total_weight_excludes_soft_deleted_nodes() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "A", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let light = rt
.create_entity_with_embedding_report(&tok, "concept", None, "light", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let heavy = rt
.create_entity_with_embedding_report(&tok, "concept", None, "heavy", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, a.id, light.id, EdgeRelation::Extends, 0.25, None)
.await
.unwrap();
rt.link(&tok, a.id, heavy.id, EdgeRelation::Extends, 0.9, None)
.await
.unwrap();
let request = || TraversalRequest {
roots: vec![a.id],
options: TraversalOptions {
max_depth: 1,
direction: Direction::Out,
relations: Some(vec![EdgeRelation::Extends]),
..Default::default()
},
include_roots: false,
include_properties: false,
execution_budget: Default::default(),
};
let before = rt.traverse(&tok, request()).await.unwrap();
assert_eq!(before.len(), 1);
assert_eq!(
before[0].total_weight, 0.9,
"baseline: the heavy neighbour sets total_weight while it is visible"
);
rt.delete_entity(&tok, heavy.id, false).await.unwrap();
let after = rt.traverse(&tok, request()).await.unwrap();
assert_eq!(after.len(), 1);
let ids: Vec<Uuid> = after[0].nodes.iter().map(|n| n.node_id).collect();
assert!(
!ids.contains(&heavy.id),
"soft-deleted node must be screened"
);
assert_eq!(
after[0].total_weight, 0.25,
"total_weight must fall to the surviving neighbour's weight, not stay \
at the soft-deleted node's 0.9"
);
}
#[tokio::test]
async fn create_note_and_list_notes() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.create_note(
&tok,
"observation",
None,
"LoRA is a fine-tuning technique",
Some(0.9),
None,
vec![],
)
.await
.unwrap();
rt.create_note(
&tok,
"observation",
None,
"QLoRA uses quantization",
Some(0.8),
None,
vec![],
)
.await
.unwrap();
rt.create_note(
&tok,
"question",
None,
"Review LoRA paper",
Some(0.7),
None,
vec![],
)
.await
.unwrap();
let observations = rt
.list_notes(&tok, Some("observation"), 50, 0)
.await
.unwrap();
assert_eq!(observations.len(), 2);
let questions = rt.list_notes(&tok, Some("question"), 50, 0).await.unwrap();
assert_eq!(questions.len(), 1);
assert_eq!(questions[0].content, "Review LoRA paper");
let all = rt.list_notes(&tok, None, 50, 0).await.unwrap();
assert_eq!(all.len(), 3);
}
#[tokio::test]
async fn create_all_note_kinds() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
for kind in [
"observation",
"insight",
"question",
"decision",
"reference",
"head",
] {
let content = if kind == "head" { "{}" } else { "content" };
rt.create_note(&tok, kind, None, content, Some(0.5), None, vec![])
.await
.unwrap();
}
let all = rt.list_notes(&tok, None, 50, 0).await.unwrap();
assert_eq!(all.len(), 6);
}
#[tokio::test]
async fn query_via_gql() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let lora = rt
.create_entity_with_embedding_report(&tok, "concept", None, "LoRA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let qlora = rt
.create_entity_with_embedding_report(&tok, "concept", None, "QLoRA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, qlora.id, lora.id, EdgeRelation::VariantOf, 1.0, None)
.await
.unwrap();
let rows = rt
.query(
&tok,
"MATCH (a:concept)-[e:variant_of]->(b:concept) RETURN a, e, b LIMIT 10",
)
.await
.unwrap();
assert_eq!(rows.len(), 1);
let first_row = &rows[0];
assert!(first_row.get("a_name").is_some() || first_row.get("a_kind").is_some());
}
#[tokio::test]
async fn query_via_gql_inline_property_map_integer_literal_matches_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({"number": 54});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"PR #54",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(&tok, "MATCH (n:artifact {number: 54}) RETURN n")
.await
.unwrap();
assert_eq!(
rows.len(),
1,
"unquoted integer literal must match the JSON-number property"
);
}
#[tokio::test]
async fn query_via_gql_inline_property_map_quoted_number_does_not_match_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({"number": 54});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"PR #54",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(&tok, "MATCH (n:artifact {number: '54'}) RETURN n")
.await
.unwrap();
assert_eq!(
rows.len(),
0,
"quoted string literal must not match a JSON-number property; \
this is the decided behavior, not a residual bug"
);
}
#[tokio::test]
async fn query_via_gql_inline_property_map_large_integer_matches_exact_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({"number": 9007199254740993i64});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"big",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(
&tok,
"MATCH (n:artifact {number: 9007199254740993}) RETURN n",
)
.await
.unwrap();
assert_eq!(
rows.len(),
1,
"2^53+1 integer literal must match the exact JSON-number property, not round to 2^53"
);
}
#[tokio::test]
async fn query_via_gql_where_equality_large_integer_matches_exact_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({"number": 9007199254740993i64});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"big",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(
&tok,
"MATCH (n:artifact) WHERE n.properties.number = 9007199254740993 RETURN n",
)
.await
.unwrap();
assert_eq!(rows.len(), 1);
}
#[tokio::test]
async fn query_via_gql_inline_property_map_i64_bounds_match_exact_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
for bound in [i64::MIN, i64::MAX] {
let props = serde_json::json!({"number": bound});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"bound",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(
&tok,
&format!("MATCH (n:artifact {{number: {bound}}}) RETURN n"),
)
.await
.unwrap();
assert_eq!(
rows.len(),
1,
"i64 bound {bound} must match its exact JSON-number property"
);
}
}
#[tokio::test]
async fn query_via_gql_where_equality_i64_bounds_match_exact_json_number() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
for bound in [i64::MIN, i64::MAX] {
let props = serde_json::json!({"number": bound});
rt.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"bound",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let rows = rt
.query(
&tok,
&format!("MATCH (n:artifact) WHERE n.properties.number = {bound} RETURN n"),
)
.await
.unwrap();
assert_eq!(
rows.len(),
1,
"i64 bound {bound} must match its exact JSON-number property via WHERE equality"
);
}
}
async fn seed_concepts(rt: &KhiveRuntime, ns: &str, n: usize) -> khive_runtime::NamespaceToken {
let tok = rt.authorize(Namespace::parse(ns).unwrap()).unwrap();
for i in 0..n {
rt.create_entity_with_embedding_report(
&tok,
"concept",
None,
&format!("seed-{i}"),
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
}
tok
}
async fn gql_page(
rt: &KhiveRuntime,
token: &khive_runtime::NamespaceToken,
query: &str,
page_size: usize,
) -> khive_runtime::QueryResult {
rt.query_with_metadata(
token,
query,
khive_query::CompileOptions {
max_limit: page_size,
..Default::default()
},
)
.await
.unwrap()
}
#[tokio::test]
async fn no_explicit_limit_under_and_at_cap_emits_no_warning() {
let rt = rt();
let tok_499 = seed_concepts(&rt, "trunc-499", 499).await;
let result_499 = rt
.query_with_metadata(
&tok_499,
"MATCH (a:concept) RETURN a",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result_499.rows.len(), 499);
assert!(
result_499.warnings.is_empty(),
"499 matches under the cap must not warn: {:?}",
result_499.warnings
);
assert!(
!result_499.truncated,
"#1247: under the cap is not truncated"
);
let tok_500 = seed_concepts(&rt, "trunc-500", 500).await;
let result_500 = rt
.query_with_metadata(
&tok_500,
"MATCH (a:concept) RETURN a",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result_500.rows.len(), 500);
assert!(
result_500.warnings.is_empty(),
"exactly 500 matches must not warn (nothing was dropped): {:?}",
result_500.warnings
);
assert!(
!result_500.truncated,
"#1247: exactly at the cap is not truncated (nothing was dropped)"
);
}
#[tokio::test]
async fn no_explicit_limit_over_cap_warns_and_strips_sentinel() {
let rt = rt();
let tok = seed_concepts(&rt, "trunc-501", 501).await;
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.rows.len(),
500,
"sentinel row must be stripped; exactly max_limit rows must be returned"
);
assert_eq!(result.warnings.len(), 1, "warnings: {:?}", result.warnings);
assert!(result.warnings[0].contains("500"), "{}", result.warnings[0]);
assert!(
result.warnings[0].contains("SKIP 500") && result.warnings[0].contains("next_offset"),
"#1601: the warning must agree with the machine-readable GQL continuation: {}",
result.warnings[0]
);
assert!(result.has_more);
assert_eq!(result.next_offset, Some(500));
assert!(
result.truncated,
"#1247: truncated must be the structural signal, independent of warnings text"
);
let names: std::collections::HashSet<_> = result
.rows
.iter()
.filter_map(|r| match r.get("a_name") {
Some(khive_storage::types::SqlValue::Text(s)) => Some(s.clone()),
_ => None,
})
.collect();
assert_eq!(names.len(), 500, "sentinel row must not leak into results");
let final_page = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a SKIP 500", 500).await;
assert_eq!((final_page.offset, final_page.page_size), (500, 500));
assert_eq!(
final_page.rows.len(),
1,
"SKIP 500 must retrieve the 501st match"
);
assert!(!final_page.has_more && !final_page.truncated);
assert_eq!(final_page.next_offset, None);
let final_name = match final_page.rows[0].get("a_name") {
Some(khive_storage::types::SqlValue::Text(name)) => name,
other => panic!("final page must project the remaining entity name, got {other:?}"),
};
assert!(
!names.contains(final_name),
"the final page must not overlap the first 500 rows"
);
}
#[tokio::test]
async fn explicit_limit_variants_against_501_matches() {
let rt = rt();
let tok = seed_concepts(&rt, "trunc-501-limits", 501).await;
let above_cap = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 600",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(above_cap.rows.len(), 500);
assert_eq!(
above_cap.warnings.len(),
1,
"LIMIT above cap with 501 real matches must warn: {:?}",
above_cap.warnings
);
assert!(above_cap.warnings[0].contains("600"));
assert!(above_cap.warnings[0].contains("500"));
assert!(
above_cap.truncated,
"#1247: LIMIT above cap must set truncated"
);
let at_cap = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 500",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(at_cap.rows.len(), 500);
assert!(
at_cap.warnings.is_empty(),
"LIMIT == cap must not warn: {:?}",
at_cap.warnings
);
assert!(
!at_cap.truncated,
"#1247: an explicit LIMIT the caller chose is not server-side truncation"
);
let below_cap = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 100",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(below_cap.rows.len(), 100);
assert!(
below_cap.warnings.is_empty(),
"LIMIT below cap must not warn: {:?}",
below_cap.warnings
);
assert!(
!below_cap.truncated,
"#1247: LIMIT below cap is not truncated"
);
}
#[tokio::test]
async fn explicit_limit_over_cap_with_few_real_matches_emits_no_warning() {
let rt = rt();
let tok = seed_concepts(&rt, "trunc-20-limit-over-cap", 20).await;
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 600",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result.rows.len(), 20);
assert!(
result.warnings.is_empty(),
"LIMIT above cap with fewer real matches than the cap must not warn: {:?}",
result.warnings
);
assert!(
!result.truncated,
"#1247: fewer real matches than the cap is not truncation"
);
}
#[tokio::test]
async fn gql_skip_pages_are_deterministic_complete_and_machine_continuable() {
let rt = rt();
let tok = seed_concepts(&rt, "gql-pages-12", 12).await;
let names = |result: &khive_runtime::QueryResult| {
result
.rows
.iter()
.filter_map(|row| match row.get("a_name") {
Some(khive_storage::types::SqlValue::Text(name)) => Some(name.clone()),
_ => None,
})
.collect::<Vec<_>>()
};
let first = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a", 5).await;
let first_repeat = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a", 5).await;
let second = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a SKIP 5", 5).await;
let third = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a SKIP 10", 5).await;
assert_eq!(
names(&first),
names(&first_repeat),
"page order must be stable"
);
assert_eq!((first.offset, first.page_size), (0, 5));
assert!(first.has_more && first.truncated);
assert_eq!(first.next_offset, Some(5));
assert!(
first
.warnings
.iter()
.any(|warning| warning.contains("SKIP 5")),
"human guidance must agree with next_offset: {:?}",
first.warnings
);
assert_eq!((second.offset, second.page_size), (5, 5));
assert!(second.has_more && second.truncated);
assert_eq!(second.next_offset, Some(10));
assert_eq!((third.offset, third.page_size), (10, 5));
assert!(!third.has_more && !third.truncated);
assert_eq!(third.next_offset, None);
assert_eq!(third.rows.len(), 2);
let all_names = names(&first)
.into_iter()
.chain(names(&second))
.chain(names(&third))
.collect::<std::collections::HashSet<_>>();
assert_eq!(
all_names.len(),
12,
"paging must retrieve the full match set"
);
}
#[tokio::test]
async fn gql_query_limit_composes_with_page_size() {
let rt = rt();
let tok = seed_concepts(&rt, "gql-page-limit-composition", 8).await;
let caller_bounded = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 3",
khive_query::CompileOptions {
max_limit: 5,
..Default::default()
},
)
.await
.unwrap();
assert_eq!(caller_bounded.rows.len(), 3);
assert_eq!(caller_bounded.page_size, 3);
assert!(!caller_bounded.has_more);
assert_eq!(caller_bounded.next_offset, None);
let server_bounded = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a LIMIT 8",
khive_query::CompileOptions {
max_limit: 5,
..Default::default()
},
)
.await
.unwrap();
assert_eq!(server_bounded.rows.len(), 5);
assert_eq!(server_bounded.page_size, 5);
assert!(server_bounded.has_more);
assert_eq!(server_bounded.next_offset, Some(5));
let terminal = rt
.query_with_metadata(
&tok,
"MATCH (a:concept) RETURN a SKIP 5 LIMIT 8",
khive_query::CompileOptions {
max_limit: 5,
..Default::default()
},
)
.await
.unwrap();
assert_eq!(terminal.rows.len(), 3);
assert!(!terminal.has_more);
assert_eq!(terminal.next_offset, None);
}
#[tokio::test]
async fn gql_query_limit_is_a_total_bound_across_skip_pages() {
let rt = rt();
let tok = seed_concepts(&rt, "gql-page-limit-total-bound", 20).await;
let mut collected_ids = std::collections::HashSet::new();
let mut page_counts = Vec::new();
let mut offset = 0usize;
loop {
let page = gql_page(
&rt,
&tok,
&format!("MATCH (a:concept) RETURN a SKIP {offset} LIMIT 8"),
5,
)
.await;
page_counts.push(page.rows.len());
for row in &page.rows {
if let Some(khive_storage::types::SqlValue::Text(id)) = row.get("a_id") {
assert!(
collected_ids.insert(id.clone()),
"row {id} returned on more than one page"
);
}
}
match page.next_offset {
Some(next) => {
assert!(page.has_more);
offset = next;
}
None => {
assert!(!page.has_more, "terminal page must not claim has_more");
break;
}
}
assert!(
page_counts.len() <= 10,
"paging did not terminate: {page_counts:?}"
);
}
assert_eq!(
collected_ids.len(),
8,
"exactly LIMIT rows must be returned across all pages, not the 20 real matches"
);
assert_eq!(
page_counts,
vec![5, 3],
"page sizes must respect the remaining LIMIT allowance, not a flat page_size"
);
let past_limit = gql_page(&rt, &tok, "MATCH (a:concept) RETURN a SKIP 8 LIMIT 8", 5).await;
assert_eq!(past_limit.rows.len(), 0);
assert!(!past_limit.has_more);
assert_eq!(past_limit.next_offset, None);
let well_past_limit =
gql_page(&rt, &tok, "MATCH (a:concept) RETURN a SKIP 20 LIMIT 8", 5).await;
assert_eq!(well_past_limit.rows.len(), 0);
assert!(!well_past_limit.has_more);
assert_eq!(well_past_limit.next_offset, None);
}
#[tokio::test]
async fn query_static_impossible_pattern_precedes_concept_concept_warns_and_returns_empty() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let d1 = rt
.create_entity_with_embedding_report(&tok, "document", None, "Doc1", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let d2 = rt
.create_entity_with_embedding_report(&tok, "document", None, "Doc2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, d1.id, d2.id, EdgeRelation::Precedes, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (n:concept)-[:precedes]->(m:concept) RETURN n.name, m.name LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.rows.is_empty(),
"no concept->concept precedes edges exist: {:?}",
result.rows
);
assert_eq!(result.warnings.len(), 1, "warnings: {:?}", result.warnings);
let w = &result.warnings[0];
assert!(w.contains("precedes"), "{w}");
assert!(w.contains("concept"), "{w}");
assert!(
w.contains("document->document"),
"must name an accepted pair for the relation: {w}"
);
}
#[tokio::test]
async fn query_static_possible_pattern_extends_concept_concept_no_warning() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "LoRA2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "QLoRA2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, b.id, a.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (n:concept)-[:extends]->(m:concept) RETURN n, m LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.rows.len(),
1,
"concept->concept is a valid base-contract pair for extends"
);
assert!(
result.warnings.is_empty(),
"a statically-valid pattern must not warn: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_check_skips_unlabeled_source_node() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (n)-[:precedes]->(m:document) RETURN n, m LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.warnings.is_empty(),
"an unlabeled endpoint names no static triple to test: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_check_honors_pack_extended_endpoint_rules() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::PartOf,
source: khive_types::EndpointKind::EntityOfKind("person"),
target: khive_types::EndpointKind::EntityOfKind("org"),
}]);
let p = rt
.create_entity_with_embedding_report(&tok, "person", None, "Ada", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let o = rt
.create_entity_with_embedding_report(&tok, "org", None, "Acme", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, p.id, o.id, EdgeRelation::PartOf, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (n:person)-[:part_of]->(m:org) RETURN n, m LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result.rows.len(), 1);
assert!(
result.warnings.is_empty(),
"a pack-declared endpoint rule must not be falsely flagged as impossible: {:?}",
result.warnings
);
}
#[tokio::test]
async fn link_and_hint_agree_special_relation_pack_rules_never_enforced() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::Supersedes,
source: khive_types::EndpointKind::EntityOfKind("person"),
target: khive_types::EndpointKind::EntityOfKind("person"),
}]);
let p1 = rt
.create_entity_with_embedding_report(&tok, "person", None, "Old Ada", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let p2 = rt
.create_entity_with_embedding_report(&tok, "person", None, "New Ada", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let link_result = rt
.link(&tok, p1.id, p2.id, EdgeRelation::Supersedes, 1.0, None)
.await;
assert!(
link_result.is_err(),
"person->person supersedes must be rejected: the special-relation branch \
never consults the installed pack rule; got {link_result:?}"
);
let query_result = rt
.query_with_metadata(
&tok,
"MATCH (a:person)-[:supersedes]->(b:person) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
query_result.warnings.len(),
1,
"the hint must agree with the validator and warn even though a pack rule \
is installed for this special relation: {:?}",
query_result.warnings
);
}
#[tokio::test]
async fn query_static_possible_pattern_inbound_introduced_by_no_warning() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let doc = rt
.create_entity_with_embedding_report(&tok, "document", None, "Paper", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let author = rt
.create_entity_with_embedding_report(&tok, "person", None, "Author", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(
&tok,
doc.id,
author.id,
EdgeRelation::IntroducedBy,
1.0,
None,
)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (p:person)<-[:introduced_by]-(d:document) RETURN d, p LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result.rows.len(), 1);
assert!(
result.warnings.is_empty(),
"an inbound-arrow pattern naming a valid base-contract triple must not warn: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_possible_pattern_entity_of_type_pack_rule_no_false_warning() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::DependsOn,
source: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "theorem",
},
target: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "definition",
},
}]);
let thm = rt
.create_entity_with_embedding_report(
&tok,
"concept",
Some("theorem"),
"T1",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let def = rt
.create_entity_with_embedding_report(
&tok,
"concept",
Some("definition"),
"D1",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, thm.id, def.id, EdgeRelation::DependsOn, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept {entity_type: 'theorem'})-[:depends_on]->\
(b:concept {entity_type: 'definition'}) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result.rows.len(), 1);
assert!(
result.warnings.is_empty(),
"the hint must consult the pack's EntityOfType rule (matched against \
each node's entity_type, not just its base kind) so this must not be \
falsely flagged: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_possible_pattern_untyped_endpoints_with_entity_of_type_rule_no_false_warning()
{
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::DependsOn,
source: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "theorem",
},
target: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "definition",
},
}]);
let thm = rt
.create_entity_with_embedding_report(
&tok,
"concept",
Some("theorem"),
"T1",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let def = rt
.create_entity_with_embedding_report(
&tok,
"concept",
Some("definition"),
"D1",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, thm.id, def.id, EdgeRelation::DependsOn, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept)-[:depends_on]->(b:concept) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.rows.len(),
1,
"the stored theorem->definition edge must be returned: {:?}",
result.rows
);
assert!(
result.warnings.is_empty(),
"an untyped pattern endpoint has not ruled out any entity_type, so it must not \
be falsely flagged as impossible just because it names no entity_type filter: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_pattern_mismatching_entity_type_still_warns() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::DependsOn,
source: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "theorem",
},
target: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "definition",
},
}]);
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept {entity_type: 'lemma'})-[:depends_on]->\
(b:concept {entity_type: 'definition'}) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.warnings.len(),
1,
"a pattern entity_type that mismatches the installed rule's subtype must still \
be flagged as impossible: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_warning_accepted_pairs_include_entity_of_type_only_relation() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.install_edge_rules(vec![khive_types::EdgeEndpointRule {
relation: EdgeRelation::DependsOn,
source: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "theorem",
},
target: khive_types::EndpointKind::EntityOfType {
kind: "concept",
entity_type: "definition",
},
}]);
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:person)-[:depends_on]->(b:person) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(result.warnings.len(), 1, "warnings: {:?}", result.warnings);
assert!(
result.warnings[0].contains("concept->concept"),
"the accepted-pairs list must include the EntityOfType-derived concept->concept \
pair for depends_on, not just base-allowlist pairs: {}",
result.warnings[0]
);
}
#[tokio::test]
async fn query_static_impossible_chained_pattern_warns_once_per_edge() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept)-[:precedes]->(b:concept)-[:precedes]->(c:concept) RETURN a, b, c LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.warnings.len(),
2,
"a chained pattern with two statically-impossible edges sharing a node \
must warn once per edge: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_check_skips_multi_relation_edge() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept)-[:precedes|contains]->(b:concept) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.warnings.is_empty(),
"a multi-relation edge names no single static triple to test: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_check_skips_variable_length_edge() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept)-[:precedes*1..2]->(b:concept) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.warnings.is_empty(),
"a variable-length hop is not a single mandatory hop and must not warn: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_check_skips_undirected_edge() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (a:concept)-[:precedes]-(b:concept) RETURN a, b LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.warnings.is_empty(),
"an undirected edge names no fixed (source, target) triple and must not warn: {:?}",
result.warnings
);
}
#[tokio::test]
async fn query_static_impossible_pattern_warning_exact_message() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.query_with_metadata(
&tok,
"MATCH (n:concept)-[:precedes]->(m:concept) RETURN n, m LIMIT 10",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert_eq!(
result.warnings,
vec![
"pattern (concept)-[:precedes]->(concept) can never match: 'precedes' does not \
accept concept->concept endpoints; accepted source->target kinds for 'precedes': \
document->document, dataset->dataset, project->project, artifact->artifact, \
service->service"
.to_string()
]
);
}
#[tokio::test]
async fn query_static_impossible_check_skips_sparql() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let d1 = rt
.create_entity_with_embedding_report(&tok, "document", None, "Doc1", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let d2 = rt
.create_entity_with_embedding_report(&tok, "document", None, "Doc2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok, d1.id, d2.id, EdgeRelation::Precedes, 1.0, None)
.await
.unwrap();
let result = rt
.query_with_metadata(
&tok,
"SELECT ?a ?b WHERE { ?a :precedes ?b . ?a a :concept . ?b a :concept . }",
khive_query::CompileOptions::default(),
)
.await
.unwrap();
assert!(
result.warnings.is_empty(),
"SPARQL must never receive the static GQL-pattern hint: {:?}",
result.warnings
);
}
#[tokio::test]
async fn namespace_isolation() {
let rt = rt();
let ns_a_tok = rt.authorize(Namespace::parse("ns-a").unwrap()).unwrap();
let ns_b_tok = rt.authorize(Namespace::parse("ns-b").unwrap()).unwrap();
rt.create_entity_with_embedding_report(
&ns_a_tok,
"concept",
None,
"EntityA",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
rt.create_entity_with_embedding_report(
&ns_b_tok,
"concept",
None,
"EntityB",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let a_entities = rt
.list_entities(&ns_a_tok, None, None, 50, 0)
.await
.unwrap();
assert_eq!(a_entities.len(), 1);
assert_eq!(a_entities[0].name, "EntityA");
let b_entities = rt
.list_entities(&ns_b_tok, None, None, 50, 0)
.await
.unwrap();
assert_eq!(b_entities.len(), 1);
assert_eq!(b_entities[0].name, "EntityB");
}
#[tokio::test]
async fn create_entity_indexes_into_text_search() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"FlashAttention",
Some("efficient attention mechanism"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let hits = rt
.hybrid_search(&tok, "FlashAttention", None, 10, None, None, &[], None)
.await
.unwrap();
assert!(
hits.iter().any(|h| h.entity_id == entity.id),
"newly created entity should be findable via hybrid_search (text path)"
);
}
#[tokio::test]
async fn create_entity_no_embedding_model_does_not_propagate_vector_error() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let tok = rt.authorize(Namespace::local()).unwrap();
let result = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"SilentVectorSkip",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record);
assert!(
result.is_ok(),
"create_entity must not propagate Unconfigured from vector store"
);
}
#[tokio::test]
async fn hybrid_search_excludes_soft_deleted_entities() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"SoftDeleteMe",
Some("entity that will be soft-deleted"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let hits_before = rt
.hybrid_search(&tok, "SoftDeleteMe", None, 10, None, None, &[], None)
.await
.unwrap();
assert!(
hits_before.iter().any(|h| h.entity_id == entity.id),
"entity should appear in hybrid_search before soft-delete"
);
rt.delete_entity(&tok, entity.id, false).await.unwrap();
let hits_after = rt
.hybrid_search(&tok, "SoftDeleteMe", None, 10, None, None, &[], None)
.await
.unwrap();
assert!(
!hits_after.iter().any(|h| h.entity_id == entity.id),
"soft-deleted entity must not appear in hybrid_search"
);
}
#[tokio::test]
async fn hybrid_search_excludes_hard_deleted_entities() {
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"HardDeleteMe",
Some("entity that will be hard-deleted"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let hits_before = rt
.hybrid_search(&tok, "HardDeleteMe", None, 10, None, None, &[], None)
.await
.unwrap();
assert!(
hits_before.iter().any(|h| h.entity_id == entity.id),
"entity should appear in hybrid_search before hard-delete"
);
rt.delete_entity(&tok, entity.id, true).await.unwrap();
let hits_after = rt
.hybrid_search(&tok, "HardDeleteMe", None, 10, None, None, &[], None)
.await
.unwrap();
assert!(
!hits_after.iter().any(|h| h.entity_id == entity.id),
"hard-deleted entity must not appear in hybrid_search"
);
}
#[tokio::test]
async fn list_notes_excludes_soft_deleted() {
use khive_storage::types::DeleteMode;
let rt = KhiveRuntime::memory().expect("in-memory runtime");
let tok = rt.authorize(Namespace::local()).unwrap();
let note = rt
.create_note(
&tok,
"observation",
None,
"soft-delete-test",
Some(0.9),
None,
vec![],
)
.await
.unwrap();
let notes_before = rt.list_notes(&tok, None, 50, 0).await.unwrap();
assert!(
notes_before.iter().any(|n| n.id == note.id),
"note should appear before soft-delete"
);
rt.notes(&tok)
.unwrap()
.delete_note(note.id, DeleteMode::Soft)
.await
.unwrap();
let notes_after = rt.list_notes(&tok, None, 50, 0).await.unwrap();
assert!(
!notes_after.iter().any(|n| n.id == note.id),
"soft-deleted note must not appear in list"
);
}
#[tokio::test]
async fn file_backed_runtime_persists() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("persist.db");
{
let config = RuntimeConfig {
wal_ceiling_bytes: 0,
wal_ceiling_configured_bytes: 0,
wal_ceiling_source: Default::default(),
wal_ceiling_env_raw: None,
web: Default::default(),
telemetry: Default::default(),
mounts: Vec::new(),
brain: Default::default(),
git_write: Default::default(),
display_timezone: chrono_tz::Tz::UTC,
events_split: None,
db_path: Some(path.clone()),
blob_hydration_bytes: khive_runtime::DEFAULT_BLOB_HYDRATION_BYTES,
default_namespace: Namespace::local(),
embedding_model: None,
gate: std::sync::Arc::new(khive_runtime::AllowAllGate),
packs: vec!["kg".to_string()],
backend_id: khive_runtime::BackendId::main(),
additional_embedding_models: vec![],
brain_profile: None,
visible_namespaces: vec![],
allowed_outbound_namespaces: vec![],
actor_id: None,
exec: Default::default(),
..khive_runtime::RuntimeConfig::no_embeddings()
};
let rt = KhiveRuntime::new_for_test(config).unwrap();
let tok = rt.authorize(Namespace::local()).unwrap();
rt.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Persistent",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
}
{
let config = RuntimeConfig {
wal_ceiling_bytes: 0,
wal_ceiling_configured_bytes: 0,
wal_ceiling_source: Default::default(),
wal_ceiling_env_raw: None,
web: Default::default(),
telemetry: Default::default(),
mounts: Vec::new(),
brain: Default::default(),
git_write: Default::default(),
display_timezone: chrono_tz::Tz::UTC,
events_split: None,
db_path: Some(path.clone()),
blob_hydration_bytes: khive_runtime::DEFAULT_BLOB_HYDRATION_BYTES,
default_namespace: Namespace::local(),
embedding_model: None,
gate: std::sync::Arc::new(khive_runtime::AllowAllGate),
packs: vec!["kg".to_string()],
backend_id: khive_runtime::BackendId::main(),
additional_embedding_models: vec![],
brain_profile: None,
visible_namespaces: vec![],
allowed_outbound_namespaces: vec![],
actor_id: None,
exec: Default::default(),
..khive_runtime::RuntimeConfig::no_embeddings()
};
let rt = KhiveRuntime::new_for_test(config).unwrap();
let tok = rt.authorize(Namespace::local()).unwrap();
let entities = rt.list_entities(&tok, None, None, 50, 0).await.unwrap();
assert_eq!(entities.len(), 1);
assert_eq!(entities[0].name, "Persistent");
}
}
#[tokio::test]
async fn synthetic_edge_observed_as_selected_returns_memory_note() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let ns = "local";
let memory_note = rt
.create_note(
&tok,
"memory",
None,
"recalled memory content",
Some(0.9),
None,
vec![],
)
.await
.unwrap();
let memory_id = memory_note.id;
let event_store = rt.events(&tok).unwrap();
let mut event = Event::new(
ns,
"search",
EventKind::SearchExecuted,
SubstrateKind::Note,
"agent:test",
);
event.payload = serde_json::json!({
"result_kind": "note",
"candidates": [],
"selected": [memory_id.to_string()]
});
event_store.append_event(event).await.unwrap();
let rows = rt
.query(
&tok,
"MATCH (ev)-[:observed_as_selected]->(m:memory) RETURN m",
)
.await
.unwrap();
assert!(
!rows.is_empty(),
"CRIT-1: synthetic edge query must return at least one row (memory note was seeded); \
got 0 rows — event_observations join is broken"
);
let memory_id_str = memory_id.to_string();
let found = rows.iter().any(|row| {
row.columns.iter().any(|col| {
if let khive_storage::types::SqlValue::Text(s) = &col.value {
s.contains(&memory_id_str)
} else {
false
}
})
});
assert!(
found,
"CRIT-1: returned rows must include the seeded memory note id {}; columns: {:?}",
memory_id,
rows.iter()
.map(|r| r
.columns
.iter()
.map(|c| (&c.name, &c.value))
.collect::<Vec<_>>())
.collect::<Vec<_>>()
);
}
#[tokio::test]
async fn update_edge_returns_surviving_canonical_id_on_conflict() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "SurvA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "SurvB", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let e1 = rt
.link(&tok, a.id, b.id, EdgeRelation::CompetesWith, 1.0, None)
.await
.unwrap();
let (src, tgt) = if a.id > b.id {
(a.id, b.id)
} else {
(b.id, a.id)
};
let e2 = rt
.link(&tok, src, tgt, EdgeRelation::Extends, 0.5, None)
.await
.unwrap();
assert_ne!(
e1.id, e2.id,
"pre-condition: E1 and E2 must be distinct edges"
);
let returned = rt
.update_edge(
&tok,
e2.id.into(),
EdgePatch {
relation: Some(EdgeRelation::CompetesWith),
weight: Some(0.9),
..Default::default()
},
)
.await
.expect("update_edge must succeed even when conflict is absorbed");
assert_eq!(
returned.id, e1.id,
"Bug 1: update_edge must return the SURVIVING canonical row id (E1={:?}), \
got E2={:?}",
e1.id, returned.id
);
let fetched = rt
.get_edge(&tok, returned.id.into())
.await
.expect("get_edge on returned id must not error")
.expect("get_edge on returned id must find a row (not 404)");
assert_eq!(
fetched.id, e1.id,
"fetched row id must match E1 (surviving canonical)"
);
let e2_lookup = rt
.get_edge(&tok, e2.id.into())
.await
.expect("get_edge on deleted id must not error");
assert!(
e2_lookup.is_none(),
"Bug 1: deleted edge E2 must not be findable after conflict absorption"
);
}
#[tokio::test]
async fn update_edge_canonical_orientation_conflict() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "CanOrA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "CanOrB", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let (canon_lo, canon_hi) = if a.id < b.id {
(a.id, b.id)
} else {
(b.id, a.id)
};
let e1 = rt
.link(
&tok,
canon_lo,
canon_hi,
EdgeRelation::CompetesWith,
1.0,
None,
)
.await
.unwrap();
let e2 = rt
.link(&tok, canon_lo, canon_hi, EdgeRelation::Extends, 0.5, None)
.await
.unwrap();
assert_ne!(
e1.id, e2.id,
"pre-condition: E1 and E2 must be distinct edges"
);
rt.update_edge(
&tok,
e2.id.into(),
EdgePatch {
relation: Some(EdgeRelation::CompetesWith),
..Default::default()
},
)
.await
.expect("Bug 2: update_edge on canonical-orientation conflict must not error");
let edges = rt
.list_edges(
&tok,
khive_runtime::EdgeListFilter {
source_id: Some(canon_lo),
target_id: Some(canon_hi),
relations: vec![EdgeRelation::CompetesWith],
..Default::default()
},
100,
0,
)
.await
.expect("list_edges must succeed");
assert_eq!(
edges.len(),
1,
"Bug 2: exactly one competes_with edge must exist after non-flipped conflict absorption; \
found {} edges: {edges:?}",
edges.len()
);
}
#[tokio::test]
async fn entity_create_blocks_secret_in_properties() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({ "api_key": "AKIAFAKEKEY1234567890" });
let result = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"TestEntity",
None,
Some(props),
vec![],
)
.await
.map(|(record, _report)| record);
assert!(
result.is_err(),
"entity create with secret in properties must be blocked"
);
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"error must be SecretDetected"
);
}
#[tokio::test]
async fn entity_create_blocks_secret_in_tags() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let tags = vec![
"type:concept".to_string(),
"AKIAFAKEKEY1234567890".to_string(),
];
let result = rt
.create_entity_with_embedding_report(&tok, "concept", None, "TestEntity", None, None, tags)
.await
.map(|(record, _report)| record);
assert!(
result.is_err(),
"entity create with secret in tags must be blocked"
);
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"error must be SecretDetected"
);
}
#[tokio::test]
async fn note_create_blocks_secret_in_properties() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let props = serde_json::json!({ "api_key": "AKIAFAKEKEY1234567890" });
let result = rt
.create_note(
&tok,
"observation",
None,
"Safe content",
None,
Some(props),
vec![],
)
.await;
assert!(
result.is_err(),
"note create with secret in properties must be blocked"
);
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"error must be SecretDetected"
);
}
#[tokio::test]
async fn note_create_blocks_hex_credential_in_content() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let content = "api key 4f9c2e8a1d3b5c7e9f0a2b4d6e8c0a2b"; let result = rt
.create_note(&tok, "observation", None, content, None, None, vec![])
.await;
assert!(
result.is_err(),
"note create with hex credential in content must be blocked; got Ok"
);
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"error must be SecretDetected"
);
}
#[tokio::test]
async fn note_create_allows_source_path_near_ordinary_key_prose() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let content = "see <a/path/to/file.py>:~97-103 lists it as a real checkpoint-supplied key";
rt.create_note(&tok, "question", None, content, None, None, vec![])
.await
.expect("source path in technical prose must be stored");
}
#[tokio::test]
async fn note_create_allows_git_revision_near_ordinary_token_prose() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let content = "revision d362950a3c9b1a4cb47d97f1623e38f1a1e6bcdf emits one extra token";
rt.create_note(&tok, "question", None, content, None, None, vec![])
.await
.expect("git revision in technical prose must be stored");
}
#[tokio::test]
async fn note_create_blocks_forty_hex_in_value_syntax_behind_vcs_marker() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let content = "api key value is commit d362950a3c9b1a4cb47d97f1623e38f1a1e6bcdf";
let result = rt
.create_note(&tok, "question", None, content, None, None, vec![])
.await;
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"a credential phrase must not be rescued by a VCS marker at the write path"
);
}
#[tokio::test]
async fn note_create_blocks_path_dressed_base64_credential_in_value_syntax() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let content = "api key value is <Xk9mZ2vQpLrT8nJwYuA/HfBsDcGiONvMabcdefgh>:~97-103";
let result = rt
.create_note(&tok, "question", None, content, None, None, vec![])
.await;
assert!(
matches!(
result.unwrap_err(),
khive_runtime::RuntimeError::SecretDetected(_)
),
"path dressing must not rescue a credential in value syntax at the write path"
);
}
mod embedder_registry_tests {
use async_trait::async_trait;
use khive_gate::AllowAllGate;
use khive_runtime::{
EmbedderProvider, KhiveRuntime, NamespaceToken, PackRuntime, RequestIdentity,
RuntimeConfig, RuntimeError, VerbRegistry, VerbRegistryBuilder,
};
use khive_types::{HandlerDef, Namespace, Pack, VerbCategory, Visibility};
use lattice_embed::{EmbeddingModel, EmbeddingService};
use serde_json::Value;
use std::sync::Arc;
struct MockEmbedderProvider {
name: String,
dims: usize,
}
impl MockEmbedderProvider {
fn new(name: &str, dims: usize) -> Self {
Self {
name: name.to_owned(),
dims,
}
}
}
struct MockEmbeddingService {
dims: usize,
}
#[async_trait]
impl EmbeddingService for MockEmbeddingService {
async fn embed(
&self,
texts: &[String],
_model: EmbeddingModel,
) -> Result<Vec<Vec<f32>>, lattice_embed::EmbedError> {
Ok(texts.iter().map(|_| vec![42.0_f32; self.dims]).collect())
}
fn supports_model(&self, _model: EmbeddingModel) -> bool {
true
}
fn name(&self) -> &'static str {
"mock-embedding-service"
}
}
#[async_trait]
impl EmbedderProvider for MockEmbedderProvider {
fn name(&self) -> &str {
&self.name
}
fn dimensions(&self) -> usize {
self.dims
}
async fn build(&self) -> Result<Arc<dyn EmbeddingService>, RuntimeError> {
tokio::time::sleep(std::time::Duration::from_millis(1)).await;
Ok(Arc::new(MockEmbeddingService { dims: self.dims }))
}
}
struct EmbedderInitPack {
runtime: KhiveRuntime,
}
impl Pack for EmbedderInitPack {
const NAME: &'static str = "embedder-init-test";
const NOTE_KINDS: &'static [&'static str] = &[];
const ENTITY_KINDS: &'static [&'static str] = &["concept"];
const HANDLERS: &'static [HandlerDef] = &[HandlerDef {
name: "initialize_embedder",
description: "initialize the test embedder",
visibility: Visibility::Verb,
category: VerbCategory::Commissive,
params: &[],
}];
}
#[async_trait]
impl PackRuntime for EmbedderInitPack {
fn name(&self) -> &str {
Self::NAME
}
fn note_kinds(&self) -> &'static [&'static str] {
Self::NOTE_KINDS
}
fn entity_kinds(&self) -> &'static [&'static str] {
Self::ENTITY_KINDS
}
fn handlers(&self) -> &'static [HandlerDef] {
Self::HANDLERS
}
async fn dispatch(
&self,
_verb: &str,
_params: Value,
_registry: &VerbRegistry,
token: &NamespaceToken,
) -> Result<Value, RuntimeError> {
self.runtime
.create_entity_with_embedding_report(
token,
"concept",
None,
"embedder initialization trigger",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)?;
Ok(Value::Null)
}
}
fn memory_rt_no_model() -> KhiveRuntime {
KhiveRuntime::new(RuntimeConfig {
wal_ceiling_bytes: 0,
wal_ceiling_configured_bytes: 0,
wal_ceiling_source: Default::default(),
wal_ceiling_env_raw: None,
web: Default::default(),
telemetry: Default::default(),
mounts: Vec::new(),
brain: Default::default(),
git_write: Default::default(),
display_timezone: chrono_tz::Tz::UTC,
events_split: None,
db_path: None,
blob_hydration_bytes: khive_runtime::DEFAULT_BLOB_HYDRATION_BYTES,
default_namespace: Namespace::local(),
embedding_model: None,
additional_embedding_models: vec![],
gate: Arc::new(AllowAllGate),
packs: vec!["kg".to_string()],
backend_id: khive_runtime::BackendId::main(),
brain_profile: None,
visible_namespaces: vec![],
allowed_outbound_namespaces: vec![],
actor_id: None,
exec: Default::default(),
..khive_runtime::RuntimeConfig::no_embeddings()
})
.expect("in-memory runtime")
}
#[tokio::test]
async fn register_embedder_and_retrieve_via_embedder_method() {
let rt = memory_rt_no_model();
rt.register_embedder(MockEmbedderProvider::new("mock", 384));
let service = rt
.embedder("mock")
.await
.expect("embedder lookup must succeed after registration");
let texts = vec!["hello world".to_string()];
let vecs = service
.embed(&texts, EmbeddingModel::AllMiniLmL6V2)
.await
.expect("mock service must embed successfully");
assert_eq!(vecs.len(), 1);
assert_eq!(vecs[0].len(), 384);
assert!(
vecs[0].iter().all(|&v| (v - 42.0_f32).abs() < 1e-6),
"mock service must return constant 42.0 vector"
);
}
#[tokio::test]
async fn embedder_initialization_writes_event() {
let rt = memory_rt_no_model();
let token = rt
.authorize(Namespace::local())
.expect("authorize local namespace");
let event_store = rt.events(&token).expect("event store must be available");
rt.register_embedder(MockEmbedderProvider::new("event-test-encoder", 64));
rt.embedder("event-test-encoder")
.await
.expect("embedder initialization must succeed");
let page = event_store
.query_events(
khive_storage::EventFilter::default(),
khive_storage::PageRequest {
limit: 10,
offset: 0,
},
)
.await
.expect("query embedder initialization event");
let event = page
.items
.iter()
.find(|event| event.verb == "embedder.init")
.expect("embedder initialization event must be written");
assert_eq!(event.kind, khive_types::EventKind::EmbedderInitialized);
assert_eq!(event.payload["model_name"], "event-test-encoder");
let duration_us = event.payload["duration_us"]
.as_i64()
.expect("duration_us must be an integer");
assert!(duration_us > 0);
assert_eq!(event.duration_us, duration_us);
}
#[tokio::test]
async fn embedder_initialization_uses_triggering_request_identity() {
let rt = memory_rt_no_model();
rt.register_embedder(MockEmbedderProvider::new("request-identity-encoder", 64));
let tenant = Namespace::parse("tenant-request").expect("valid tenant namespace");
let event_store = rt
.events(&rt.authorize(tenant.clone()).expect("authorize tenant"))
.expect("event store must be available");
let mut builder = VerbRegistryBuilder::new();
builder.register(EmbedderInitPack {
runtime: rt.clone(),
});
builder.with_default_namespace("baked-namespace");
builder.with_actor_id(Some("baked-actor".to_string()));
let registry = builder.build().expect("registry builds");
registry
.dispatch_with_identity(
"initialize_embedder",
serde_json::json!({"namespace": tenant.as_str()}),
Some(RequestIdentity {
namespace: tenant.as_str().to_string(),
actor_id: Some("request-actor".to_string()),
visible_namespaces: vec![],
process_ref: None,
request_id: None,
}),
)
.await
.expect("request-triggered embedder initialization must succeed");
let page = event_store
.query_events(
khive_storage::EventFilter::default(),
khive_storage::PageRequest {
limit: 10,
offset: 0,
},
)
.await
.expect("query embedder initialization event");
let event = page
.items
.iter()
.find(|event| event.verb == "embedder.init")
.expect("embedder initialization event must be written");
assert_eq!(event.namespace, tenant.as_str());
assert_eq!(event.actor, "actor:request-actor");
}
#[tokio::test]
async fn registered_names_includes_custom_provider() {
let rt = memory_rt_no_model();
rt.register_embedder(MockEmbedderProvider::new("my-encoder", 128));
let names = rt.registered_embedding_model_names();
assert!(
names.contains(&"my-encoder".to_string()),
"registered_embedding_model_names must include custom provider 'my-encoder'; got {names:?}"
);
}
#[tokio::test]
async fn dual_embedding_regression_both_models_registered() {
use khive_runtime::RuntimeConfig;
let rt = KhiveRuntime::new(RuntimeConfig {
wal_ceiling_bytes: 0,
wal_ceiling_configured_bytes: 0,
wal_ceiling_source: Default::default(),
wal_ceiling_env_raw: None,
web: Default::default(),
telemetry: Default::default(),
mounts: Vec::new(),
brain: Default::default(),
git_write: Default::default(),
display_timezone: chrono_tz::Tz::UTC,
events_split: None,
db_path: None,
blob_hydration_bytes: khive_runtime::DEFAULT_BLOB_HYDRATION_BYTES,
default_namespace: Namespace::local(),
embedding_model: Some(EmbeddingModel::AllMiniLmL6V2),
additional_embedding_models: vec![EmbeddingModel::ParaphraseMultilingualMiniLmL12V2],
gate: Arc::new(AllowAllGate),
packs: vec!["kg".to_string()],
backend_id: khive_runtime::BackendId::main(),
brain_profile: None,
visible_namespaces: vec![],
allowed_outbound_namespaces: vec![],
actor_id: None,
exec: Default::default(),
..khive_runtime::RuntimeConfig::no_embeddings()
})
.expect("runtime with two models");
let names = rt.registered_embedding_model_names();
assert!(
names.contains(&"all-minilm-l6-v2".to_string()),
"MiniLM must be registered; names: {names:?}"
);
assert!(
names.contains(&"paraphrase-multilingual-minilm-l12-v2".to_string()),
"paraphrase must be registered; names: {names:?}"
);
rt.resolve_embedding_model(Some("all-minilm-l6-v2"))
.expect("MiniLM must resolve");
rt.resolve_embedding_model(Some("paraphrase"))
.expect("paraphrase alias must resolve");
}
#[tokio::test]
async fn embedder_unknown_name_returns_error() {
let rt = memory_rt_no_model();
let err = rt
.embedder("no-such-model")
.await
.err()
.expect("expected Err for unknown embedder name, got Ok");
assert!(
matches!(err, RuntimeError::UnknownModel(ref n) if n == "no-such-model"),
"expected UnknownModel for unregistered name; got {err:?}"
);
}
#[tokio::test]
async fn pack_register_embedders_hook_makes_provider_reachable() {
use async_trait::async_trait;
use khive_runtime::pack::HandlerDef;
use khive_runtime::NamespaceToken;
use khive_runtime::{PackRuntime, VerbRegistry, VerbRegistryBuilder};
use khive_types::Pack;
use serde_json::Value;
struct EmbedderPack;
impl Pack for EmbedderPack {
const NAME: &'static str = "embedder-test-pack";
const NOTE_KINDS: &'static [&'static str] = &[];
const ENTITY_KINDS: &'static [&'static str] = &[];
const HANDLERS: &'static [HandlerDef] = &[];
}
#[async_trait]
impl PackRuntime for EmbedderPack {
fn name(&self) -> &str {
Self::NAME
}
fn note_kinds(&self) -> &'static [&'static str] {
Self::NOTE_KINDS
}
fn entity_kinds(&self) -> &'static [&'static str] {
Self::ENTITY_KINDS
}
fn handlers(&self) -> &'static [HandlerDef] {
Self::HANDLERS
}
fn register_embedders(&self, runtime: &KhiveRuntime) {
runtime.register_embedder(MockEmbedderProvider::new("pack-custom-encoder", 256));
}
async fn dispatch(
&self,
_verb: &str,
_params: Value,
_registry: &VerbRegistry,
_token: &NamespaceToken,
) -> Result<Value, khive_runtime::RuntimeError> {
Ok(Value::Null)
}
}
let rt = memory_rt_no_model();
let mut builder = VerbRegistryBuilder::new();
builder.register(EmbedderPack);
let registry = builder.build().expect("registry builds");
registry.call_register_embedders(&rt);
let service = rt
.embedder("pack-custom-encoder")
.await
.expect("pack-contributed provider must be reachable after call_register_embedders");
let texts = vec!["test sentence".to_string()];
let vecs = service
.embed(&texts, EmbeddingModel::AllMiniLmL6V2)
.await
.expect("custom service must embed without error");
assert_eq!(vecs.len(), 1);
assert_eq!(
vecs[0].len(),
256,
"dims must match provider declaration (256)"
);
}
#[tokio::test]
async fn failing_provider_build_returns_err_not_panic() {
struct FailingProvider;
#[async_trait]
impl EmbedderProvider for FailingProvider {
fn name(&self) -> &str {
"failing-provider"
}
fn dimensions(&self) -> usize {
128
}
async fn build(&self) -> Result<Arc<dyn EmbeddingService>, RuntimeError> {
Err(RuntimeError::Internal(
"simulated provider construction failure".into(),
))
}
}
let rt = memory_rt_no_model();
rt.register_embedder(FailingProvider);
let result = rt.embedder("failing-provider").await;
assert!(
result.is_err(),
"embedder() must return Err when build() fails, not panic; got Ok"
);
let err = result.err().expect("checked above");
let msg = err.to_string();
assert!(
msg.contains("simulated provider construction failure")
|| msg.contains("build() failed")
|| msg.contains("Internal"),
"error must carry build failure context; got: {msg}"
);
}
}
#[tokio::test]
async fn link_concept_concept_supports_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Finding A", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim B", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, a.id, b.id, EdgeRelation::Supports, 0.8, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Supports);
assert_eq!(edge.source_id, a.id);
assert_eq!(edge.target_id, b.id);
}
#[tokio::test]
async fn link_document_concept_supports_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let doc = rt
.create_entity_with_embedding_report(&tok, "document", None, "Paper X", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Hypothesis Y",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, doc.id, claim.id, EdgeRelation::Supports, 0.9, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Supports);
}
#[tokio::test]
async fn link_concept_concept_refutes_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let a = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Counter-evidence",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let b = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim B", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, a.id, b.id, EdgeRelation::Refutes, 0.7, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Refutes);
}
#[tokio::test]
async fn link_document_concept_refutes_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let doc = rt
.create_entity_with_embedding_report(
&tok,
"document",
None,
"Negative study",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim C", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, doc.id, claim.id, EdgeRelation::Refutes, 0.85, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Refutes);
}
#[tokio::test]
async fn link_note_note_supports_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let finding = rt
.create_note(
&tok,
"observation",
Some("Finding note"),
"experiment shows positive result",
Some(0.8),
None,
vec![],
)
.await
.unwrap();
let claim = rt
.create_note(
&tok,
"question",
Some("Claim note"),
"does intervention work?",
Some(0.7),
None,
vec![],
)
.await
.unwrap();
let edge = rt
.link(
&tok,
finding.id,
claim.id,
EdgeRelation::Supports,
0.9,
None,
)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Supports);
assert_eq!(edge.source_id, finding.id);
assert_eq!(edge.target_id, claim.id);
}
#[tokio::test]
async fn link_note_note_refutes_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let counter = rt
.create_note(
&tok,
"observation",
Some("Counter finding"),
"null result from replication",
Some(0.6),
None,
vec![],
)
.await
.unwrap();
let hypothesis = rt
.create_note(
&tok,
"insight",
Some("Hypothesis"),
"the intervention increases outcome",
Some(0.7),
None,
vec![],
)
.await
.unwrap();
let edge = rt
.link(
&tok,
counter.id,
hypothesis.id,
EdgeRelation::Refutes,
0.75,
None,
)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Refutes);
}
#[tokio::test]
async fn link_note_entity_supports_rejected() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let note = rt
.create_note(
&tok,
"observation",
None,
"finding note",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Some concept",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let result = rt
.link(&tok, note.id, entity.id, EdgeRelation::Supports, 0.8, None)
.await;
assert!(
matches!(result, Err(khive_runtime::RuntimeError::InvalidInput(_))),
"note→entity supports must be rejected (cross-substrate); got {result:?}"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("supports"),
"error message must name the relation 'supports'; got: {msg}"
);
}
#[tokio::test]
async fn link_entity_note_refutes_rejected() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(&tok, "concept", None, "A concept", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let note = rt
.create_note(
&tok,
"observation",
None,
"some note",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
let result = rt
.link(&tok, entity.id, note.id, EdgeRelation::Refutes, 0.5, None)
.await;
assert!(
matches!(result, Err(khive_runtime::RuntimeError::InvalidInput(_))),
"entity→note refutes must be rejected (cross-substrate); got {result:?}"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("refutes"),
"error message must name the relation 'refutes'; got: {msg}"
);
}
#[tokio::test]
async fn link_person_concept_supports_rejected_with_relation_name() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let person = rt
.create_entity_with_embedding_report(
&tok,
"person",
None,
"Researcher A",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Hypothesis Z",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let result = rt
.link(&tok, person.id, claim.id, EdgeRelation::Supports, 0.5, None)
.await;
assert!(
matches!(result, Err(khive_runtime::RuntimeError::InvalidInput(_))),
"person→concept supports is not in base allowlist; got {result:?}"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("supports"),
"error message must name the relation 'supports'; got: {msg}"
);
}
#[tokio::test]
async fn link_dataset_concept_supports_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let ds = rt
.create_entity_with_embedding_report(&tok, "dataset", None, "Bench-X", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Hypothesis Q",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, ds.id, claim.id, EdgeRelation::Supports, 0.8, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Supports);
}
#[tokio::test]
async fn link_artifact_concept_refutes_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let art = rt
.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"Checkpoint-v1",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim R", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, art.id, claim.id, EdgeRelation::Refutes, 0.7, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Refutes);
}
#[tokio::test]
async fn link_artifact_concept_supports_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let art = rt
.create_entity_with_embedding_report(
&tok,
"artifact",
None,
"Checkpoint-v2",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim T", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, art.id, claim.id, EdgeRelation::Supports, 0.8, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Supports);
assert_eq!(edge.source_id, art.id);
assert_eq!(edge.target_id, claim.id);
}
#[tokio::test]
async fn link_dataset_concept_refutes_accepted() {
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let ds = rt
.create_entity_with_embedding_report(&tok, "dataset", None, "Bench-Y", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Hypothesis W",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok, ds.id, claim.id, EdgeRelation::Refutes, 0.75, None)
.await
.unwrap();
assert_eq!(edge.relation, EdgeRelation::Refutes);
assert_eq!(edge.source_id, ds.id);
assert_eq!(edge.target_id, claim.id);
}
#[tokio::test]
async fn update_edge_to_supports_on_legal_entity_pair_accepted() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let evidence = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Evidence concept",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let claim = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Hypothesis H",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(
&tok,
evidence.id,
claim.id,
EdgeRelation::Extends,
0.9,
None,
)
.await
.unwrap();
let updated = rt
.update_edge(
&tok,
edge.id.into(),
EdgePatch {
relation: Some(EdgeRelation::Supports),
..Default::default()
},
)
.await
.expect("update_edge to supports on concept→concept must be accepted");
assert_eq!(updated.relation, EdgeRelation::Supports);
}
#[tokio::test]
async fn update_edge_to_supports_on_disallowed_entity_pair_rejected() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let person = rt
.create_entity_with_embedding_report(
&tok,
"person",
None,
"Researcher B",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let concept = rt
.create_entity_with_embedding_report(&tok, "concept", None, "Claim S", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(
&tok,
person.id,
concept.id,
EdgeRelation::InstanceOf,
1.0,
None,
)
.await
.unwrap();
let result = rt
.update_edge(
&tok,
edge.id.into(),
EdgePatch {
relation: Some(EdgeRelation::Supports),
..Default::default()
},
)
.await;
assert!(
matches!(result, Err(khive_runtime::RuntimeError::InvalidInput(_))),
"update_edge to supports on person→concept must be rejected; got {result:?}"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("supports"),
"error message must name the relation 'supports'; got: {msg}"
);
assert!(
msg.contains(
"currently legal relations for person -> concept under the loaded endpoint rules: instance_of"
),
"error message must expose the derived legal-relations set from the shared validator \
update_edge reaches at operations.rs; got: {msg}"
);
}
#[tokio::test]
async fn update_edge_annotates_to_supports_rejected_cross_substrate() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let entity = rt
.create_entity_with_embedding_report(
&tok,
"concept",
None,
"Target concept",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let note = rt
.create_note(
&tok,
"observation",
None,
"some observation",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
let edge = rt
.link(&tok, note.id, entity.id, EdgeRelation::Annotates, 1.0, None)
.await
.unwrap();
let result = rt
.update_edge(
&tok,
edge.id.into(),
EdgePatch {
relation: Some(EdgeRelation::Supports),
..Default::default()
},
)
.await;
assert!(
matches!(result, Err(khive_runtime::RuntimeError::InvalidInput(_))),
"update_edge note→entity annotates → supports must be rejected; got {result:?}"
);
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("supports"),
"error message must name the relation 'supports'; got: {msg}"
);
}
#[tokio::test]
async fn update_edge_note_note_to_refutes_accepted() {
use khive_runtime::EdgePatch;
let rt = rt();
let tok = rt.authorize(Namespace::local()).unwrap();
let note_a = rt
.create_note(
&tok,
"observation",
None,
"prior finding",
Some(0.6),
None,
vec![],
)
.await
.unwrap();
let note_b = rt
.create_note(
&tok,
"insight",
None,
"derived claim",
Some(0.7),
None,
vec![],
)
.await
.unwrap();
let edge = rt
.link(
&tok,
note_a.id,
note_b.id,
EdgeRelation::Supports,
0.8,
None,
)
.await
.unwrap();
let updated = rt
.update_edge(
&tok,
edge.id.into(),
EdgePatch {
relation: Some(EdgeRelation::Refutes),
..Default::default()
},
)
.await
.expect("update_edge note→note supports → refutes must be accepted");
assert_eq!(updated.relation, EdgeRelation::Refutes);
}
#[tokio::test]
async fn visible_set_reads_primary_and_extra_not_third() {
let rt = rt();
let tok_a = rt.authorize(Namespace::parse("vis-a").unwrap()).unwrap();
let tok_b = rt.authorize(Namespace::parse("vis-b").unwrap()).unwrap();
let tok_c = rt.authorize(Namespace::parse("vis-c").unwrap()).unwrap();
let entity_a = rt
.create_entity_with_embedding_report(&tok_a, "concept", None, "EntityA", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let entity_b = rt
.create_entity_with_embedding_report(&tok_b, "concept", None, "EntityB", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let entity_c = rt
.create_entity_with_embedding_report(&tok_c, "concept", None, "EntityC", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let note_a = rt
.create_note(&tok_a, "observation", None, "NoteA", None, None, vec![])
.await
.unwrap();
let note_b = rt
.create_note(&tok_b, "observation", None, "NoteB", None, None, vec![])
.await
.unwrap();
let note_c = rt
.create_note(&tok_c, "observation", None, "NoteC", None, None, vec![])
.await
.unwrap();
let vis_tok = rt
.authorize_with_visibility(
Namespace::parse("vis-a").unwrap(),
vec![Namespace::parse("vis-b").unwrap()],
)
.unwrap();
let visible_entities = rt.list_entities(&vis_tok, None, None, 50, 0).await.unwrap();
let entity_names: Vec<&str> = visible_entities.iter().map(|e| e.name.as_str()).collect();
assert!(entity_names.contains(&"EntityA"), "EntityA must be visible");
assert!(entity_names.contains(&"EntityB"), "EntityB must be visible");
assert!(
!entity_names.contains(&"EntityC"),
"EntityC must NOT be visible"
);
let visible_notes = rt.list_notes(&vis_tok, None, 50, 0).await.unwrap();
let note_contents: Vec<&str> = visible_notes.iter().map(|n| n.content.as_str()).collect();
assert!(note_contents.contains(&"NoteA"), "NoteA must be visible");
assert!(note_contents.contains(&"NoteB"), "NoteB must be visible");
assert!(
!note_contents.contains(&"NoteC"),
"NoteC must NOT be visible"
);
rt.get_entity(&vis_tok, entity_a.id)
.await
.expect("get entity_a must succeed");
rt.get_entity(&vis_tok, entity_b.id)
.await
.expect("get entity_b (visible non-primary) must succeed");
rt.get_entity(&vis_tok, entity_c.id)
.await
.expect("get entity_c by UUID succeeds — visible-set gate removed in PR-A1");
let fetched_note_a = rt
.get_note_including_deleted(&vis_tok, note_a.id)
.await
.expect("call must not error");
assert!(
fetched_note_a.is_some(),
"note_a (primary namespace) must be returned"
);
let fetched_note_b = rt
.get_note_including_deleted(&vis_tok, note_b.id)
.await
.expect("call must not error");
assert!(
fetched_note_b.is_some(),
"note_b (visible non-primary) must be returned"
);
let fetched_note_c = rt
.get_note_including_deleted(&vis_tok, note_c.id)
.await
.expect("call must not error");
assert!(
fetched_note_c.is_some(),
"note_c (outside visible set) must be returned by UUID via PR-A1 by-ID contract"
);
assert_eq!(
fetched_note_c.as_ref().unwrap().namespace.as_str(),
"vis-c",
"fetched note_c must preserve its stored namespace"
);
let written = rt
.create_entity_with_embedding_report(
&vis_tok,
"concept",
None,
"WrittenViaVisToken",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
assert_eq!(
written.namespace.as_str(),
"vis-a",
"write must stamp primary namespace, not any extra-visible one"
);
let b_entities = rt.list_entities(&tok_b, None, None, 50, 0).await.unwrap();
let b_names: Vec<&str> = b_entities.iter().map(|e| e.name.as_str()).collect();
assert!(
!b_names.contains(&"WrittenViaVisToken"),
"write must NOT appear in vis-b"
);
let _ = note_a;
let _ = note_c;
let _ = entity_c;
}
#[tokio::test]
async fn namespace_isolation_backward_compat() {
let rt = rt();
let ns_a_tok = rt.authorize(Namespace::parse("bc-a").unwrap()).unwrap();
let ns_b_tok = rt.authorize(Namespace::parse("bc-b").unwrap()).unwrap();
rt.create_entity_with_embedding_report(
&ns_a_tok,
"concept",
None,
"EntityA",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
rt.create_entity_with_embedding_report(
&ns_b_tok,
"concept",
None,
"EntityB",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let a_entities = rt
.list_entities(&ns_a_tok, None, None, 50, 0)
.await
.unwrap();
assert_eq!(a_entities.len(), 1);
assert_eq!(a_entities[0].name, "EntityA");
let b_entities = rt
.list_entities(&ns_b_tok, None, None, 50, 0)
.await
.unwrap();
assert_eq!(b_entities.len(), 1);
assert_eq!(b_entities[0].name, "EntityB");
}
#[test]
fn mint_with_visibility_empty_extra_yields_primary_only() {
let rt = rt();
let tok = rt
.authorize_with_visibility(Namespace::parse("ns-primary-only").unwrap(), vec![])
.unwrap();
let vis = tok.visible_namespaces();
assert_eq!(vis.len(), 1, "primary only when no extras given");
assert_eq!(vis[0].as_str(), "ns-primary-only");
assert_eq!(tok.namespace().as_str(), "ns-primary-only");
}
#[test]
fn mint_with_visibility_deduplicates_primary_in_extras() {
let rt = rt();
let tok = rt
.authorize_with_visibility(
Namespace::parse("ns-dedup").unwrap(),
vec![
Namespace::parse("ns-dedup").unwrap(),
Namespace::parse("ns-extra").unwrap(),
],
)
.unwrap();
let vis = tok.visible_namespaces();
assert_eq!(vis.len(), 2, "primary counted once, one distinct extra");
assert_eq!(vis[0].as_str(), "ns-dedup");
assert_eq!(vis[1].as_str(), "ns-extra");
}
#[tokio::test]
async fn resolve_uses_visible_set_for_note_in_extra_namespace() {
let rt = rt();
let _tok_a = rt.authorize(Namespace::parse("res-a").unwrap()).unwrap();
let tok_b = rt.authorize(Namespace::parse("res-b").unwrap()).unwrap();
let note_b = rt
.create_note(&tok_b, "observation", None, "NoteInB", None, None, vec![])
.await
.unwrap();
let vis_tok = rt
.authorize_with_visibility(
Namespace::parse("res-a").unwrap(),
vec![Namespace::parse("res-b").unwrap()],
)
.unwrap();
let fetched = rt
.get_note_including_deleted(&vis_tok, note_b.id)
.await
.expect("call must not error");
assert!(
fetched.is_some(),
"note in extra-visible namespace must be readable via visible-set token"
);
assert_eq!(fetched.unwrap().content, "NoteInB");
}
#[tokio::test]
async fn link_target_in_visible_but_not_primary_namespace_succeeds() {
let rt = rt();
let tok_a = rt
.authorize(Namespace::parse("link-mut-a").unwrap())
.unwrap();
let tok_b = rt
.authorize(Namespace::parse("link-mut-b").unwrap())
.unwrap();
let entity_a = rt
.create_entity_with_embedding_report(
&tok_a,
"concept",
None,
"SrcEntity",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let entity_b = rt
.create_entity_with_embedding_report(
&tok_b,
"concept",
None,
"TgtEntity",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let vis_tok = rt
.authorize_with_visibility(
Namespace::parse("link-mut-a").unwrap(),
vec![Namespace::parse("link-mut-b").unwrap()],
)
.unwrap();
let result = rt
.link(
&vis_tok,
entity_a.id,
entity_b.id,
EdgeRelation::Extends,
1.0,
None,
)
.await;
assert!(
result.is_ok(),
"link with target in visible-only namespace must succeed (#631), got {result:?}"
);
}
#[tokio::test]
async fn create_note_annotates_target_in_visible_only_namespace_succeeds() {
let rt = rt();
let _tok_a = rt
.authorize(Namespace::parse("ann-mut-a").unwrap())
.unwrap();
let tok_b = rt
.authorize(Namespace::parse("ann-mut-b").unwrap())
.unwrap();
let entity_b = rt
.create_entity_with_embedding_report(
&tok_b,
"concept",
None,
"AnnotTarget",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let vis_tok = rt
.authorize_with_visibility(
Namespace::parse("ann-mut-a").unwrap(),
vec![Namespace::parse("ann-mut-b").unwrap()],
)
.unwrap();
let result = rt
.create_note(
&vis_tok,
"observation",
None,
"AnnotNote",
None,
None,
vec![entity_b.id],
)
.await;
assert!(
result.is_ok(),
"annotates with target in visible-only namespace must succeed (#631), got {result:?}"
);
}
#[tokio::test]
async fn hybrid_search_surfaces_all_visible_namespaces() {
let rt = rt();
let ns_primary = Namespace::parse("hs-primary-ns").unwrap();
let ns_extra = Namespace::parse("hs-extra-ns").unwrap();
let tok_primary = rt.authorize(ns_primary.clone()).unwrap();
let tok_extra = rt.authorize(ns_extra.clone()).unwrap();
let entity_in_primary = rt
.create_entity_with_embedding_report(
&tok_primary,
"concept",
None,
"StellarPrimary",
Some("unique stellar primary concept"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let entity_in_extra = rt
.create_entity_with_embedding_report(
&tok_extra,
"concept",
None,
"StellarExtra",
Some("unique stellar extra concept"),
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let vis_tok = rt
.authorize_with_visibility(ns_primary.clone(), vec![ns_extra.clone()])
.unwrap();
let hits = rt
.hybrid_search(&vis_tok, "stellar", None, 20, None, None, &[], None)
.await
.unwrap();
let hit_ids: Vec<Uuid> = hits.iter().map(|h| h.entity_id).collect();
assert!(
hit_ids.contains(&entity_in_primary.id),
"hybrid_search must return entity from primary namespace; \
expected entity_id={}, got: {hit_ids:?}",
entity_in_primary.id,
);
assert!(
hit_ids.contains(&entity_in_extra.id),
"hybrid_search must return entity from visible extra namespace; \
entity_id={} missing from: {hit_ids:?}",
entity_in_extra.id,
);
let fetched = rt
.get_entity(&vis_tok, entity_in_extra.id)
.await
.expect("get_entity via visible-set token must return extra-namespace entity");
assert_eq!(
fetched.id, entity_in_extra.id,
"visible-set read of extra-namespace entity must succeed"
);
}
#[tokio::test]
async fn update_note_cross_namespace_succeeds() {
use khive_runtime::NotePatch;
let rt = rt();
let tok_a = rt
.authorize(Namespace::parse("note-ns-a").unwrap())
.unwrap();
let tok_b = rt
.authorize(Namespace::parse("note-ns-b").unwrap())
.unwrap();
let note = rt
.create_note(
&tok_a,
"observation",
None,
"original content",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
assert_eq!(note.namespace.as_str(), "note-ns-a");
let patch = NotePatch::new(None, Some("updated content".to_string()), None, None, None);
let updated = rt
.update_note_with_embedding_report(&tok_b, note.id, patch)
.await
.map(|(record, _report)| record);
assert!(
updated.is_ok(),
"update_note from foreign token must succeed; got {:?}",
updated
);
let updated = updated.unwrap();
assert_eq!(updated.content, "updated content");
assert_eq!(
updated.namespace.as_str(),
"note-ns-a",
"namespace must remain the record's stored namespace after cross-ns update"
);
}
#[tokio::test]
async fn delete_note_cross_namespace_succeeds() {
let rt = rt();
let tok_a = rt.authorize(Namespace::parse("del-ns-a").unwrap()).unwrap();
let tok_b = rt.authorize(Namespace::parse("del-ns-b").unwrap()).unwrap();
let note_soft = rt
.create_note(
&tok_a,
"observation",
None,
"soft target",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
let soft_result = rt.delete_note(&tok_b, note_soft.id, false).await;
assert!(
soft_result.unwrap(),
"cross-namespace soft delete_note must return true"
);
let after_soft = rt
.get_note_including_deleted(&tok_a, note_soft.id)
.await
.unwrap();
assert!(
after_soft.is_some(),
"soft-deleted note must still appear in including_deleted"
);
let note_hard = rt
.create_note(
&tok_a,
"observation",
None,
"hard target",
Some(0.5),
None,
vec![],
)
.await
.unwrap();
let hard_result = rt.delete_note(&tok_b, note_hard.id, true).await;
assert!(
hard_result.unwrap(),
"cross-namespace hard delete_note must return true"
);
let after_hard = rt
.get_note_including_deleted(&tok_a, note_hard.id)
.await
.unwrap();
assert!(
after_hard.is_none(),
"hard-deleted note must not appear even via including_deleted"
);
}
#[tokio::test]
async fn delete_edge_cross_namespace_audit_uses_record_namespace_soft() {
let rt = rt();
let tok_owner = rt.authorize(Namespace::parse("ns-owner").unwrap()).unwrap();
let tok_caller = rt
.authorize(Namespace::parse("ns-caller").unwrap())
.unwrap();
let src = rt
.create_entity_with_embedding_report(
&tok_owner,
"concept",
None,
"AuditSrcSoft",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let tgt = rt
.create_entity_with_embedding_report(
&tok_owner,
"concept",
None,
"AuditTgtSoft",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok_owner, src.id, tgt.id, EdgeRelation::Extends, 0.5, None)
.await
.unwrap();
let edge_id: Uuid = edge.id.into();
let deleted = rt.delete_edge(&tok_caller, edge_id, false).await.unwrap();
assert!(deleted, "cross-namespace soft delete_edge must return true");
let live = rt.get_edge(&tok_owner, edge_id).await.unwrap();
assert!(
live.is_none(),
"soft-deleted edge must not appear in live get_edge"
);
let incl = rt
.get_edge_including_deleted(&tok_owner, edge_id)
.await
.unwrap();
assert!(
incl.is_some(),
"soft-deleted edge must appear via get_edge_including_deleted"
);
let events = rt
.list_events(
&tok_owner,
EventFilter {
kinds: vec![EventKind::EdgeDeleted],
..Default::default()
},
PageRequest::default(),
)
.await
.unwrap();
let delete_event = events
.items
.iter()
.find(|e| e.target_id == Some(edge_id))
.expect("EdgeDeleted event must exist for the deleted edge");
assert_eq!(
delete_event.namespace, "ns-owner",
"EdgeDeleted event namespace must be the record's namespace (ns-owner), not the caller's"
);
assert_eq!(
delete_event
.payload
.get("namespace")
.and_then(|v| v.as_str()),
Some("ns-owner"),
"EdgeDeleted payload.namespace must be the record's namespace (ns-owner)"
);
}
#[tokio::test]
async fn delete_edge_cross_namespace_audit_uses_record_namespace_hard() {
let rt = rt();
let tok_owner = rt
.authorize(Namespace::parse("ns-owner-hard").unwrap())
.unwrap();
let tok_caller = rt
.authorize(Namespace::parse("ns-caller-hard").unwrap())
.unwrap();
let src = rt
.create_entity_with_embedding_report(
&tok_owner,
"concept",
None,
"AuditSrcHard",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let tgt = rt
.create_entity_with_embedding_report(
&tok_owner,
"concept",
None,
"AuditTgtHard",
None,
None,
vec![],
)
.await
.map(|(record, _report)| record)
.unwrap();
let edge = rt
.link(&tok_owner, src.id, tgt.id, EdgeRelation::Extends, 0.6, None)
.await
.unwrap();
let edge_id: Uuid = edge.id.into();
let deleted = rt.delete_edge(&tok_caller, edge_id, true).await.unwrap();
assert!(deleted, "cross-namespace hard delete_edge must return true");
let incl = rt
.get_edge_including_deleted(&tok_owner, edge_id)
.await
.unwrap();
assert!(
incl.is_none(),
"hard-deleted edge must not appear via get_edge_including_deleted"
);
let events = rt
.list_events(
&tok_owner,
EventFilter {
kinds: vec![EventKind::EdgeDeleted],
..Default::default()
},
PageRequest::default(),
)
.await
.unwrap();
let delete_event = events
.items
.iter()
.find(|e| e.target_id == Some(edge_id))
.expect("EdgeDeleted event must exist for the hard-deleted edge");
assert_eq!(
delete_event.namespace, "ns-owner-hard",
"EdgeDeleted event namespace must be record's namespace (ns-owner-hard), not caller's"
);
assert_eq!(
delete_event
.payload
.get("namespace")
.and_then(|v| v.as_str()),
Some("ns-owner-hard"),
"EdgeDeleted payload.namespace must be the record's namespace (ns-owner-hard)"
);
}
#[tokio::test]
async fn stats_totals_match_list_walk_across_visible_namespaces() {
use khive_runtime::EdgeListFilter;
let rt = KhiveRuntime::new(RuntimeConfig {
db_path: None,
packs: vec!["kg".to_string()],
brain_profile: None,
actor_id: Some("lambda:stats-711-test".to_string()),
..RuntimeConfig::no_embeddings()
})
.expect("in-memory runtime with actor identity");
let ns_a = Namespace::parse("stats-711-a").unwrap();
let ns_b = Namespace::parse("stats-711-b").unwrap();
let tok_a = rt.authorize(ns_a.clone()).unwrap();
let tok_b = rt.authorize(ns_b.clone()).unwrap();
let a1 = rt
.create_entity_with_embedding_report(&tok_a, "concept", None, "StatsA1", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let a2 = rt
.create_entity_with_embedding_report(&tok_a, "concept", None, "StatsA2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
let b1 = rt
.create_entity_with_embedding_report(&tok_b, "concept", None, "StatsB1", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok_a, a1.id, a2.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
rt.link(&tok_a, a2.id, a1.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
let b2 = rt
.create_entity_with_embedding_report(&tok_b, "concept", None, "StatsB2", None, None, vec![])
.await
.map(|(record, _report)| record)
.unwrap();
rt.link(&tok_b, b1.id, b2.id, EdgeRelation::Enables, 1.0, None)
.await
.unwrap();
let deleted_edge = rt
.link(&tok_b, b2.id, b1.id, EdgeRelation::Extends, 1.0, None)
.await
.unwrap();
assert!(rt
.delete_edge(&tok_b, deleted_edge.id.into(), false)
.await
.unwrap());
rt.create_note(&tok_a, "observation", None, "NoteInA", None, None, vec![])
.await
.unwrap();
rt.create_note(&tok_b, "observation", None, "NoteInB", None, None, vec![])
.await
.unwrap();
let deleted_note = rt
.create_note(
&tok_a,
"observation",
None,
"DeletedNoteInA",
None,
None,
vec![],
)
.await
.unwrap();
assert!(rt
.delete_note(&tok_a, deleted_note.id, false)
.await
.unwrap());
let per_namespace_edges = rt
.count_edges(&tok_a, EdgeListFilter::default())
.await
.unwrap()
+ rt.count_edges(&tok_b, EdgeListFilter::default())
.await
.unwrap();
let mut per_namespace_relations = rt.count_edges_by_relation(&tok_a).await.unwrap();
for (relation, count) in rt.count_edges_by_relation(&tok_b).await.unwrap() {
*per_namespace_relations.entry(relation).or_insert(0) += count;
}
let per_namespace_notes =
rt.count_notes(&tok_a, None).await.unwrap() + rt.count_notes(&tok_b, None).await.unwrap();
let vis_tok = rt
.authorize_with_visibility(ns_a.clone(), vec![ns_b.clone()])
.unwrap();
let full_entities = rt
.list_entities(&vis_tok, None, None, 500, 0)
.await
.unwrap();
let stats_entities = rt.count_entities(&vis_tok, None).await.unwrap();
assert_eq!(
stats_entities,
full_entities.len() as u64,
"stats() entity total must equal a full list keyset walk under the same identity"
);
assert_eq!(stats_entities, 4);
let full_edges = rt
.list_edges(&vis_tok, EdgeListFilter::default(), 1000, 0)
.await
.unwrap();
let stats_edges = rt
.count_edges(&vis_tok, EdgeListFilter::default())
.await
.unwrap();
assert_eq!(
stats_edges,
full_edges.len() as u64,
"stats() edge total must equal a full list keyset walk under the same identity"
);
assert_eq!(stats_edges, per_namespace_edges);
assert_eq!(stats_edges, 3);
let edges_by_relation = rt.count_edges_by_relation(&vis_tok).await.unwrap();
assert_eq!(edges_by_relation, per_namespace_relations);
let relation_sum: u64 = edges_by_relation.values().sum();
assert_eq!(
relation_sum, stats_edges,
"edges_by_relation must sum to the edges scalar under all identities"
);
let full_notes = rt.list_notes(&vis_tok, None, 200, 0).await.unwrap();
let stats_notes = rt.count_notes(&vis_tok, None).await.unwrap();
assert_eq!(
stats_notes,
full_notes.len() as u64,
"stats() note total must equal a full list keyset walk under the same identity"
);
assert_eq!(stats_notes, per_namespace_notes);
assert_eq!(stats_notes, 2);
}
#[test]
#[serial_test::serial]
fn test_harness_refuses_runtime_default_store() {
if test_process::run_in_child() {
return;
}
struct HomeGuard(Option<std::ffi::OsString>);
impl Drop for HomeGuard {
fn drop(&mut self) {
match self.0.take() {
Some(home) => std::env::set_var("HOME", home),
None => std::env::remove_var("HOME"),
}
}
}
assert_eq!(
std::env::var("KHIVE_TEST_HARNESS").as_deref(),
Ok("1"),
"the workspace Cargo harness must mark integration-test processes"
);
let fake_home = tempfile::tempdir().expect("temporary HOME");
let _home_guard = HomeGuard(std::env::var_os("HOME"));
std::env::set_var("HOME", fake_home.path());
let expected_path = fake_home.path().join(".khive/khive.db");
let config = RuntimeConfig::no_embeddings();
assert_eq!(config.db_path.as_deref(), Some(expected_path.as_path()));
let error = match KhiveRuntime::new(config) {
Ok(_) => panic!("test harness opened the default home-directory store"),
Err(error) => error,
};
assert!(
error.to_string().contains("test harness refused"),
"unexpected guard error: {error}"
);
assert!(
!expected_path.exists(),
"the guarded runtime must refuse the path before SQLite creates it"
);
}