use std::collections::HashMap;
use std::str::FromStr;
use chrono::Utc;
use serde::Serialize;
use uuid::Uuid;
use khive_score::DeterministicScore;
use khive_storage::entity::EntityTypeCounts;
use khive_storage::graph::{CommitAnnotationGuard, CommitAnnotationInsertOutcome};
use khive_storage::note::Note;
use khive_storage::types::{
DeleteMode, DirectedNeighborHit, Direction, EdgeSortField, EdgeUpsertDisposition,
EdgeUpsertRefusal, EdgeUpsertRequest, EdgeUpsertResult, GraphPath, GuardedEdgeUpsertOutcome,
LinkId, NeighborCursor, NeighborHit, NeighborQuery, Page, PageRequest, SeekCursor, SortOrder,
SqlRow, SqlStatement, SqlValue, TextFilter, TextQueryMode, TextSearchRequest, TraversalRequest,
};
use khive_storage::{
Attachment, AttachmentSubstrate, Edge, EdgeRelation, Entity, EntityFilter, Event, EventFilter,
NewAttachment,
};
use khive_types::{EdgeEndpointRule, EndpointKind, EventKind, KhiveError, SubstrateKind};
use khive_db::stores::entity::{entity_hard_delete_statement, entity_upsert_statement};
use khive_db::stores::event::hard_delete_lineage_warning_statements;
use khive_db::stores::graph::{
compose_graph_mutation_events, edge_hard_delete_statement, purge_incident_edges_statement,
GraphMutationOutcome, GraphMutationPreconditions, GraphMutationRequest,
};
use khive_db::stores::note::note_hard_delete_statement;
use khive_db::stores::text::insert_document_statements;
use khive_db::{pool::RuntimeWriteOperation, SqliteError};
use rusqlite::OptionalExtension;
#[cfg(test)]
mod batch_edge_tests;
struct EdgeReadWindow {
outcomes: Vec<Option<RuntimeResult<Option<Edge>>>>,
groups: Vec<(khive_types::Namespace, Vec<usize>)>,
}
fn restore_reindex_failed(kind: &str, id: Uuid, error: RuntimeError) -> RuntimeError {
RuntimeError::Internal(format!(
"{kind} {id} is restored and text-indexed, but its embedding rebuild failed \
and will be retried by the next reindex: {error}"
))
}
fn merge_tombstone_restore_refused(id: Uuid, kept_id: impl std::fmt::Display) -> RuntimeError {
KhiveError::conflict(format!(
"merge_tombstone: {id} was merged into {kept_id}; a merge tombstone is not restorable, query the kept id"
))
.with_details(khive_types::Details::new_owned([
("reason", "merge_tombstone".into()),
("merged_into", kept_id.to_string()),
]))
.into()
}
#[derive(Debug, PartialEq, Eq)]
pub struct EntityStatsCounts {
pub entities: u64,
pub entities_by_type: Option<EntityTypeCounts>,
}
pub async fn entity_stats_counts(
store: &dyn khive_storage::EntityStore,
token: &NamespaceToken,
) -> RuntimeResult<EntityStatsCounts> {
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
match store.count_entities_by_type(&namespaces).await? {
Some(groups) => {
let entities = groups.iter().try_fold(0_u64, |total, (_, count)| {
total.checked_add(*count).ok_or_else(|| {
RuntimeError::Internal(
"entity type counts exceed the scalar count range".into(),
)
})
})?;
Ok(EntityStatsCounts {
entities,
entities_by_type: Some(groups),
})
}
None => {
let entities = store
.count_entities(
token.namespace().as_str(),
EntityFilter {
namespaces,
..EntityFilter::default()
},
)
.await?;
Ok(EntityStatsCounts {
entities,
entities_by_type: None,
})
}
}
}
pub struct EntityClaimSpec {
pub id: Uuid,
pub kind: String,
pub entity_type: Option<String>,
pub name: String,
pub description: Option<String>,
pub properties: Option<serde_json::Value>,
pub tags: Vec<String>,
pub identity_tag: String,
}
fn live_merged_entity_refused(id: Uuid, kept_id: impl std::fmt::Display) -> RuntimeError {
KhiveError::conflict(format!(
"live_merged_entity: {id} is live but still carries merged_into {kept_id}; a row an \
earlier restore left live over its merge is not restorable, re-tombstone it or query the \
kept id"
))
.with_details(khive_types::Details::new_owned([
("reason", "live_merged_entity".into()),
("merged_into", kept_id.to_string()),
]))
.into()
}
fn restore_key_conflict(key: &str, holder: &Note) -> RuntimeError {
KhiveError::conflict(format!(
"restore_key_conflict: key {key:?} is already held by live note {}",
holder.id
))
.with_details(khive_types::Details::new_owned([
("reason", "restore_key_conflict".into()),
("key", key.to_owned()),
("existing_id", holder.id.to_string()),
]))
.into()
}
use crate::atomic_plan::{
AddEntityPlan, AffectedRowGuard, DeletePlan, PlanStatement, PostCommitEffect, UpdatePlan,
};
use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicOpPlan, AtomicRunOutcome};
use crate::curation::{entity_fts_document, note_embedding_text_ref, note_fts_document};
use crate::error::{GuardedWriteFailure, RuntimeError, RuntimeResult};
use crate::runtime::{KhiveRuntime, NamespaceToken};
#[derive(Clone, Debug, Serialize)]
pub struct PostCommitDegradation {
pub stage: &'static str,
pub error: String,
}
impl PostCommitDegradation {
fn new(stage: &'static str, error: impl ToString) -> Self {
Self {
stage,
error: error.to_string(),
}
}
}
macro_rules! conditional_insert_stages {
($($stage:ident => $label:literal),+ $(,)?) => {
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ConditionalInsertStage {
$($stage,)+
}
impl ConditionalInsertStage {
pub const ALL: &'static [Self] = &[$(Self::$stage,)+];
pub const fn label(self) -> &'static str {
match self {
$(Self::$stage => $label,)+
}
}
}
};
}
conditional_insert_stages! {
FtsAcquisition => "fts_acquisition",
FtsUpsert => "fts_upsert",
Embedding => "embedding",
VectorAcquisition => "vector_acquisition",
VectorInsert => "vector_insert",
}
impl ConditionalInsertStage {
pub fn from_label(label: &str) -> Option<Self> {
Self::ALL
.iter()
.copied()
.find(|stage| stage.label() == label)
}
}
fn record_conditional_insert_degradation(
degradations: &mut Vec<PostCommitDegradation>,
id: Uuid,
stage: ConditionalInsertStage,
error: impl ToString,
) {
record_post_commit_degradation(degradations, "try_create_note", id, stage.label(), error);
}
fn record_post_commit_degradation(
degradations: &mut Vec<PostCommitDegradation>,
operation: &'static str,
id: Uuid,
stage: &'static str,
error: impl ToString,
) {
let degradation = PostCommitDegradation::new(stage, error);
tracing::warn!(%operation, %id, stage, error = %degradation.error,
"substrate mutation committed with post-commit degradation");
degradations.push(degradation);
}
fn legacy_post_commit_result<T>(
operation: &'static str,
id: Uuid,
value: T,
degradations: Vec<PostCommitDegradation>,
) -> RuntimeResult<T> {
if degradations.is_empty() {
return Ok(value);
}
let failures = serde_json::Value::Array(
degradations
.iter()
.map(|failure| {
serde_json::json!({"stage": failure.stage, "error": failure.error.as_str()})
})
.collect(),
);
Err(KhiveError::internal(format!(
"{operation} committed record {id}, but post-commit work failed; do not retry the mutation; reconcile by record_id"
))
.with_details(khive_types::Details::new_owned([
("reason", "post_commit_degraded".to_string()),
("operation", operation.to_string()),
("record_id", id.to_string()),
("committed", "true".to_string()),
("retryable", "false".to_string()),
("post_commit_degradations", failures.to_string()),
]))
.into())
}
pub(crate) fn legacy_post_commit_result_with_embedding<T>(
operation: &'static str,
id: Uuid,
value: T,
embedding: crate::retrieval::EmbeddingTruncationReport,
degradations: Vec<PostCommitDegradation>,
) -> RuntimeResult<T> {
if !embedding.any_truncated() {
return legacy_post_commit_result(operation, id, value, degradations);
}
let failures = serde_json::Value::Array(
degradations
.iter()
.map(|failure| {
serde_json::json!({"stage": failure.stage, "error": failure.error.as_str()})
})
.collect(),
);
Err(KhiveError::internal(format!(
"{operation} committed record {id}, but embedding input was truncated; do not retry the mutation; reconcile by record_id"
))
.with_details(khive_types::Details::new_owned([
("reason", "embedding_input_truncated".to_string()),
("operation", operation.to_string()),
("record_id", id.to_string()),
("committed", "true".to_string()),
("retryable", "false".to_string()),
(
"embedding_truncation_report",
serde_json::json!(embedding).to_string(),
),
("post_commit_degradations", failures.to_string()),
]))
.into())
}
#[cfg(test)]
std::thread_local! {
static LINK_FAIL_AFTER: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
}
#[cfg(any(test, feature = "fault-injection"))]
std::thread_local! {
static VECTOR_FAIL_AFTER: std::cell::Cell<Option<usize>> =
const { std::cell::Cell::new(None) };
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_vector_fail_after(n: usize) {
VECTOR_FAIL_AFTER.with(|cell| cell.set(Some(n)));
}
#[cfg(any(test, feature = "fault-injection"))]
type FaultArmSet = std::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<()>>>;
#[cfg(any(test, feature = "fault-injection"))]
const MAX_FAULT_ARMS: usize = 64;
#[cfg(any(test, feature = "fault-injection"))]
static FTS_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
static VECTOR_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
static ENTITY_COMPENSATION_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
static FTS_FAIL_MANY_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
static FTS_FAIL_MANY_PARTIAL_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
static PREFIX_RESOLVE_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
#[cfg(any(test, feature = "fault-injection"))]
#[must_use = "the fault injection is disarmed when this guard is dropped"]
pub struct FaultInjectionArm {
namespace: String,
token: std::sync::Arc<()>,
arms: &'static FaultArmSet,
}
#[cfg(any(test, feature = "fault-injection"))]
impl Drop for FaultInjectionArm {
fn drop(&mut self) {
let mut arms = self.arms.lock().unwrap();
if arms
.get(&self.namespace)
.is_some_and(|token| std::sync::Arc::ptr_eq(token, &self.token))
{
arms.remove(&self.namespace);
}
}
}
#[cfg(any(test, feature = "fault-injection"))]
fn arm_fault(arms: &'static FaultArmSet, namespace: &str, max_arms: usize) -> FaultInjectionArm {
let token = std::sync::Arc::new(());
let refusal = {
let mut active = arms.lock().unwrap();
if active.contains_key(namespace) {
Some("the namespace is already armed")
} else if active.len() >= max_arms {
Some("the arm set is at capacity")
} else {
active.insert(namespace.to_string(), std::sync::Arc::clone(&token));
None
}
};
if let Some(reason) = refusal {
panic!("cannot arm fault injection for namespace `{namespace}`: {reason}");
}
FaultInjectionArm {
namespace: namespace.to_string(),
token,
arms,
}
}
#[cfg(any(test, feature = "fault-injection"))]
fn consume_fault(arms: &FaultArmSet, namespace: &str) -> bool {
arms.lock().unwrap().remove(namespace).is_some()
}
#[cfg(any(test, feature = "fault-injection"))]
static FTS_SEARCH_FAIL_NS: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_fts_fail_scoped(ns: &str) -> FaultInjectionArm {
arm_fault(&FTS_FAIL_NS, ns, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_fts_fail_many_scoped(ns: &str) -> FaultInjectionArm {
arm_fault(&FTS_FAIL_MANY_NS, ns, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_fts_fail_many_partial_scoped(ns: &str) -> FaultInjectionArm {
arm_fault(&FTS_FAIL_MANY_PARTIAL_NS, ns, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_fts_search_fail(ns: &str) {
*FTS_SEARCH_FAIL_NS.lock().unwrap() = Some(ns.to_string());
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_vector_fail_scoped(ns: &str) -> FaultInjectionArm {
arm_fault(&VECTOR_FAIL_NS, ns, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_entity_compensation_fail_scoped(ns: &str) -> FaultInjectionArm {
arm_fault(&ENTITY_COMPENSATION_FAIL_NS, ns, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_prefix_resolve_fail_scoped(prefix: &str) -> FaultInjectionArm {
arm_fault(&PREFIX_RESOLVE_FAIL_NS, prefix, MAX_FAULT_ARMS)
}
#[cfg(any(test, feature = "fault-injection"))]
static ROLLBACK_CLEANUP_FAIL_NS: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
#[cfg(any(test, feature = "fault-injection"))]
pub fn arm_rollback_cleanup_fail(ns: &str) {
*ROLLBACK_CLEANUP_FAIL_NS.lock().unwrap() = Some(ns.to_string());
}
#[cfg(any(test, feature = "fault-injection"))]
pub(crate) fn consume_fts_fail_fault(ns: &str) -> bool {
consume_fault(&FTS_FAIL_NS, ns)
}
#[cfg(any(test, feature = "fault-injection"))]
pub(crate) fn consume_vector_fail_fault(ns: &str) -> bool {
consume_fault(&VECTOR_FAIL_NS, ns)
}
#[derive(Clone, Debug)]
pub struct NoteSearchHit {
pub note_id: Uuid,
pub score: DeterministicScore,
pub rank_score_kind: crate::RankScoreKind,
pub signals: crate::SearchSignals,
pub source: crate::SearchSource,
pub title: Option<String>,
pub snippet: Option<String>,
}
fn salience_weighted_rank(score: DeterministicScore, salience: Option<f64>) -> DeterministicScore {
const SCALE_RAW: i128 = 1_i128 << 32;
let salience = DeterministicScore::from_f64(salience.unwrap_or(0.5));
let weight_raw = SCALE_RAW / 2 + i128::from(salience.to_raw()) / 2;
let weighted_raw = i128::from(score.to_raw()) * weight_raw / SCALE_RAW;
DeterministicScore::from_raw(weighted_raw.clamp(
i128::from(DeterministicScore::NEG_INF.to_raw()),
i128::from(DeterministicScore::MAX.to_raw()),
) as i64)
}
#[derive(Clone, Debug)]
pub struct NoteSearchOutcome {
pub hits: Vec<NoteSearchHit>,
pub vector_error: Option<String>,
}
pub fn hex_prefix_to_uuid_pattern(prefix: &str) -> String {
if prefix.contains('-') {
return prefix.to_string();
}
const BOUNDARIES: [usize; 4] = [8, 13, 18, 23]; let mut out = String::with_capacity(36);
for c in prefix.chars() {
if BOUNDARIES.contains(&out.len()) {
out.push('-');
}
out.push(c);
}
out
}
pub fn uuid_prefix_bounds(prefix: &str) -> Option<(String, String)> {
const HYPHEN_POSITIONS: [usize; 4] = [8, 13, 18, 23];
let compact = if prefix.contains('-') {
if prefix.len() > 36 {
return None;
}
let mut compact = String::with_capacity(32);
for (index, byte) in prefix.bytes().enumerate() {
if HYPHEN_POSITIONS.contains(&index) {
if byte != b'-' {
return None;
}
} else if byte.is_ascii_hexdigit() {
compact.push(char::from(byte.to_ascii_lowercase()));
} else {
return None;
}
}
compact
} else {
if prefix.is_empty()
|| prefix.len() > 32
|| !prefix.bytes().all(|byte| byte.is_ascii_hexdigit())
{
return None;
}
prefix.to_ascii_lowercase()
};
if compact.is_empty() || compact.len() > 32 {
return None;
}
let lower = hex_prefix_to_uuid_pattern(&compact);
let mut successor = compact.into_bytes();
let mut carried_past_start = true;
for index in (0..successor.len()).rev() {
let next = match successor[index] {
b'0'..=b'8' | b'a'..=b'e' => Some(successor[index] + 1),
b'9' => Some(b'a'),
b'f' => None,
_ => return None,
};
if let Some(next) = next {
successor[index] = next;
successor.truncate(index + 1);
carried_past_start = false;
break;
}
}
let upper = if carried_past_start {
"g".to_string()
} else {
let compact_upper = String::from_utf8(successor).ok()?;
hex_prefix_to_uuid_pattern(&compact_upper)
};
Some((lower, upper))
}
fn resolve_prefix_statement(
table: &str,
has_deleted_at: bool,
include_deleted: bool,
namespaces: Option<&[String]>,
lower: &str,
upper: &str,
) -> SqlStatement {
let namespace_clause = namespaces.map(|namespaces| {
let placeholders: Vec<String> = (0..namespaces.len())
.map(|index| format!("?{}", index + 3))
.collect();
format!(" AND namespace IN ({})", placeholders.join(", "))
});
let deleted_filter = if has_deleted_at && !include_deleted {
" AND deleted_at IS NULL"
} else {
""
};
let mut params = vec![
SqlValue::Text(lower.to_owned()),
SqlValue::Text(upper.to_owned()),
];
if let Some(namespaces) = namespaces {
params.extend(
namespaces
.iter()
.map(|namespace| SqlValue::Text(namespace.clone())),
);
}
SqlStatement {
sql: format!(
"SELECT id FROM {table} \
WHERE id >= ?1 AND id < ?2{namespace_clause}{deleted_filter} ORDER BY id LIMIT 2",
namespace_clause = namespace_clause.as_deref().unwrap_or("")
),
params,
label: Some("resolve_prefix".into()),
}
}
fn text_preview(text: &str, max_chars: usize) -> Option<String> {
let trimmed = text.trim();
if trimmed.is_empty() {
None
} else {
Some(trimmed.chars().take(max_chars).collect())
}
}
fn normalize_symmetric_direction(
direction: Direction,
relations: Option<&[EdgeRelation]>,
) -> Direction {
let Some(rels) = relations else {
return direction;
};
if rels.is_empty() {
return direction;
}
let all_symmetric = rels
.iter()
.all(|r| matches!(r, EdgeRelation::CompetesWith | EdgeRelation::ComposedWith));
if all_symmetric {
Direction::Both
} else {
direction
}
}
fn direction_sort_rank(direction: &Direction) -> u8 {
match direction {
Direction::Out => 0,
Direction::In => 1,
Direction::Both => 2,
}
}
fn note_title(note: &Note) -> Option<String> {
note.name
.clone()
.filter(|s| !s.trim().is_empty())
.or_else(|| Some(format!("[{}]", note.kind.as_str())))
}
fn note_snippet(note: &Note) -> Option<String> {
text_preview(¬e.content, 200)
}
#[derive(Clone, Debug)]
pub enum Resolved {
Entity(Entity),
Note(Note),
Event(Event),
PackRecord {
pack: String,
kind: String,
data: serde_json::Value,
},
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum EdgeEndpointKind {
Entity,
Note,
Event,
Edge,
}
impl EdgeEndpointKind {
pub const fn name(self) -> &'static str {
match self {
Self::Entity => "entity",
Self::Note => "note",
Self::Event => "event",
Self::Edge => "edge",
}
}
}
fn resolved_pair(r: Option<&Resolved>) -> Option<(&'static str, &str, Option<&str>)> {
match r? {
Resolved::Entity(e) => Some(("entity", e.kind.as_str(), e.entity_type.as_deref())),
Resolved::Note(n) => Some(("note", n.kind.as_str(), None)),
Resolved::Event(_) => None,
Resolved::PackRecord { .. } => None,
}
}
pub fn endpoint_matches(
spec: &EndpointKind,
substrate: &str,
kind: &str,
entity_type: Option<&str>,
) -> bool {
match spec {
EndpointKind::EntityOfKind(k) => substrate == "entity" && *k == kind,
EndpointKind::NoteOfKind(k) => substrate == "note" && *k == kind,
EndpointKind::EntityOfType {
kind: k,
entity_type: t,
} => substrate == "entity" && *k == kind && entity_type == Some(*t),
}
}
fn pattern_endpoint_matches(
spec: &EndpointKind,
substrate: &str,
kind: &str,
entity_type: Option<&str>,
) -> bool {
match spec {
EndpointKind::EntityOfType {
kind: k,
entity_type: t,
} => substrate == "entity" && *k == kind && entity_type.is_none_or(|et| et == *t),
_ => endpoint_matches(spec, substrate, kind, entity_type),
}
}
pub fn accepted_pack_relations_for_entities(
rules: &[EdgeEndpointRule],
src_kind: &str,
src_entity_type: Option<&str>,
tgt_kind: &str,
tgt_entity_type: Option<&str>,
) -> Vec<EdgeRelation> {
let mut relations: Vec<EdgeRelation> = rules
.iter()
.filter(|r| {
endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
&& endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
})
.map(|r| r.relation)
.collect();
relations.sort_by_key(|r| r.as_str());
relations.dedup();
relations
}
pub fn accepted_entity_relations_for_entities(
rules: &[EdgeEndpointRule],
src_kind: &str,
src_entity_type: Option<&str>,
tgt_kind: &str,
tgt_entity_type: Option<&str>,
) -> Vec<EdgeRelation> {
let mut relations: Vec<EdgeRelation> = BASE_ENTITY_ENDPOINT_RULES
.iter()
.filter(|(src, _relation, tgt)| (*src == "*" || *src == src_kind) && *tgt == tgt_kind)
.map(|(_src, relation, _tgt)| *relation)
.collect();
relations.extend(
accepted_pack_relations_for_entities(
rules,
src_kind,
src_entity_type,
tgt_kind,
tgt_entity_type,
)
.into_iter()
.filter(|relation| {
*relation != EdgeRelation::Annotates && !crate::pack::is_special_relation(*relation)
}),
);
relations.sort_by_key(|relation| relation.as_str());
relations.dedup();
relations
}
fn accepted_entity_relations_description(
rules: &[EdgeEndpointRule],
src_kind: &str,
src_entity_type: Option<&str>,
tgt_kind: &str,
tgt_entity_type: Option<&str>,
) -> String {
let relations = accepted_entity_relations_for_entities(
rules,
src_kind,
src_entity_type,
tgt_kind,
tgt_entity_type,
);
if relations.is_empty() {
"none".to_string()
} else {
relations
.iter()
.map(EdgeRelation::as_str)
.collect::<Vec<_>>()
.join(", ")
}
}
fn accepted_pack_relations_for_pattern_entities(
rules: &[EdgeEndpointRule],
src_kind: &str,
src_entity_type: Option<&str>,
tgt_kind: &str,
tgt_entity_type: Option<&str>,
) -> Vec<EdgeRelation> {
let mut relations: Vec<EdgeRelation> = rules
.iter()
.filter(|r| {
pattern_endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
&& pattern_endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
})
.map(|r| r.relation)
.collect();
relations.sort_by_key(|r| r.as_str());
relations.dedup();
relations
}
fn accepted_entity_kind_pairs_for_relation(
pack_rules: &[EdgeEndpointRule],
relation: EdgeRelation,
) -> Vec<(&'static str, &'static str)> {
let mut pairs = Vec::new();
for src in khive_types::EntityKind::ALL {
for tgt in khive_types::EntityKind::ALL {
let allowed = base_entity_rule_allows(src.name(), relation, tgt.name())
|| (!crate::pack::is_special_relation(relation)
&& accepted_pack_relations_for_pattern_entities(
pack_rules,
src.name(),
None,
tgt.name(),
None,
)
.contains(&relation));
if allowed {
pairs.push((src.name(), tgt.name()));
}
}
}
pairs
}
fn static_impossible_edge_pattern_warnings(
language: khive_query::QueryLanguage,
pattern: &khive_query::ast::MatchPattern,
pack_rules: &[EdgeEndpointRule],
) -> Vec<String> {
use khive_query::ast::{EdgeDirection, PatternElement};
if language != khive_query::QueryLanguage::Gql {
return Vec::new();
}
let elements = &pattern.elements;
let mut warnings = Vec::new();
for (i, el) in elements.iter().enumerate() {
let PatternElement::Edge(edge) = el else {
continue;
};
if edge.relations.len() != 1 || edge.min_hops != 1 || edge.max_hops != 1 {
continue;
}
let (left, right) = match (elements.get(i.wrapping_sub(1)), elements.get(i + 1)) {
(Some(PatternElement::Node(l)), Some(PatternElement::Node(r))) => (l, r),
_ => continue,
};
let (src_node, tgt_node) = match edge.direction {
EdgeDirection::Out => (left, right),
EdgeDirection::In => (right, left),
EdgeDirection::Both => continue,
};
let (Some(src_raw), Some(tgt_raw)) = (src_node.kind.as_deref(), tgt_node.kind.as_deref())
else {
continue;
};
let (Ok(src_kind), Ok(tgt_kind)) = (
src_raw.parse::<khive_types::EntityKind>(),
tgt_raw.parse::<khive_types::EntityKind>(),
) else {
continue;
};
let Ok(relation) = edge.relations[0].parse::<EdgeRelation>() else {
continue;
};
let possible = base_entity_rule_allows(src_kind.name(), relation, tgt_kind.name())
|| (!crate::pack::is_special_relation(relation)
&& accepted_pack_relations_for_pattern_entities(
pack_rules,
src_kind.name(),
src_node.entity_type.as_deref(),
tgt_kind.name(),
tgt_node.entity_type.as_deref(),
)
.contains(&relation));
if possible {
continue;
}
let accepted = accepted_entity_kind_pairs_for_relation(pack_rules, relation);
let accepted_str = if accepted.is_empty() {
"none".to_string()
} else {
accepted
.iter()
.map(|(s, t)| format!("{s}->{t}"))
.collect::<Vec<_>>()
.join(", ")
};
warnings.push(format!(
"pattern ({src})-[:{relation}]->({tgt}) can never match: '{relation}' does not accept \
{src}->{tgt} endpoints; accepted source->target kinds for '{relation}': {accepted_str}",
src = src_kind.name(),
tgt = tgt_kind.name(),
));
}
warnings
}
fn pack_rule_allows(
rules: &[EdgeEndpointRule],
relation: EdgeRelation,
src: Option<&Resolved>,
tgt: Option<&Resolved>,
) -> bool {
let Some((src_sub, src_kind, src_type)) = resolved_pair(src) else {
return false;
};
let Some((tgt_sub, tgt_kind, tgt_type)) = resolved_pair(tgt) else {
return false;
};
rules.iter().any(|r| {
r.relation == relation
&& endpoint_matches(&r.source, src_sub, src_kind, src_type)
&& endpoint_matches(&r.target, tgt_sub, tgt_kind, tgt_type)
})
}
pub const BASE_ENTITY_ENDPOINT_RULES: &[(&str, EdgeRelation, &str)] = &[
("concept", EdgeRelation::Contains, "concept"),
("project", EdgeRelation::Contains, "project"),
("project", EdgeRelation::Contains, "artifact"),
("org", EdgeRelation::Contains, "project"),
("org", EdgeRelation::Contains, "service"),
("concept", EdgeRelation::PartOf, "concept"),
("project", EdgeRelation::PartOf, "project"),
("project", EdgeRelation::PartOf, "org"),
("*", EdgeRelation::InstanceOf, "concept"),
("service", EdgeRelation::InstanceOf, "project"),
("document", EdgeRelation::LinksTo, "document"),
("concept", EdgeRelation::LocatedIn, "concept"),
("org", EdgeRelation::LocatedIn, "concept"),
("person", EdgeRelation::Owns, "org"),
("org", EdgeRelation::Owns, "org"),
("concept", EdgeRelation::Extends, "concept"),
("concept", EdgeRelation::VariantOf, "concept"),
("artifact", EdgeRelation::VariantOf, "artifact"),
("concept", EdgeRelation::IntroducedBy, "document"),
("concept", EdgeRelation::IntroducedBy, "person"),
("artifact", EdgeRelation::IntroducedBy, "document"),
("project", EdgeRelation::IntroducedBy, "document"),
("service", EdgeRelation::IntroducedBy, "document"),
("document", EdgeRelation::IntroducedBy, "person"),
("document", EdgeRelation::IntroducedBy, "org"),
("concept", EdgeRelation::IntroducedBy, "org"),
("artifact", EdgeRelation::DerivedFrom, "dataset"),
("artifact", EdgeRelation::DerivedFrom, "document"),
("artifact", EdgeRelation::DerivedFrom, "project"),
("artifact", EdgeRelation::DerivedFrom, "artifact"),
("document", EdgeRelation::DerivedFrom, "document"),
("document", EdgeRelation::Precedes, "document"),
("dataset", EdgeRelation::Precedes, "dataset"),
("artifact", EdgeRelation::Precedes, "artifact"),
("service", EdgeRelation::Precedes, "service"),
("project", EdgeRelation::Precedes, "project"),
("project", EdgeRelation::DependsOn, "project"),
("service", EdgeRelation::DependsOn, "project"),
("service", EdgeRelation::DependsOn, "service"),
("service", EdgeRelation::DependsOn, "artifact"),
("service", EdgeRelation::DependsOn, "dataset"),
("artifact", EdgeRelation::DependsOn, "project"),
("artifact", EdgeRelation::DependsOn, "service"),
("document", EdgeRelation::DependsOn, "document"),
("concept", EdgeRelation::Enables, "concept"),
("service", EdgeRelation::Enables, "concept"),
("dataset", EdgeRelation::Enables, "concept"),
("project", EdgeRelation::Implements, "concept"),
("service", EdgeRelation::Implements, "concept"),
("concept", EdgeRelation::CompetesWith, "concept"),
("project", EdgeRelation::CompetesWith, "project"),
("service", EdgeRelation::CompetesWith, "service"),
("org", EdgeRelation::CompetesWith, "org"),
("concept", EdgeRelation::ComposedWith, "concept"),
("project", EdgeRelation::ComposedWith, "project"),
("concept", EdgeRelation::Supersedes, "concept"),
("document", EdgeRelation::Supersedes, "document"),
("artifact", EdgeRelation::Supersedes, "artifact"),
("service", EdgeRelation::Supersedes, "service"),
("dataset", EdgeRelation::Supersedes, "dataset"),
("concept", EdgeRelation::Supports, "concept"),
("document", EdgeRelation::Supports, "concept"),
("dataset", EdgeRelation::Supports, "concept"),
("artifact", EdgeRelation::Supports, "concept"),
("concept", EdgeRelation::Refutes, "concept"),
("document", EdgeRelation::Refutes, "concept"),
("dataset", EdgeRelation::Refutes, "concept"),
("artifact", EdgeRelation::Refutes, "concept"),
];
pub fn base_entity_endpoint_rules() -> &'static [(&'static str, EdgeRelation, &'static str)] {
BASE_ENTITY_ENDPOINT_RULES
}
pub fn base_entity_rule_allows(src_kind: &str, relation: EdgeRelation, tgt_kind: &str) -> bool {
BASE_ENTITY_ENDPOINT_RULES.iter().any(|(src, rel, tgt)| {
*rel == relation && (*src == "*" || *src == src_kind) && *tgt == tgt_kind
})
}
pub(crate) fn canonical_edge_endpoints(
relation: EdgeRelation,
source_id: Uuid,
target_id: Uuid,
) -> (Uuid, Uuid) {
relation.canonical_endpoints(source_id, target_id)
}
pub(crate) fn canonical_edge_endpoint_kinds(
requested_source_id: Uuid,
canonical_source_id: Uuid,
source_kind: EdgeEndpointKind,
target_kind: EdgeEndpointKind,
) -> (EdgeEndpointKind, EdgeEndpointKind) {
if requested_source_id == canonical_source_id {
(source_kind, target_kind)
} else {
(target_kind, source_kind)
}
}
pub(crate) fn infer_dependency_kind(src_kind: &str, tgt_kind: &str) -> Option<&'static str> {
match (src_kind, tgt_kind) {
("project", "project") => Some("build"),
("service", "service") => Some("runtime"),
("service", "dataset") => Some("data"),
("service", "artifact") => Some("artifact"),
("artifact", "project") | ("artifact", "service") => Some("tooling"),
("document", "document") => Some("normative"),
_ => None,
}
}
pub(crate) fn merge_dependency_kind(
src_kind: &str,
tgt_kind: &str,
metadata: Option<serde_json::Value>,
) -> Option<serde_json::Value> {
let metadata = metadata.filter(|value| !value.is_null());
if let Some(ref m) = metadata {
if m.get("dependency_kind").is_some() {
return metadata;
}
}
let Some(inferred) = infer_dependency_kind(src_kind, tgt_kind) else {
return metadata;
};
let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
if let Some(o) = obj.as_object_mut() {
o.insert("dependency_kind".to_string(), serde_json::json!(inferred));
}
Some(obj)
}
pub fn merge_entry_metadata(
metadata: Option<serde_json::Value>,
dependency_kind: Option<String>,
) -> RuntimeResult<Option<serde_json::Value>> {
validate_metadata_shape(metadata.as_ref())?;
let metadata = metadata.filter(|value| !value.is_null());
let Some(dk) = dependency_kind else {
return Ok(metadata);
};
let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
let map = obj
.as_object_mut()
.ok_or_else(|| RuntimeError::InvalidInput("metadata must be a JSON object".into()))?;
map.entry("dependency_kind".to_string())
.or_insert_with(|| serde_json::json!(dk));
Ok(Some(obj))
}
const VALID_DEPENDENCY_KINDS: &[&str] = &[
"build",
"runtime",
"data",
"artifact",
"tooling",
"normative",
];
pub(crate) fn validate_edge_weight(weight: f64) -> RuntimeResult<()> {
if !weight.is_finite() || !(0.0..=1.0).contains(&weight) {
return Err(RuntimeError::InvalidInput(format!(
"edge weight must be finite and in [0.0, 1.0], got {weight}"
)));
}
Ok(())
}
fn validate_metadata_shape(metadata: Option<&serde_json::Value>) -> RuntimeResult<()> {
if metadata.is_some_and(|value| !value.is_null() && !value.is_object()) {
return Err(RuntimeError::InvalidInput(
"metadata must be a JSON object".into(),
));
}
Ok(())
}
pub(crate) fn validate_edge_metadata(
relation: EdgeRelation,
metadata: Option<&serde_json::Value>,
) -> RuntimeResult<()> {
validate_metadata_shape(metadata)?;
let Some(meta) = metadata.filter(|value| !value.is_null()) else {
return Ok(());
};
let object = meta.as_object().expect("validated metadata object");
if object
.get("optional")
.is_some_and(|value| !value.is_boolean())
{
return Err(RuntimeError::InvalidInput(
"metadata.optional must be a boolean".into(),
));
}
if let Some(dk) = meta.get("dependency_kind") {
if relation != EdgeRelation::DependsOn {
return Err(RuntimeError::InvalidInput(format!(
"dependency_kind is only valid on depends_on edges (got {})",
relation.as_str()
)));
}
let dk_str = dk
.as_str()
.ok_or_else(|| RuntimeError::InvalidInput("dependency_kind must be a string".into()))?;
if !VALID_DEPENDENCY_KINDS.contains(&dk_str) {
return Err(RuntimeError::InvalidInput(format!(
"unknown dependency_kind {dk_str:?}; valid: {}",
VALID_DEPENDENCY_KINDS.join(" | ")
)));
}
}
Ok(())
}
fn note_props_match(note_props: Option<&serde_json::Value>, filter: &serde_json::Value) -> bool {
let required = match filter.as_object() {
Some(obj) if !obj.is_empty() => obj,
_ => return true,
};
let actual = match note_props.and_then(serde_json::Value::as_object) {
Some(obj) => obj,
None => return false,
};
required
.iter()
.all(|(k, v)| actual.get(k).is_some_and(|av| av == v))
}
fn note_graph_name(note: &Note) -> String {
note.name
.as_deref()
.filter(|name| !name.trim().is_empty())
.map(str::to_owned)
.unwrap_or_else(|| format!("[{}]", note.kind))
}
fn merge_traversal_paths_by_root(paths: Vec<GraphPath>, limit: Option<u32>) -> Vec<GraphPath> {
let mut order: Vec<Uuid> = Vec::new();
let mut merged: HashMap<Uuid, GraphPath> = HashMap::new();
let mut node_index: HashMap<Uuid, HashMap<Uuid, usize>> = HashMap::new();
for path in paths {
let existing = merged.entry(path.root_id).or_insert_with(|| {
order.push(path.root_id);
GraphPath {
root_id: path.root_id,
nodes: Vec::new(),
total_weight: 0.0,
}
});
let index = node_index.entry(path.root_id).or_default();
for node in path.nodes {
match index.get(&node.node_id) {
Some(&i) => {
if node.depth < existing.nodes[i].depth {
existing.nodes[i] = node;
}
}
None => {
index.insert(node.node_id, existing.nodes.len());
existing.nodes.push(node);
}
}
}
}
order
.into_iter()
.filter_map(|root_id| merged.remove(&root_id))
.map(|mut path| {
path.nodes.sort_by_key(|n| n.depth);
if let Some(lim) = limit {
let lim = lim as usize;
let mut non_root_kept = 0usize;
path.nodes.retain(|n| {
if n.depth == 0 {
return true;
}
if non_root_kept < lim {
non_root_kept += 1;
true
} else {
false
}
});
}
recompute_total_weight(&mut path);
path
})
.collect()
}
fn recompute_total_weight(path: &mut GraphPath) {
path.total_weight = path.nodes.iter().map(|n| n.weight).fold(0.0_f64, f64::max);
}
async fn drain_embed_join_set<T: Send + 'static>(
mut join_set: tokio::task::JoinSet<(usize, RuntimeResult<T>)>,
model_count: usize,
) -> RuntimeResult<Vec<T>> {
let mut vectors: Vec<Option<T>> = (0..model_count).map(|_| None).collect();
while let Some(joined) = join_set.join_next().await {
match joined {
Ok((idx, Ok(vector))) => vectors[idx] = Some(vector),
Ok((_idx, Err(e))) => {
join_set.abort_all();
return Err(e);
}
Err(join_err) => {
join_set.abort_all();
return Err(RuntimeError::Internal(format!(
"embed task panicked: {join_err}"
)));
}
}
}
Ok(vectors
.into_iter()
.map(|v| v.expect("every model index observed exactly once by join_set drain"))
.collect())
}
impl KhiveRuntime {
async fn compensate_entity_create(
&self,
token: &NamespaceToken,
entity_id: Uuid,
namespace: &str,
vector_models: &[String],
) -> Vec<String> {
let mut cleanup_errors = Vec::new();
#[cfg(any(test, feature = "fault-injection"))]
let entity_delete_injected = consume_fault(&ENTITY_COMPENSATION_FAIL_NS, namespace);
#[cfg(not(any(test, feature = "fault-injection")))]
let entity_delete_injected = false;
if entity_delete_injected {
cleanup_errors.push("entity row delete: injected compensation failure".to_string());
} else {
match self.entities(token) {
Ok(store) => {
if let Err(error) = store.delete_entity(entity_id, DeleteMode::Hard).await {
cleanup_errors.push(format!("entity row delete: {error}"));
}
}
Err(error) => cleanup_errors.push(format!("entity store access: {error}")),
}
}
match self.text(token) {
Ok(fts) => {
if let Err(error) = fts.delete_document(namespace, entity_id).await {
cleanup_errors.push(format!("FTS document delete: {error}"));
}
}
Err(error) => cleanup_errors.push(format!("FTS store access: {error}")),
}
for model_name in vector_models {
match self.vectors_for_model(token, model_name) {
Ok(vectors) => {
if let Err(error) = vectors.delete(entity_id).await {
cleanup_errors
.push(format!("vector delete for model {model_name}: {error}"));
}
}
Err(error) => cleanup_errors.push(format!(
"vector store access for model {model_name}: {error}"
)),
}
}
cleanup_errors
}
fn entity_create_failure(
entity_id: Uuid,
primary: RuntimeError,
cleanup_errors: Vec<String>,
) -> RuntimeError {
if cleanup_errors.is_empty() {
primary
} else {
RuntimeError::Khive(KhiveError::internal(format!(
"create_entity indexing failed for record {entity_id}; primary failure: \
{primary}; compensation failure(s): {}; partial persistence is possible; \
inspect and reconcile this record before retrying",
cleanup_errors.join("; ")
)))
}
}
pub async fn claim_entity_if_absent(
&self,
token: &NamespaceToken,
spec: EntityClaimSpec,
) -> RuntimeResult<(Entity, bool)> {
self.validate_entity_kind(&spec.kind)?;
let entity_type =
self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
crate::secret_gate::check_at(&spec.name, "entity", "name")?;
if let Some(description) = &spec.description {
crate::secret_gate::check_at(description, "entity", "description")?;
}
if let Some(properties) = &spec.properties {
crate::secret_gate::check_json_at(properties, "entity", "properties")?;
}
crate::secret_gate::check_tags_at(&spec.tags, "entity", "tags")?;
let mut proposed = Entity::new(token.namespace().as_str(), &spec.kind, &spec.name);
proposed.id = spec.id;
proposed.entity_type = entity_type.clone();
proposed.description = spec.description;
proposed.properties = spec.properties;
proposed.tags = spec.tags;
let store = self.entities(token)?;
let inserted = store.insert_entity_if_absent(proposed.clone()).await?;
let entity = if inserted {
proposed
} else {
store
.get_entity_including_deleted(spec.id)
.await?
.ok_or_else(|| {
RuntimeError::Internal(format!(
"entity claim {} lost but the winning row is missing",
spec.id
))
})?
};
if entity.deleted_at.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"entity claim {} is soft-deleted; restore it explicitly",
entity.id
)));
}
if entity.namespace != token.namespace().as_str()
|| entity.kind != spec.kind
|| entity.entity_type.as_deref() != entity_type.as_deref()
|| !entity.name.eq_ignore_ascii_case(&spec.name)
|| !entity
.tags
.iter()
.any(|tag| tag.eq_ignore_ascii_case(&spec.identity_tag))
{
return Err(RuntimeError::InvalidInput(format!(
"entity claim {} belongs to a different record",
entity.id
)));
}
self.ensure_claimed_entity_create_event(token, &entity)
.await?;
self.reindex_claimed_entity(token, &entity).await?;
Ok((entity, inserted))
}
pub async fn ensure_claimed_entity_create_event(
&self,
token: &NamespaceToken,
entity: &Entity,
) -> RuntimeResult<()> {
if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"entity {} is not a live row in the write namespace",
entity.id
)));
}
let events = self.events(token).map_err(|error| {
RuntimeError::Internal(format!(
"entity {} persists but its create event store is unavailable: {error}",
entity.id
))
})?;
let filter = EventFilter {
target_id: Some(entity.id),
kinds: vec![EventKind::EntityCreated],
verbs: vec!["create".into()],
substrates: vec![SubstrateKind::Entity],
after: Some(entity.created_at.saturating_sub(1)),
..EventFilter::default()
};
let page = PageRequest {
offset: 0,
limit: 1,
};
if !events
.query_events(filter.clone(), page.clone())
.await?
.items
.is_empty()
{
return Ok(());
}
let mut event = Event::new(
entity.namespace.clone(),
"create",
EventKind::EntityCreated,
SubstrateKind::Entity,
"",
)
.with_target(entity.id)
.with_payload(serde_json::json!({
"id": entity.id,
"namespace": &entity.namespace,
"kind": &entity.kind,
}));
let event_seed = Uuid::new_v5(&Uuid::NAMESPACE_URL, b"khive:claimed-entity-create:v1");
let mut event_key = Vec::with_capacity(24);
event_key.extend_from_slice(entity.id.as_bytes());
event_key.extend_from_slice(&entity.created_at.to_be_bytes());
event.id = Uuid::new_v5(&event_seed, &event_key);
if let Err(error) = events.append_event(event).await {
if events.query_events(filter, page).await?.items.is_empty() {
return Err(RuntimeError::Internal(format!(
"entity {} persists but its create event failed: {error}",
entity.id
)));
}
}
Ok(())
}
pub async fn reindex_claimed_entity(
&self,
token: &NamespaceToken,
entity: &Entity,
) -> RuntimeResult<()> {
if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"entity {} is not a live row in the write namespace",
entity.id
)));
}
let doc = entity_fts_document(entity);
let embed_body = doc.body.clone();
#[cfg(any(test, feature = "fault-injection"))]
let fts_inject = consume_fault(&FTS_FAIL_NS, &entity.namespace);
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_inject = false;
let fts_result = if fts_inject {
Err(RuntimeError::Internal("injected FTS failure".into()))
} else {
match self.text(token) {
Ok(text) => text.upsert_document(doc).await.map_err(Into::into),
Err(error) => Err(error),
}
};
fts_result.map_err(|error| {
RuntimeError::Internal(format!(
"entity {} persists but its text index failed: {error}",
entity.id
))
})?;
for model_name in self.registered_embedding_model_names() {
let outcome = self
.embed_document_with_model_outcome_for_token(token, &model_name, &embed_body)
.await
.map_err(|error| {
RuntimeError::Internal(format!(
"entity {} persists but model {model_name} embedding failed: {error}",
entity.id
))
})?;
#[cfg(any(test, feature = "fault-injection"))]
let vector_inject = consume_fault(&VECTOR_FAIL_NS, &entity.namespace);
#[cfg(not(any(test, feature = "fault-injection")))]
let vector_inject = false;
if vector_inject {
return Err(RuntimeError::Internal(format!(
"entity {} persists but model {model_name} vector indexing failed: injected vector failure",
entity.id
)));
}
self.vectors_for_model(token, &model_name)
.map_err(|error| {
RuntimeError::Internal(format!(
"entity {} persists but model {model_name} vector store is unavailable: {error}",
entity.id
))
})?
.insert(
entity.id,
SubstrateKind::Entity,
&entity.namespace,
"entity.body",
vec![outcome.vector],
)
.await
.map_err(|error| {
RuntimeError::Internal(format!(
"entity {} persists but model {model_name} vector indexing failed: {error}",
entity.id
))
})?;
}
Ok(())
}
#[allow(clippy::too_many_arguments)]
#[cfg(test)]
pub(crate) async fn create_entity(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
) -> RuntimeResult<Entity> {
let (entity, _, degradations) = self
.create_entity_with_embedding_report_inner(
token,
kind,
entity_type,
name,
description,
properties,
tags,
Vec::new(),
)
.await?;
legacy_post_commit_result("create_entity", entity.id, entity, degradations)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_entity_with_attachments(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
attachments: Vec<NewAttachment>,
) -> RuntimeResult<Entity> {
let (entity, embedding, degradations) = self
.create_entity_with_attachments_inner(
token,
kind,
entity_type,
name,
description,
properties,
tags,
attachments,
)
.await?;
legacy_post_commit_result_with_embedding(
"create_entity_with_attachments",
entity.id,
entity,
embedding,
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_entity_with_attachments_and_report(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
attachments: Vec<NewAttachment>,
) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
let (entity, embedding, degradations) = self
.create_entity_with_attachments_inner(
token,
kind,
entity_type,
name,
description,
properties,
tags,
attachments,
)
.await?;
legacy_post_commit_result(
"create_entity_with_attachments_and_report",
entity.id,
(entity, embedding),
degradations,
)
}
#[allow(clippy::too_many_arguments)]
async fn create_entity_with_attachments_inner(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
attachments: Vec<NewAttachment>,
) -> RuntimeResult<(
Entity,
crate::retrieval::EmbeddingTruncationReport,
Vec<PostCommitDegradation>,
)> {
drop(self.attachments()?);
let blob_store = self.blob_store().ok_or_else(|| {
RuntimeError::Unconfigured(
"create_entity_with_attachments requires an installed BlobStore".to_string(),
)
})?;
let mut roles = std::collections::HashSet::with_capacity(attachments.len());
for attachment in &attachments {
attachment.validate()?;
if !roles.insert(attachment.role.as_str()) {
return Err(RuntimeError::InvalidInput(format!(
"duplicate attachment role {:?}",
attachment.role
)));
}
}
for attachment in &attachments {
if !blob_store.exists(&attachment.content_ref).await? {
return Err(RuntimeError::InvalidInput(format!(
"create_entity_with_attachments requires a published blob; no object exists for {}",
attachment.content_ref
)));
}
}
let validated_type = self.validate_entity_type_for_kind(kind, entity_type)?;
let (entity, embedding, degradations) = self
.create_entity_with_embedding_report_inner(
token,
kind,
validated_type.as_deref(),
name,
description,
properties,
tags,
attachments,
)
.await?;
Ok((entity, embedding, degradations))
}
#[allow(clippy::too_many_arguments)]
pub async fn create_entity_with_embedding_report(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
let (entity, embedding, degradations) = self
.create_entity_with_embedding_report_inner(
token,
kind,
entity_type,
name,
description,
properties,
tags,
Vec::new(),
)
.await?;
legacy_post_commit_result(
"create_entity_with_embedding_report",
entity.id,
(entity, embedding),
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_entity_with_post_commit_report(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
) -> RuntimeResult<(
Entity,
crate::retrieval::EmbeddingTruncationReport,
Vec<PostCommitDegradation>,
)> {
self.create_entity_with_embedding_report_inner(
token,
kind,
entity_type,
name,
description,
properties,
tags,
Vec::new(),
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn create_entity_with_embedding_report_inner(
&self,
token: &NamespaceToken,
kind: &str,
entity_type: Option<&str>,
name: &str,
description: Option<&str>,
properties: Option<serde_json::Value>,
tags: Vec<String>,
attachments: Vec<NewAttachment>,
) -> RuntimeResult<(
Entity,
crate::retrieval::EmbeddingTruncationReport,
Vec<PostCommitDegradation>,
)> {
self.validate_entity_kind(kind)?;
crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
crate::secret_gate::check_at(name, "entity", "name")?;
if let Some(d) = description {
crate::secret_gate::check_at(d, "entity", "description")?;
}
if let Some(ref p) = properties {
crate::secret_gate::check_json_at(p, "entity", "properties")?;
}
crate::secret_gate::check_tags_at(&tags, "entity", "tags")?;
let ns = token.namespace().as_str();
let mut entity = Entity::new(ns, kind, name).with_entity_type(entity_type);
if let Some(d) = description {
entity = entity.with_description(d);
}
if let Some(p) = properties {
entity = entity.with_properties(p);
}
if !tags.is_empty() {
entity = entity.with_tags(tags);
}
let projected_content_ref = attachments
.iter()
.find(|attachment| attachment.role == "content")
.map(|attachment| attachment.content_ref.to_string());
let attachment_rows = attachments
.into_iter()
.map(|attachment| {
Attachment::from_new(
entity.id,
AttachmentSubstrate::Entity,
attachment,
entity.created_at,
)
})
.collect();
self.entities(token)?
.upsert_entity_with_attachments(entity.clone(), attachment_rows)
.await?;
entity.content_ref = projected_content_ref;
let doc = entity_fts_document(&entity);
let embed_body = doc.body.clone();
{
#[cfg(any(test, feature = "fault-injection"))]
let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_inject = false;
let fts_result: RuntimeResult<()> = if fts_inject {
Err(RuntimeError::Internal("injected FTS failure".to_string()))
} else {
match self.text(token) {
Ok(fts) => fts.upsert_document(doc).await.map_err(RuntimeError::from),
Err(e) => Err(e),
}
};
if let Err(e) = fts_result {
let cleanup_errors = self
.compensate_entity_create(token, entity.id, ns, &[])
.await;
return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
}
}
let embed_model_names = {
let names = self.registered_embedding_model_names();
if names.is_empty() {
vec![]
} else {
names
}
};
let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
if embed_model_names.len() == 1 {
let model_name = &embed_model_names[0];
let vec_result = self
.embed_document_with_model_outcome_for_token(token, model_name, &embed_body)
.await;
#[cfg(any(test, feature = "fault-injection"))]
let vec_inject = consume_fault(&VECTOR_FAIL_NS, ns);
#[cfg(not(any(test, feature = "fault-injection")))]
let vec_inject = false;
let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
if vec_inject {
Err(RuntimeError::Internal(
"injected vector failure".to_string(),
))
} else {
vec_result
};
let single_result: RuntimeResult<()> = match vec_result {
Ok(outcome) => {
embedding_report.observe(&outcome);
match self.vectors_for_model(token, model_name) {
Ok(vs) => vs
.insert(
entity.id,
SubstrateKind::Entity,
ns,
"entity.body",
vec![outcome.vector],
)
.await
.map_err(RuntimeError::from),
Err(e) => Err(e),
}
}
Err(e) => Err(e),
};
if let Err(e) = single_result {
let cleanup_errors = self
.compensate_entity_create(
token,
entity.id,
ns,
std::slice::from_ref(model_name),
)
.await;
return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
}
} else if !embed_model_names.is_empty() {
let rt_clone = self.clone();
let body_owned = embed_body.clone();
let usage_ctx = crate::usage::current();
let mut join_set = tokio::task::JoinSet::new();
for (idx, model_name) in embed_model_names.iter().enumerate() {
let rt = rt_clone.clone();
let text = body_owned.clone();
let name = model_name.clone();
let ctx = usage_ctx.clone();
let token = (*token).clone();
join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
let fut = rt.embed_document_with_model_outcome_for_token(&token, &name, &text);
let result = match ctx {
Some(ctx) => crate::usage::scope(ctx, fut).await,
None => fut.await,
};
(idx, result)
}));
}
let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
Ok(outcomes) => outcomes,
Err(e) => {
let cleanup_errors = self
.compensate_entity_create(token, entity.id, ns, &[])
.await;
return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
}
};
let mut inserted_models: Vec<String> = Vec::with_capacity(embed_model_names.len());
for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
embedding_report.observe(&outcome);
#[cfg(any(test, feature = "fault-injection"))]
let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
Some(0) => {
cell.set(None);
true
}
Some(n) => {
cell.set(Some(n - 1));
false
}
None => false,
});
#[cfg(not(any(test, feature = "fault-injection")))]
let count_inject = false;
let insert_result = if count_inject {
Err(RuntimeError::Internal(
"injected vector insert failure".to_string(),
))
} else {
match self.vectors_for_model(token, model_name) {
Ok(vs) => vs
.insert(
entity.id,
SubstrateKind::Entity,
ns,
"entity.body",
vec![outcome.vector],
)
.await
.map_err(RuntimeError::from),
Err(e) => Err(e),
}
};
if let Err(e) = insert_result {
let mut cleanup_models = inserted_models.clone();
cleanup_models.push(model_name.clone());
let cleanup_errors = self
.compensate_entity_create(token, entity.id, ns, &cleanup_models)
.await;
return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
}
inserted_models.push(model_name.clone());
}
}
let created_event = khive_storage::event::Event::new(
entity.namespace.clone(),
"create",
EventKind::EntityCreated,
SubstrateKind::Entity,
"",
)
.with_target(entity.id)
.with_payload(serde_json::json!({
"id": entity.id,
"namespace": entity.namespace,
"kind": entity.kind,
}));
let event_result = match self.events(token) {
Ok(store) => store
.append_event(created_event)
.await
.map_err(RuntimeError::from),
Err(error) => Err(error),
};
let mut degradations = Vec::new();
if let Err(error) = event_result {
record_post_commit_degradation(
&mut degradations,
"create_entity",
entity.id,
"event_append",
error,
);
}
Ok((entity, embedding_report, degradations))
}
pub async fn get_entity(&self, token: &NamespaceToken, id: Uuid) -> RuntimeResult<Entity> {
let store = self.entities(token)?;
if let Some(entity) = store.get_entity(id).await? {
return Ok(entity);
}
if let Some(tombstone) = store.get_entity_including_deleted(id).await? {
if let Some(kept_id) = tombstone.merged_into {
return Err(RuntimeError::NotFound(format!(
"{id} was merged into {kept_id}; query the kept id"
)));
}
}
Err(RuntimeError::NotFound(format!("entity {id}")))
}
pub async fn get_entity_including_deleted(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Entity>> {
self.entities(token)?
.get_entity_including_deleted(id)
.await
.map_err(Into::into)
}
pub async fn get_note_including_deleted(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<khive_storage::note::Note>> {
self.notes(token)?
.get_note_including_deleted(id)
.await
.map_err(Into::into)
}
pub async fn get_entities_by_ids(
&self,
token: &NamespaceToken,
ids: &[Uuid],
) -> RuntimeResult<Vec<Entity>> {
if ids.is_empty() {
return Ok(vec![]);
}
let filter = EntityFilter {
ids: ids.to_vec(),
..Default::default()
};
let page = self
.entities(token)?
.query_entities(
token.namespace().as_str(),
filter,
PageRequest {
offset: 0,
limit: ids.len() as u32,
},
)
.await?;
Ok(page.items)
}
async fn get_entities_by_ids_visible(
&self,
token: &NamespaceToken,
ids: &[Uuid],
) -> RuntimeResult<Vec<Entity>> {
if ids.is_empty() {
return Ok(vec![]);
}
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_owned())
.collect();
let filter = EntityFilter {
ids: ids.to_vec(),
namespaces,
..Default::default()
};
let page = self
.entities(token)?
.query_entities(
token.namespace().as_str(),
filter,
PageRequest {
offset: 0,
limit: ids.len() as u32,
},
)
.await?;
Ok(page.items)
}
pub(crate) fn ensure_namespace(record_ns: &str, caller_primary_ns: &str) -> RuntimeResult<()> {
if record_ns == caller_primary_ns {
return Ok(());
}
Err(RuntimeError::NotFound("not found in this namespace".into()))
}
pub(crate) fn ensure_namespace_visible(
record_ns: &str,
token: &NamespaceToken,
) -> RuntimeResult<()> {
for ns in token.visible_namespaces() {
if record_ns == ns.as_str() {
return Ok(());
}
}
Err(RuntimeError::NotFound("not found in this namespace".into()))
}
pub async fn list_entities(
&self,
token: &NamespaceToken,
kind: Option<&str>,
entity_type: Option<&str>,
limit: u32,
offset: u32,
) -> RuntimeResult<Vec<Entity>> {
let filter = EntityFilter {
kinds: kind
.map(|value| vec![value.to_string()])
.unwrap_or_default(),
entity_types: entity_type
.map(|value| vec![value.to_string()])
.unwrap_or_default(),
legacy_entity_type_fallback: true,
..Default::default()
};
self.list_entities_filtered(token, filter, limit, offset)
.await
}
pub async fn list_entities_filtered(
&self,
token: &NamespaceToken,
mut filter: EntityFilter,
limit: u32,
offset: u32,
) -> RuntimeResult<Vec<Entity>> {
filter.namespaces = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
let page = self
.entities(token)?
.query_entities_count_free(
token.namespace().as_str(),
filter,
PageRequest {
offset: offset.into(),
limit,
},
)
.await?;
Ok(page.items)
}
pub async fn list_entities_after(
&self,
token: &NamespaceToken,
kind: Option<&str>,
entity_type: Option<&str>,
tags_any: &[String],
after: Option<Uuid>,
limit: u32,
) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
let filter = EntityFilter {
kinds: kind
.map(|value| vec![value.to_string()])
.unwrap_or_default(),
entity_types: entity_type
.map(|value| vec![value.to_string()])
.unwrap_or_default(),
legacy_entity_type_fallback: true,
tags_any: tags_any.to_vec(),
..Default::default()
};
self.list_entities_after_filtered(token, filter, after, limit)
.await
}
pub async fn list_entities_after_filtered(
&self,
token: &NamespaceToken,
mut filter: EntityFilter,
after: Option<Uuid>,
limit: u32,
) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
let store = self.entities(token)?;
let after = match after {
Some(id) => {
let entity = self
.get_entity_including_deleted(token, id)
.await?
.ok_or_else(|| RuntimeError::NotFound(format!("entity cursor {id}")))?;
Self::ensure_namespace_visible(&entity.namespace, token)?;
let sequence = store.entity_sequence(id).await?.ok_or_else(|| {
RuntimeError::Internal(format!(
"entity cursor {id} has no insertion-sequence ledger row"
))
})?;
Some(SeekCursor { sequence, id })
}
None => None,
};
filter.namespaces = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
let page = store
.query_entities_after(token.namespace().as_str(), filter, after, limit)
.await?;
Ok((page.items, page.next_after.map(|cursor| cursor.id)))
}
pub async fn list_entities_tagged(
&self,
token: &NamespaceToken,
kind: Option<&str>,
domain_tag: Option<&str>,
limit: u32,
offset: u32,
) -> RuntimeResult<Vec<Entity>> {
let ns_strs: Vec<String> = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_owned())
.collect();
let filter = EntityFilter {
kinds: match kind {
Some(k) => vec![k.to_string()],
None => vec![],
},
tags_any: match domain_tag {
Some(t) if !t.is_empty() => vec![t.to_string()],
_ => vec![],
},
namespaces: ns_strs,
..Default::default()
};
let page = self
.entities(token)?
.query_entities_count_free(
token.namespace().as_str(),
filter,
PageRequest {
offset: offset.into(),
limit,
},
)
.await?;
Ok(page.items)
}
pub async fn count_entities_tagged(
&self,
token: &NamespaceToken,
kind: Option<&str>,
domain_tag: Option<&str>,
) -> RuntimeResult<u64> {
let ns_strs: Vec<String> = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_owned())
.collect();
let filter = EntityFilter {
kinds: match kind {
Some(k) => vec![k.to_string()],
None => vec![],
},
tags_any: match domain_tag {
Some(t) if !t.is_empty() => vec![t.to_string()],
_ => vec![],
},
namespaces: ns_strs,
..Default::default()
};
Ok(self
.entities(token)?
.count_entities(token.namespace().as_str(), filter)
.await?)
}
pub async fn list_events(
&self,
token: &NamespaceToken,
filter: EventFilter,
page: PageRequest,
) -> RuntimeResult<Page<Event>> {
self.events(token)?
.query_events(filter, page)
.await
.map_err(Into::into)
}
pub(crate) async fn validate_edge_relation_endpoints(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
) -> RuntimeResult<(EdgeEndpointKind, EdgeEndpointKind)> {
if source_id == target_id {
return Err(RuntimeError::InvalidInput(
"self-loop edges are not allowed: source_id and target_id must be different".into(),
));
}
if relation == EdgeRelation::Annotates {
match self.resolve_edge_endpoint(token, source_id).await? {
Some(Resolved::Note(_)) => {}
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"annotates source {source_id} must be a note"
)));
}
None => {
if self.get_edge(token, source_id).await?.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"annotates source {source_id} must be a note"
)));
}
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
}
let target_kind = match self.resolve_edge_endpoint(token, target_id).await? {
Some(Resolved::Entity(_)) => EdgeEndpointKind::Entity,
Some(Resolved::Note(_)) => EdgeEndpointKind::Note,
Some(Resolved::Event(_)) => EdgeEndpointKind::Event,
Some(Resolved::PackRecord { .. }) => {
return Err(RuntimeError::InvalidInput(
"pack-private record is not a valid edge endpoint for annotates".into(),
));
}
None => match self.get_edge(token, target_id).await {
Ok(Some(_)) => EdgeEndpointKind::Edge,
Ok(None) | Err(RuntimeError::NotFound(_)) => {
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
Err(error) => return Err(error),
},
};
return Ok((EdgeEndpointKind::Note, target_kind));
} else if crate::pack::is_special_relation(relation) {
let rel_name = relation.as_str();
let src = match self.resolve_edge_endpoint(token, source_id).await? {
Some(r) => r,
None => {
if self.get_edge(token, source_id).await?.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} source {source_id} must be a note or entity (got edge)"
)));
}
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
};
let tgt = match self.resolve_edge_endpoint(token, target_id).await? {
Some(r) => r,
None => {
if self.get_edge(token, target_id).await?.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} target {target_id} must be a note or entity (got edge)"
)));
}
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
};
return match (&src, &tgt) {
(Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
let legal_relations = accepted_entity_relations_description(
&self.pack_edge_rules(),
&src_e.kind,
src_e.entity_type.as_deref(),
&tgt_e.kind,
tgt_e.entity_type.as_deref(),
);
let rule_hint = match relation {
EdgeRelation::Supports | EdgeRelation::Refutes => {
"requires concept|document|dataset|artifact -> concept \
(or same-substrate note -> note)"
}
_ => "requires same-kind entity endpoints",
};
return Err(RuntimeError::InvalidInput(format!(
"({}) -[{rel_name}]-> ({}) is not in the base endpoint \
allowlist; {rel_name} {rule_hint}; currently legal relations for \
{} -> {} under the loaded endpoint rules: {legal_relations}",
src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
)));
}
Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
}
(Resolved::Note(_), Resolved::Note(_)) => {
Ok((EdgeEndpointKind::Note, EdgeEndpointKind::Note))
}
(Resolved::Event(_), _) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} does not apply to events; source {source_id} is an event"
)));
}
(_, Resolved::Event(_)) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} does not apply to events; target {target_id} is an event"
)));
}
(Resolved::Entity(_), Resolved::Note(_)) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
got source={source_id} (entity) target={target_id} (note)"
)));
}
(Resolved::Note(_), Resolved::Entity(_)) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
got source={source_id} (note) target={target_id} (entity)"
)));
}
(Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
return Err(RuntimeError::InvalidInput(format!(
"pack-private record is not a valid edge endpoint for {rel_name}"
)));
}
};
} else {
let src_res = self.resolve_edge_endpoint(token, source_id).await?;
let tgt_res = self.resolve_edge_endpoint(token, target_id).await?;
let pack_rules = self.pack_edge_rules();
if pack_rule_allows(&pack_rules, relation, src_res.as_ref(), tgt_res.as_ref()) {
let kind = |resolved: Option<&Resolved>| match resolved {
Some(Resolved::Entity(_)) => Some(EdgeEndpointKind::Entity),
Some(Resolved::Note(_)) => Some(EdgeEndpointKind::Note),
_ => None,
};
return match (kind(src_res.as_ref()), kind(tgt_res.as_ref())) {
(Some(source_kind), Some(target_kind)) => Ok((source_kind, target_kind)),
_ => Err(RuntimeError::Internal(
"pack endpoint rule admitted an unsupported substrate".into(),
)),
};
}
let (src_kind, src_entity_type) = match src_res.as_ref() {
Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"link source {source_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
None => {
if self.get_edge(token, source_id).await?.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"link source {source_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
};
let (tgt_kind, tgt_entity_type) = match tgt_res.as_ref() {
Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"link target {target_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
None => {
if self.get_edge(token, target_id).await?.is_some() {
return Err(RuntimeError::InvalidInput(format!(
"link target {target_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
};
if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
let legal_relations = accepted_entity_relations_description(
&pack_rules,
src_kind,
src_entity_type,
tgt_kind,
tgt_entity_type,
);
return Err(RuntimeError::InvalidInput(format!(
"({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
allowlist; use pack EDGE_RULES to extend the allowlist; currently legal \
relations for {src_kind} -> {tgt_kind} under the loaded endpoint rules: \
{legal_relations}",
relation.as_str()
)));
}
}
Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
}
pub async fn validate_link_endpoints(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
) -> RuntimeResult<()> {
self.validate_edge_relation_endpoints(token, source_id, target_id, relation)
.await
.map(|_| ())
}
pub fn validate_link_endpoints_by_resolved(
&self,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
src: Option<&Resolved>,
tgt: Option<&Resolved>,
) -> RuntimeResult<()> {
if source_id == target_id {
return Err(RuntimeError::InvalidInput(
"self-loop edges are not allowed: source_id and target_id must be different".into(),
));
}
if relation == EdgeRelation::Annotates {
match src {
Some(Resolved::Note(_)) => {}
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"annotates source {source_id} must be a note"
)));
}
None => {
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
}
if tgt.is_none() {
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
return Ok(());
}
if crate::pack::is_special_relation(relation) {
let rel_name = relation.as_str();
let src = src.ok_or_else(|| {
RuntimeError::NotFound(format!("link source {source_id} not found"))
})?;
let tgt = tgt.ok_or_else(|| {
RuntimeError::NotFound(format!("link target {target_id} not found"))
})?;
match (src, tgt) {
(Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
let legal_relations = accepted_entity_relations_description(
&self.pack_edge_rules(),
&src_e.kind,
src_e.entity_type.as_deref(),
&tgt_e.kind,
tgt_e.entity_type.as_deref(),
);
let rule_hint = match relation {
EdgeRelation::Supports | EdgeRelation::Refutes => {
"requires concept|document|dataset|artifact -> concept \
(or same-substrate note -> note)"
}
_ => "requires same-kind entity endpoints",
};
return Err(RuntimeError::InvalidInput(format!(
"({}) -[{rel_name}]-> ({}) is not in the base endpoint \
allowlist; {rel_name} {rule_hint}; currently legal relations for \
{} -> {} under the loaded endpoint rules: {legal_relations}",
src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
)));
}
}
(Resolved::Note(_), Resolved::Note(_)) => {}
(Resolved::Entity(_), Resolved::Note(_)) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} endpoints must be the same substrate \
(note→note or entity→entity); got source={source_id} (entity) \
target={target_id} (note)"
)));
}
(Resolved::Note(_), Resolved::Entity(_)) => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} endpoints must be the same substrate \
(note→note or entity→entity); got source={source_id} (note) \
target={target_id} (entity)"
)));
}
(Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
return Err(RuntimeError::InvalidInput(format!(
"pack-private record is not a valid edge endpoint for {rel_name}"
)));
}
_ => {
return Err(RuntimeError::InvalidInput(format!(
"{rel_name} endpoints must be notes or entities (not events)"
)));
}
}
return Ok(());
}
let pack_rules = self.pack_edge_rules();
if pack_rule_allows(&pack_rules, relation, src, tgt) {
return Ok(());
}
let (src_kind, src_entity_type) = match src {
Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"link source {source_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
None => {
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
};
let (tgt_kind, tgt_entity_type) = match tgt {
Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"link target {target_id} must be an entity for relation {relation:?} \
(only `annotates` crosses substrates)"
)));
}
None => {
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
};
if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
let legal_relations = accepted_entity_relations_description(
&pack_rules,
src_kind,
src_entity_type,
tgt_kind,
tgt_entity_type,
);
return Err(RuntimeError::InvalidInput(format!(
"({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
allowlist; use pack EDGE_RULES to extend the allowlist; currently legal relations \
for {src_kind} -> {tgt_kind} under the loaded endpoint rules: {legal_relations}",
relation.as_str()
)));
}
Ok(())
}
pub fn validate_annotates_endpoint_kinds(
&self,
source_id: Uuid,
target_id: Uuid,
source: Option<EdgeEndpointKind>,
target: Option<EdgeEndpointKind>,
) -> RuntimeResult<()> {
if source_id == target_id {
return Err(RuntimeError::InvalidInput(
"self-loop edges are not allowed: source_id and target_id must be different".into(),
));
}
match source {
Some(EdgeEndpointKind::Note) => {}
Some(_) => {
return Err(RuntimeError::InvalidInput(format!(
"annotates source {source_id} must be a note"
)));
}
None => {
return Err(RuntimeError::NotFound(format!(
"link source {source_id} not found"
)));
}
}
if target.is_none() {
return Err(RuntimeError::NotFound(format!(
"link target {target_id} not found"
)));
}
Ok(())
}
pub async fn link(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
weight: f64,
metadata: Option<serde_json::Value>,
) -> RuntimeResult<Edge> {
self.link_observed(
token, source_id, target_id, relation, weight, metadata, false,
)
.await
.map(|result| result.edge)
}
#[allow(clippy::too_many_arguments)]
pub async fn link_observed(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
weight: f64,
metadata: Option<serde_json::Value>,
resurrect: bool,
) -> RuntimeResult<EdgeUpsertResult> {
validate_edge_weight(weight)?;
validate_edge_metadata(relation, metadata.as_ref())?;
let (source_kind, target_kind) = self
.validate_edge_relation_endpoints(token, source_id, target_id, relation)
.await?;
let (canonical_source, canonical_target) =
canonical_edge_endpoints(relation, source_id, target_id);
let (source_kind, target_kind) =
canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
let (source_id, target_id) = (canonical_source, canonical_target);
let metadata = if relation == EdgeRelation::DependsOn {
match (
self.resolve_edge_endpoint(token, source_id).await?,
self.resolve_edge_endpoint(token, target_id).await?,
) {
(Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
merge_dependency_kind(&src_e.kind, &tgt_e.kind, metadata)
}
_ => metadata,
}
} else {
metadata
};
validate_edge_metadata(relation, metadata.as_ref())?;
let now = chrono::Utc::now();
let ns = token.namespace().as_str();
let edge = Edge {
id: LinkId::from(Uuid::new_v4()),
namespace: ns.to_string(),
source_id,
target_id,
relation,
weight,
created_at: now,
updated_at: now,
deleted_at: None,
metadata,
target_backend: None,
};
let attribution = crate::EventAttribution::from_token(token);
let outcome = compose_graph_mutation_events(
self.backend(),
GraphMutationRequest::Single {
request: EdgeUpsertRequest { edge, resurrect },
guard_endpoints: true,
},
GraphMutationPreconditions::default(),
Vec::new(),
move |outcome| match &outcome.mutation {
GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
Ok(vec![Self::link_mutation_event(
&attribution,
result,
source_kind,
target_kind,
)])
}
_ => Err(Self::link_composition_shape_error(
"expected a written singleton",
)),
},
)
.await?;
let GraphMutationOutcome::Single(outcome) = outcome.mutation else {
return Err(RuntimeError::Internal(
"link: unexpected composition outcome".into(),
));
};
let result = match outcome {
GuardedEdgeUpsertOutcome::Written(result) => result,
GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::MissingEndpoints(missing)) => {
return Err(RuntimeError::GuardedWriteFailed(GuardedWriteFailure {
entry_index: None,
missing_source: missing.source.then_some(source_id),
missing_target: missing.target.then_some(target_id),
}));
}
GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::ResurrectionRequired { edge }) => {
return Err(RuntimeError::InvalidInput(format!(
"edge {} is soft-deleted; pass resurrect=true to link explicitly",
edge.id
)))
}
};
Ok(result)
}
#[allow(clippy::too_many_arguments)]
pub async fn link_with_target_backend(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
source_kind: EdgeEndpointKind,
target_kind: EdgeEndpointKind,
relation: EdgeRelation,
weight: f64,
metadata: Option<serde_json::Value>,
target_backend: Option<String>,
) -> RuntimeResult<Edge> {
self.link_with_target_backend_observed(
token,
source_id,
target_id,
source_kind,
target_kind,
relation,
weight,
metadata,
target_backend,
false,
)
.await
.map(|result| result.edge)
}
#[allow(clippy::too_many_arguments)]
pub async fn link_with_target_backend_observed(
&self,
token: &NamespaceToken,
source_id: Uuid,
target_id: Uuid,
source_kind: EdgeEndpointKind,
target_kind: EdgeEndpointKind,
relation: EdgeRelation,
weight: f64,
metadata: Option<serde_json::Value>,
target_backend: Option<String>,
resurrect: bool,
) -> RuntimeResult<EdgeUpsertResult> {
validate_edge_weight(weight)?;
let (canonical_source, canonical_target) =
canonical_edge_endpoints(relation, source_id, target_id);
let (source_kind, target_kind) =
canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
let (source_id, target_id) = (canonical_source, canonical_target);
validate_edge_metadata(relation, metadata.as_ref())?;
let now = chrono::Utc::now();
let ns = token.namespace().as_str();
let edge = Edge {
id: LinkId::from(Uuid::new_v4()),
namespace: ns.to_string(),
source_id,
target_id,
relation,
weight,
created_at: now,
updated_at: now,
deleted_at: None,
metadata,
target_backend,
};
let attribution = crate::EventAttribution::from_token(token);
let outcome = compose_graph_mutation_events(
self.backend(),
GraphMutationRequest::Single {
request: EdgeUpsertRequest { edge, resurrect },
guard_endpoints: false,
},
GraphMutationPreconditions::default(),
Vec::new(),
move |outcome| match &outcome.mutation {
GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
Ok(vec![Self::link_mutation_event(
&attribution,
result,
source_kind,
target_kind,
)])
}
_ => Err(Self::link_composition_shape_error(
"expected a written singleton",
)),
},
)
.await?;
match outcome.mutation {
GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => Ok(result),
GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Refused(
EdgeUpsertRefusal::ResurrectionRequired { edge },
)) => {
let error = khive_storage::StorageError::Conflict {
capability: khive_storage::StorageCapability::Graph,
operation: "upsert_edge_observed".into(),
message: format!(
"edge {} is soft-deleted; explicit resurrection is required",
edge.id,
),
};
Err(RuntimeError::InvalidInput(format!(
"edge natural key is soft-deleted; pass resurrect=true to link explicitly: {error}"
)))
}
_ => Err(RuntimeError::Internal(
"link: unexpected composition outcome".into(),
)),
}
}
fn link_mutation_event(
attribution: &crate::EventAttribution,
result: &EdgeUpsertResult,
source_kind: EdgeEndpointKind,
target_kind: EdgeEndpointKind,
) -> Event {
let kind = match result.disposition {
EdgeUpsertDisposition::Created => EventKind::LinkCreated,
EdgeUpsertDisposition::Updated | EdgeUpsertDisposition::Resurrected => {
EventKind::EdgeUpdated
}
};
let edge_id = Uuid::from(result.edge.id);
let mut payload = serde_json::json!({
"id": edge_id,
"namespace": result.edge.namespace,
"mutation": result.disposition.name(),
"source_id": result.edge.source_id,
"target_id": result.edge.target_id,
"relation": result.edge.relation,
"weight": result.edge.weight,
"metadata": result.edge.metadata,
"previous": result.previous,
});
if kind == EventKind::LinkCreated {
payload["source_kind"] = serde_json::json!(source_kind.name());
payload["target_kind"] = serde_json::json!(target_kind.name());
}
attribution.stamp(
Event::new(
result.edge.namespace.clone(),
"link",
kind,
SubstrateKind::Entity,
"",
)
.with_target(edge_id)
.with_payload(payload),
)
}
fn link_composition_shape_error(message: &'static str) -> khive_storage::StorageError {
khive_storage::StorageError::InvalidInput {
capability: khive_storage::StorageCapability::Graph,
operation: "link_mutation_event".into(),
message: message.into(),
}
}
pub(crate) async fn substrate_exists_in_ns(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<bool> {
if self.resolve(token, id).await?.is_some() {
return Ok(true);
}
match self.get_edge_visible(token, id).await {
Ok(Some(_)) => Ok(true),
Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
Err(err) => Err(err),
}
}
pub(crate) async fn substrate_exists_by_id(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<bool> {
if self.resolve_edge_endpoint(token, id).await?.is_some() {
return Ok(true);
}
match self.get_edge(token, id).await {
Ok(Some(_)) => Ok(true),
Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
Err(err) => Err(err),
}
}
pub async fn latest_annotating_note(
&self,
token: &NamespaceToken,
node_id: Uuid,
kind: &str,
tag: &str,
) -> RuntimeResult<Option<Uuid>> {
self.latest_annotating_note_inner(token, node_id, kind, tag, None)
.await
}
pub async fn latest_annotating_note_with_property(
&self,
token: &NamespaceToken,
node_id: Uuid,
kind: &str,
tag: &str,
property_key: &str,
property_value: &str,
) -> RuntimeResult<Option<Uuid>> {
self.latest_annotating_note_inner(
token,
node_id,
kind,
tag,
Some((property_key, property_value)),
)
.await
}
async fn latest_annotating_note_inner(
&self,
token: &NamespaceToken,
node_id: Uuid,
kind: &str,
tag: &str,
required_property: Option<(&str, &str)>,
) -> RuntimeResult<Option<Uuid>> {
if !self.substrate_exists_in_ns(token, node_id).await? {
return Ok(None);
}
let mut latest: Option<(Uuid, i64)> = None;
for namespace in token.visible_namespaces() {
let scoped = NamespaceToken::for_namespace(namespace.clone());
let graph = self.graph(&scoped)?;
let candidate = match required_property {
Some((key, value)) => {
graph
.latest_annotating_note_with_property(node_id, kind, tag, key, value)
.await?
}
None => graph.latest_annotating_note(node_id, kind, tag).await?,
};
if let Some(candidate) = candidate {
if latest.is_none_or(|(id, created_at)| {
candidate.1 > created_at || (candidate.1 == created_at && candidate.0 < id)
}) {
latest = Some(candidate);
}
}
}
Ok(latest.map(|(id, _)| id))
}
pub async fn neighbors(
&self,
token: &NamespaceToken,
node_id: Uuid,
direction: Direction,
limit: Option<u32>,
relations: Option<Vec<EdgeRelation>>,
) -> RuntimeResult<Vec<NeighborHit>> {
self.neighbors_with_query(
token,
node_id,
NeighborQuery {
direction,
relations,
limit,
min_weight: None,
},
)
.await
}
pub async fn neighbors_with_query(
&self,
token: &NamespaceToken,
node_id: Uuid,
query: NeighborQuery,
) -> RuntimeResult<Vec<NeighborHit>> {
self.neighbors_with_query_page(token, node_id, query, None, None, true)
.await
}
pub async fn neighbors_with_query_page(
&self,
token: &NamespaceToken,
node_id: Uuid,
query: NeighborQuery,
after: Option<NeighborCursor>,
neighbor_kinds: Option<Vec<String>>,
enrich: bool,
) -> RuntimeResult<Vec<NeighborHit>> {
if !self.substrate_exists_by_id(token, node_id).await? {
return Err(RuntimeError::NotFound(format!(
"neighbor anchor {node_id} not found"
)));
}
self.neighbors_for_resolved_kg_read(
token,
node_id,
crate::KgNeighborRead {
query,
after,
neighbor_kinds,
enrich,
namespace: None,
},
)
.await
}
pub async fn neighbors_for_resolved_kg_read(
&self,
token: &NamespaceToken,
node_id: Uuid,
options: crate::KgNeighborRead,
) -> RuntimeResult<Vec<NeighborHit>> {
self.neighbors_for_resolved_kg_read_inner(token, node_id, options, false)
.await
.map(|(hits, _)| hits)
}
pub async fn neighbors_for_resolved_kg_read_with_entity_kinds(
&self,
token: &NamespaceToken,
node_id: Uuid,
options: crate::KgNeighborRead,
) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
self.neighbors_for_resolved_kg_read_inner(token, node_id, options, true)
.await
}
async fn neighbors_for_resolved_kg_read_inner(
&self,
token: &NamespaceToken,
node_id: Uuid,
options: crate::KgNeighborRead,
with_entity_kinds: bool,
) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
let crate::KgNeighborRead {
mut query,
after,
neighbor_kinds,
enrich,
namespace,
} = options;
let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace.as_ref())?;
query.direction =
normalize_symmetric_direction(query.direction, query.relations.as_deref());
let mut hits = Vec::new();
for ns in namespaces {
let temp = NamespaceToken::for_namespace(ns.clone());
let mut ns_hits = self
.graph(&temp)?
.neighbors_page(node_id, query.clone(), after, neighbor_kinds.clone())
.await?;
hits.append(&mut ns_hits);
}
hits.sort_by_key(|h| (h.node_id, h.edge_id));
hits.dedup_by_key(|h| (h.node_id, h.edge_id));
if enrich {
self.enrich_neighbor_hits(token, &mut hits).await;
}
let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.node_id).collect();
let (deleted, entity_kinds) = self
.neighbor_node_screen(
candidate_ids,
(with_entity_kinds && !enrich).then_some(token),
)
.await?;
if !deleted.is_empty() {
hits.retain(|h| !deleted.contains(&h.node_id));
}
hits.sort_by(|a, b| {
b.weight
.partial_cmp(&a.weight)
.unwrap_or(std::cmp::Ordering::Equal)
.then(a.node_id.cmp(&b.node_id))
.then(a.edge_id.cmp(&b.edge_id))
});
Ok((hits, entity_kinds))
}
pub async fn annotation_neighbors_by_target_id(
&self,
target_id: Uuid,
) -> RuntimeResult<Vec<NeighborHit>> {
let mut reader = self.sql().reader().await?;
let rows = reader
.query_all(SqlStatement {
sql: "SELECT source_id, id, weight FROM graph_edges \
WHERE target_id = ?1 AND relation = ?2 AND deleted_at IS NULL \
ORDER BY weight DESC, source_id ASC"
.to_string(),
params: vec![
SqlValue::Text(target_id.to_string()),
SqlValue::Text(EdgeRelation::Annotates.to_string()),
],
label: Some("annotations.by_target_id_unfiltered".into()),
})
.await?;
rows.into_iter()
.map(|row| {
let parse_uuid = |name: &str| match row.get(name) {
Some(SqlValue::Text(value)) => Uuid::from_str(value).map_err(|error| {
RuntimeError::Internal(format!("graph_edges.{name} is not a UUID: {error}"))
}),
Some(value) => Err(RuntimeError::Internal(format!(
"graph_edges.{name} has unexpected SQL value {value:?}"
))),
None => Err(RuntimeError::Internal(format!(
"graph_edges row missing {name}"
))),
};
let weight = match row.get("weight") {
Some(SqlValue::Float(value)) => Ok(*value),
Some(value) => Err(RuntimeError::Internal(format!(
"graph_edges.weight has unexpected SQL value {value:?}"
))),
None => Err(RuntimeError::Internal(
"graph_edges row missing weight".into(),
)),
}?;
Ok(NeighborHit {
node_id: parse_uuid("source_id")?,
edge_id: parse_uuid("id")?,
relation: EdgeRelation::Annotates,
weight,
name: None,
kind: None,
entity_type: None,
})
})
.collect()
}
pub async fn neighbors_with_query_directed(
&self,
token: &NamespaceToken,
node_id: Uuid,
query: NeighborQuery,
) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
if !self.substrate_exists_by_id(token, node_id).await? {
return Err(RuntimeError::NotFound(format!(
"neighbor anchor {node_id} not found"
)));
}
self.directed_neighbors_for_resolved_kg_read(token, node_id, query, None)
.await
}
pub(crate) async fn directed_neighbors_for_resolved_kg_read(
&self,
token: &NamespaceToken,
node_id: Uuid,
query: NeighborQuery,
namespace: Option<&crate::Namespace>,
) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace)?;
let mut hits: Vec<DirectedNeighborHit> = Vec::new();
for ns in namespaces {
let temp = NamespaceToken::for_namespace(ns.clone());
let mut ns_hits = self
.graph(&temp)?
.neighbors_both_directions(node_id, query.clone())
.await?;
hits.append(&mut ns_hits);
}
hits.sort_by_key(|h| {
(
h.hit.node_id,
h.hit.edge_id,
direction_sort_rank(&h.direction),
)
});
hits.dedup_by_key(|h| {
(
h.hit.node_id,
h.hit.edge_id,
direction_sort_rank(&h.direction),
)
});
let mut plain_hits: Vec<NeighborHit> = hits.iter().map(|h| h.hit.clone()).collect();
self.enrich_neighbor_hits(token, &mut plain_hits).await;
for (dh, enriched) in hits.iter_mut().zip(plain_hits) {
dh.hit = enriched;
}
let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.hit.node_id).collect();
let deleted = self.deleted_entity_ids(candidate_ids).await?;
if !deleted.is_empty() {
hits.retain(|h| !deleted.contains(&h.hit.node_id));
}
hits.sort_by(|a, b| {
b.hit
.weight
.partial_cmp(&a.hit.weight)
.unwrap_or(std::cmp::Ordering::Equal)
.then(a.hit.node_id.cmp(&b.hit.node_id))
.then(a.hit.edge_id.cmp(&b.hit.edge_id))
});
Ok(hits.into_iter().map(|h| (h.hit, h.direction)).collect())
}
pub async fn traverse(
&self,
token: &NamespaceToken,
request: TraversalRequest,
) -> RuntimeResult<Vec<GraphPath>> {
let mut request = request;
request.validate().map_err(RuntimeError::InvalidInput)?;
let mut roots = Vec::with_capacity(request.roots.len());
let mut seen_roots = std::collections::HashSet::with_capacity(request.roots.len());
for root in request.roots.drain(..) {
if seen_roots.insert(root) {
if !self.substrate_exists_by_id(token, root).await? {
return Err(RuntimeError::NotFound(format!(
"traverse root {root} not found"
)));
}
roots.push(root);
}
}
request.roots = roots;
if request.roots.is_empty() {
return Ok(Vec::new());
}
let mut paths = Vec::new();
for ns in token.visible_namespaces() {
let temp = NamespaceToken::for_namespace(ns.clone());
let mut ns_paths = self.graph(&temp)?.traverse(request.clone()).await?;
paths.append(&mut ns_paths);
}
let mut paths =
merge_traversal_paths_by_root(paths, Some(request.options.effective_limit()));
self.enrich_path_nodes(token, &mut paths, request.include_properties)
.await;
let all_node_ids: Vec<Uuid> = paths
.iter()
.flat_map(|p| p.nodes.iter().map(|n| n.node_id))
.collect();
let deleted = self.deleted_entity_ids(all_node_ids).await?;
if !deleted.is_empty() {
for path in paths.iter_mut() {
path.nodes.retain(|n| !deleted.contains(&n.node_id));
recompute_total_weight(path);
}
paths.retain(|p| !p.nodes.is_empty());
}
Ok(paths)
}
async fn deleted_entity_ids(
&self,
ids: Vec<Uuid>,
) -> RuntimeResult<std::collections::HashSet<Uuid>> {
self.neighbor_node_screen(ids, None)
.await
.map(|(deleted, _)| deleted)
}
async fn neighbor_node_screen(
&self,
ids: Vec<Uuid>,
kind_token: Option<&NamespaceToken>,
) -> RuntimeResult<(std::collections::HashSet<Uuid>, HashMap<Uuid, String>)> {
if ids.is_empty() {
return Ok((std::collections::HashSet::new(), HashMap::new()));
}
let id_strs: Vec<String> = ids.iter().map(|u| u.to_string()).collect();
let n = id_strs.len();
let entities_placeholders = (0..n)
.map(|i| format!("?{}", i + 1))
.collect::<Vec<_>>()
.join(",");
let notes_placeholders = (0..n)
.map(|i| format!("?{}", n + i + 1))
.collect::<Vec<_>>()
.join(",");
let sql_str = if kind_token.is_some() {
format!(
"SELECT id, kind, namespace, deleted_at IS NOT NULL AS is_deleted \
FROM entities WHERE id IN ({entities_placeholders}) \
UNION ALL \
SELECT id, NULL, NULL, 1 FROM notes \
WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
)
} else {
format!(
"SELECT id FROM entities WHERE id IN ({entities_placeholders}) AND deleted_at IS NOT NULL \
UNION \
SELECT id FROM notes WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
)
};
let params: Vec<SqlValue> = id_strs
.iter()
.chain(id_strs.iter())
.cloned()
.map(SqlValue::Text)
.collect();
let stmt = SqlStatement {
sql: sql_str,
params,
label: Some("deleted_entity_ids".into()),
};
let mut out = std::collections::HashSet::new();
let mut entity_kinds = HashMap::new();
let sql = self.sql();
let mut reader = sql.reader().await?;
let rows = reader.query_all(stmt).await?;
for row in rows {
if let Some(col) = row.columns.first() {
if let SqlValue::Text(s) = &col.value {
if let Ok(u) = s.parse::<Uuid>() {
if kind_token.is_none()
|| matches!(
row.columns.get(3).map(|col| &col.value),
Some(SqlValue::Integer(1))
)
{
out.insert(u);
} else if let (
Some(token),
Some(SqlValue::Text(kind)),
Some(SqlValue::Text(namespace)),
Some(SqlValue::Integer(0)),
) = (
kind_token,
row.columns.get(1).map(|col| &col.value),
row.columns.get(2).map(|col| &col.value),
row.columns.get(3).map(|col| &col.value),
) {
if token
.visible_namespaces()
.iter()
.any(|ns| ns.as_str() == namespace.as_str())
{
entity_kinds.insert(u, kind.clone());
}
}
}
}
}
}
Ok((out, entity_kinds))
}
async fn enrich_neighbor_hits(&self, token: &NamespaceToken, hits: &mut [NeighborHit]) {
if hits.is_empty() {
return;
}
let unique_ids: Vec<Uuid> = {
let mut seen = std::collections::HashSet::new();
hits.iter()
.filter_map(|h| {
if seen.insert(h.node_id) {
Some(h.node_id)
} else {
None
}
})
.collect()
};
let entity_map: HashMap<Uuid, Entity> = self
.get_entities_by_ids_visible(token, &unique_ids)
.await
.unwrap_or_default()
.into_iter()
.map(|e| (e.id, e))
.collect();
let residual_ids: Vec<Uuid> = unique_ids
.iter()
.filter(|id| !entity_map.contains_key(id))
.copied()
.collect();
let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
if let Ok(store) = self.notes(token) {
store
.get_notes_batch(&residual_ids)
.await
.unwrap_or_default()
.into_iter()
.map(|n| (n.id, n))
.collect()
} else {
HashMap::new()
}
} else {
HashMap::new()
};
for hit in hits.iter_mut() {
if let Some(entity) = entity_map.get(&hit.node_id) {
hit.name = Some(entity.name.clone());
hit.kind = Some(entity.kind.clone());
hit.entity_type = entity.entity_type.clone();
} else if let Some(note) = note_map.get(&hit.node_id) {
hit.name = Some(note_graph_name(note));
hit.kind = Some(note.kind.clone());
}
}
}
async fn enrich_path_nodes(
&self,
token: &NamespaceToken,
paths: &mut [GraphPath],
include_properties: bool,
) {
if paths.is_empty() {
return;
}
let unique_ids: Vec<Uuid> = {
let mut seen = std::collections::HashSet::new();
paths
.iter()
.flat_map(|p| p.nodes.iter())
.filter_map(|n| {
if seen.insert(n.node_id) {
Some(n.node_id)
} else {
None
}
})
.collect()
};
let entity_map: HashMap<Uuid, Entity> = self
.get_entities_by_ids_visible(token, &unique_ids)
.await
.unwrap_or_default()
.into_iter()
.map(|e| (e.id, e))
.collect();
let residual_ids: Vec<Uuid> = unique_ids
.iter()
.filter(|id| !entity_map.contains_key(id))
.copied()
.collect();
let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
if let Ok(store) = self.notes(token) {
store
.get_notes_batch(&residual_ids)
.await
.unwrap_or_default()
.into_iter()
.map(|n| (n.id, n))
.collect()
} else {
HashMap::new()
}
} else {
HashMap::new()
};
for path in paths.iter_mut() {
for node in path.nodes.iter_mut() {
if let Some(entity) = entity_map.get(&node.node_id) {
node.name = Some(entity.name.clone());
node.kind = Some(entity.kind.clone());
if include_properties {
node.properties = entity.properties.clone();
}
} else if let Some(note) = note_map.get(&node.node_id) {
node.name = Some(note_graph_name(note));
node.kind = Some(note.kind.clone());
}
}
}
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<Note> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token, kind, name, content, None, salience, None, properties, annotates, None,
false, false,
)
.await?;
legacy_post_commit_result_with_embedding(
"create_note",
note.id,
note,
embedding,
degradations,
)
}
pub async fn create_web_receipt_note(
&self,
token: &NamespaceToken,
summary: &str,
request: serde_json::Value,
annotates: Vec<Uuid>,
) -> RuntimeResult<Note> {
let properties = serde_json::json!({
"tags": ["web.receipt"],
"request": request,
});
let (note, _, degradations, _) = self
.create_note_inner(
token,
"observation",
None,
summary,
None,
None,
None,
Some(properties),
annotates,
None,
false,
true,
)
.await?;
legacy_post_commit_result("create_web_receipt_note", note.id, note, degradations)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_embedding_content(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
embedding_content: Option<&str>,
salience: Option<f64>,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<Note> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token,
kind,
name,
content,
embedding_content,
salience,
None,
properties,
annotates,
None,
false,
false,
)
.await?;
legacy_post_commit_result_with_embedding(
"create_note_with_embedding_content",
note.id,
note,
embedding,
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_embedding_content_and_report(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
embedding_content: Option<&str>,
salience: Option<f64>,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token,
kind,
name,
content,
embedding_content,
salience,
None,
properties,
annotates,
None,
false,
false,
)
.await?;
legacy_post_commit_result(
"create_note_with_embedding_content_and_report",
note.id,
(note, embedding),
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_embedding_content_and_post_commit_report(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
embedding_content: Option<&str>,
salience: Option<f64>,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<(
Note,
crate::retrieval::EmbeddingTruncationReport,
Vec<PostCommitDegradation>,
)> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token,
kind,
name,
content,
embedding_content,
salience,
None,
properties,
annotates,
None,
false,
false,
)
.await?;
Ok((note, embedding, degradations))
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<Note> {
self.create_note_with_decay_for_embedding_model(
token,
kind,
name,
content,
salience,
decay_factor,
properties,
annotates,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay_and_report(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
self.create_note_with_decay_for_embedding_model_and_report(
token,
kind,
name,
content,
salience,
decay_factor,
properties,
annotates,
None,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay_for_embedding_model(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
embedding_model: Option<&str>,
) -> RuntimeResult<Note> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token,
kind,
name,
content,
None,
salience,
Some(decay_factor),
properties,
annotates,
embedding_model,
false,
false,
)
.await?;
legacy_post_commit_result_with_embedding(
"create_note_with_decay_for_embedding_model",
note.id,
note,
embedding,
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay_for_embedding_model_and_report(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
embedding_model: Option<&str>,
) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
let (note, embedding, degradations, _) = self
.create_note_inner(
token,
kind,
name,
content,
None,
salience,
Some(decay_factor),
properties,
annotates,
embedding_model,
false,
false,
)
.await?;
legacy_post_commit_result(
"create_note_with_decay_for_embedding_model_and_report",
note.id,
(note, embedding),
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay_for_embedding_model_with_visibility(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
embedding_model: Option<&str>,
) -> RuntimeResult<(Note, Vec<(String, u64)>)> {
let (note, _, degradations, fences) = self
.create_note_inner(
token,
kind,
name,
content,
None,
salience,
Some(decay_factor),
properties,
annotates,
embedding_model,
true,
false,
)
.await?;
legacy_post_commit_result(
"create_note_with_decay_for_embedding_model_with_visibility",
note.id,
(note, fences),
degradations,
)
}
#[allow(clippy::too_many_arguments)]
pub async fn create_note_with_decay_for_embedding_model_with_visibility_and_report(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
salience: Option<f64>,
decay_factor: f64,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
embedding_model: Option<&str>,
) -> RuntimeResult<(
Note,
Vec<(String, u64)>,
crate::retrieval::EmbeddingTruncationReport,
)> {
let (note, embedding, degradations, fences) = self
.create_note_inner(
token,
kind,
name,
content,
None,
salience,
Some(decay_factor),
properties,
annotates,
embedding_model,
true,
false,
)
.await?;
legacy_post_commit_result(
"create_note_with_decay_for_embedding_model_with_visibility_and_report",
note.id,
(note, fences, embedding),
degradations,
)
}
pub async fn try_create_note(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
properties: Option<serde_json::Value>,
) -> RuntimeResult<Option<Note>> {
self.try_create_note_impl(token, kind, name, content, properties, false, None, None)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn try_create_note_as_trusted_ingest(
&self,
_capability: &crate::pack::ChannelIngestCapability,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
properties: Option<serde_json::Value>,
expires_after: Option<std::time::Duration>,
) -> RuntimeResult<Option<Note>> {
self.try_create_note_impl(
token,
kind,
name,
content,
properties,
true,
None,
expires_after,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn try_create_note_as_trusted_ingest_with_attachment(
&self,
_capability: &crate::pack::ChannelIngestCapability,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
properties: Option<serde_json::Value>,
attachment: NewAttachment,
expires_after: Option<std::time::Duration>,
) -> RuntimeResult<Option<Note>> {
self.try_create_note_impl(
token,
kind,
name,
content,
properties,
true,
Some(attachment),
expires_after,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn try_create_note_impl(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
properties: Option<serde_json::Value>,
allow_transport_owned_message_properties: bool,
attachment: Option<NewAttachment>,
expires_after: Option<std::time::Duration>,
) -> RuntimeResult<Option<Note>> {
self.validate_note_kind(kind)?;
crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
crate::secret_gate::check_at(content, "note", "content")?;
if let Some(n) = name {
crate::secret_gate::check_at(n, "note", "name")?;
}
if let Some(ref p) = properties {
crate::secret_gate::check_json_at(p, "note", "properties")?;
}
if let Some(ref attachment) = attachment {
drop(self.attachments()?);
attachment.validate()?;
let blob_store = self.blob_store().ok_or_else(|| {
RuntimeError::Unconfigured(
"trusted ingest attachment requires an installed BlobStore".to_string(),
)
})?;
if !blob_store.exists(&attachment.content_ref).await? {
return Err(RuntimeError::InvalidInput(format!(
"trusted ingest attachment refers to an unpublished blob: {}",
attachment.content_ref
)));
}
}
if !allow_transport_owned_message_properties && kind == "message" {
if let Some(key) = properties
.as_ref()
.and_then(serde_json::Value::as_object)
.and_then(|supplied| {
crate::curation::kind_owned_properties("message")
.iter()
.copied()
.find(|key| supplied.contains_key(*key))
})
{
return Err(RuntimeError::InvalidInput(format!(
"`{key}` is transport-owned on a `message` note and cannot be supplied \
through `try_create_note`; only the trusted channel-ingest path may \
establish quarantine disposition and channel provenance"
)));
}
}
let ns = token.namespace().as_str();
let mut note = Note::new(ns, kind, content);
if let Some(retention) = expires_after {
let duration_us = i64::try_from(retention.as_micros()).map_err(|_| {
RuntimeError::InvalidInput(
"trusted ingest retention exceeds i64 microseconds".into(),
)
})?;
note.expires_at = Some(note.created_at.checked_add(duration_us).ok_or_else(|| {
RuntimeError::InvalidInput("trusted ingest expiry exceeds i64 microseconds".into())
})?);
}
if let Some(n) = name {
note = note.with_name(n);
}
if let Some(p) = properties {
note = note.with_properties(p);
}
let inserted = if let Some(attachment) = attachment {
self.raw_notes(token)?
.try_insert_note_with_attachments(
note.clone(),
vec![Attachment::from_new(
note.id,
AttachmentSubstrate::Note,
attachment,
note.created_at,
)],
)
.await?
} else {
self.raw_notes(token)?.try_insert_note(note.clone()).await?
};
if !inserted {
return Ok(None);
}
let mut degradations = Vec::new();
match self.text_for_notes(token) {
Ok(fts) => {
if let Err(error) = fts.upsert_document(note_fts_document(¬e)).await {
record_conditional_insert_degradation(
&mut degradations,
note.id,
ConditionalInsertStage::FtsUpsert,
error,
);
}
}
Err(error) => record_conditional_insert_degradation(
&mut degradations,
note.id,
ConditionalInsertStage::FtsAcquisition,
error,
),
}
let embed_model_names = self.embedding_models_for_note_kind(kind);
for model_name in &embed_model_names {
match self
.embed_document_with_model_outcome_for_token(
token,
model_name,
note_embedding_text_ref(¬e),
)
.await
{
Ok(outcome) => {
if outcome.truncated {
tracing::warn!(
note_id = %note.id,
model = %outcome.model_name,
source_bytes = outcome.source_bytes,
embedded_bytes = outcome.embedded_bytes,
"try_create_note: embedding input truncated; full content stored unchanged"
);
}
match self.vectors_for_model(token, model_name) {
Ok(vs) => {
if let Err(error) = vs
.insert(
note.id,
SubstrateKind::Note,
ns,
"note.content",
vec![outcome.vector],
)
.await
{
record_conditional_insert_degradation(
&mut degradations,
note.id,
ConditionalInsertStage::VectorInsert,
format!("model {model_name}: {error}"),
);
}
}
Err(error) => record_conditional_insert_degradation(
&mut degradations,
note.id,
ConditionalInsertStage::VectorAcquisition,
format!("model {model_name}: {error}"),
),
}
}
Err(error) => record_conditional_insert_degradation(
&mut degradations,
note.id,
ConditionalInsertStage::Embedding,
format!("model {model_name}: {error}"),
),
}
}
legacy_post_commit_result("try_create_note", note.id, Some(note), degradations)
}
#[allow(clippy::too_many_arguments)]
async fn create_note_inner(
&self,
token: &NamespaceToken,
kind: &str,
name: Option<&str>,
content: &str,
embedding_content: Option<&str>,
salience: Option<f64>,
decay_factor: Option<f64>,
properties: Option<serde_json::Value>,
annotates: Vec<Uuid>,
embedding_model: Option<&str>,
capture_visibility: bool,
web_receipt: bool,
) -> RuntimeResult<(
Note,
crate::retrieval::EmbeddingTruncationReport,
Vec<PostCommitDegradation>,
Vec<(String, u64)>,
)> {
self.validate_note_kind(kind)?;
let mut properties = self.derive_note_write_properties(kind, token, properties)?;
crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
if web_receipt {
let map = properties
.as_mut()
.and_then(serde_json::Value::as_object_mut)
.expect("web receipt properties are constructed as an object");
map.insert(
crate::secret_gate::RESERVED_WEB_RECEIPT_KEY.to_string(),
serde_json::Value::String(
crate::secret_gate::WEB_RECEIPT_PROVENANCE_VALUE.to_string(),
),
);
}
crate::secret_gate::check_at(content, "note", "content")?;
if let Some(n) = name {
crate::secret_gate::check_at(n, "note", "name")?;
}
if let Some(ref p) = properties {
crate::secret_gate::check_json_at(p, "note", "properties")?;
}
if let Some(ec) = embedding_content {
if ec.is_empty() {
return Err(RuntimeError::InvalidInput(
"embedding_content must not be empty".into(),
));
}
if ec.len() >= content.len() || !content.starts_with(ec) {
return Err(RuntimeError::InvalidInput(
"embedding_content must be a proper prefix of content".into(),
));
}
crate::secret_gate::check_at(ec, "note", "embedding_content")?;
}
let ns = token.namespace().as_str();
for &target_id in &annotates {
if !self.substrate_exists_by_id(token, target_id).await? {
return Err(RuntimeError::NotFound(format!(
"create_note annotates target {target_id} not found"
)));
}
}
if let Some(s) = salience {
if !s.is_finite() || !(0.0..=1.0).contains(&s) {
return Err(RuntimeError::InvalidInput(format!(
"salience must be a finite value in [0.0, 1.0]; got {s}"
)));
}
}
if let Some(d) = decay_factor {
if !d.is_finite() || d < 0.0 {
return Err(RuntimeError::InvalidInput(format!(
"decay_factor must be a finite value >= 0.0; got {d}"
)));
}
}
if let Some(model_name) = embedding_model {
self.resolve_embedding_model(Some(model_name))?;
}
let mut note = Note::new(ns, kind, content);
if let Some(s) = salience {
note = note.with_salience(s);
}
if let Some(df) = decay_factor {
note = note.with_decay(df);
}
if let Some(n) = name {
note = note.with_name(n);
}
if let Some(p) = properties {
note = note.with_properties(p);
}
let notes = if web_receipt {
self.raw_notes(token)?
} else {
self.notes(token)?
};
notes.upsert_note(note.clone()).await?;
let embed_model_names: Vec<String> = if let Some(m) = embedding_model {
vec![m.to_string()]
} else {
self.embedding_models_for_note_kind(kind)
};
{
#[cfg(any(test, feature = "fault-injection"))]
let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_inject = false;
let fts_result: RuntimeResult<()> = if fts_inject {
Err(RuntimeError::Internal("injected FTS failure".to_string()))
} else {
let statements =
khive_db::stores::text::delete_document_statements("fts_notes", ns, note.id)
.into_iter()
.chain(khive_db::stores::text::insert_document_statements(
"fts_notes",
¬e_fts_document(¬e),
))
.collect();
self.apply_note_index_revision(¬e, statements)
.await
.map(|_| ())
};
if let Err(e) = fts_result {
self.compensate_note_creation(¬e).await;
return Err(e);
}
}
let canonical_embed_text = note_embedding_text_ref(¬e);
let embed_text = embedding_content.unwrap_or(canonical_embed_text);
let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
let mut vector_fences = Vec::with_capacity(embed_model_names.len());
if embed_model_names.len() == 1 {
let model_name = &embed_model_names[0];
let vec_result = self
.embed_document_with_model_outcome_for_token(token, model_name, embed_text)
.await;
#[cfg(any(test, feature = "fault-injection"))]
let vec_inject = {
let ns_inject = consume_fault(&VECTOR_FAIL_NS, ns);
let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
Some(0) => {
cell.set(None);
true
}
Some(n) => {
cell.set(Some(n - 1));
false
}
None => false,
});
ns_inject || count_inject
};
#[cfg(not(any(test, feature = "fault-injection")))]
let vec_inject = false;
let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
if vec_inject {
Err(RuntimeError::Internal(
"injected vector failure".to_string(),
))
} else {
vec_result
};
let single_model_result: RuntimeResult<()> = match vec_result {
Ok(outcome) => {
embedding_report.observe(&outcome);
if capture_visibility {
match self
.publish_note_vector_revision_with_seq(
token,
¬e,
model_name,
&outcome.vector,
)
.await
{
Ok(Some(seq)) => {
vector_fences.push((model_name.clone(), seq));
Ok(())
}
Ok(None) => Ok(()),
Err(error) => Err(error),
}
} else {
self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
.await
.map(|_| ())
}
}
Err(e) => Err(e),
};
if let Err(e) = single_model_result {
self.compensate_note_creation(¬e).await;
return Err(e);
}
} else if !embed_model_names.is_empty() {
let rt_clone = self.clone();
let content_owned: std::sync::Arc<str> = std::sync::Arc::from(embed_text);
let usage_ctx = crate::usage::current();
let mut join_set = tokio::task::JoinSet::new();
for (idx, model_name) in embed_model_names.iter().enumerate() {
let rt = rt_clone.clone();
let text = std::sync::Arc::clone(&content_owned);
let name = model_name.clone();
let ctx = usage_ctx.clone();
let token = (*token).clone();
join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
let fut = rt.embed_document_with_model_outcome_for_token(
&token,
&name,
text.as_ref(),
);
let result = match ctx {
Some(ctx) => crate::usage::scope(ctx, fut).await,
None => fut.await,
};
(idx, result)
}));
}
let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
Ok(outcomes) => outcomes,
Err(e) => {
self.compensate_note_creation(¬e).await;
return Err(e);
}
};
for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
embedding_report.observe(&outcome);
let insert_result = if capture_visibility {
self.publish_note_vector_revision_with_seq(
token,
¬e,
model_name,
&outcome.vector,
)
.await
.map(|seq| {
if let Some(seq) = seq {
vector_fences.push((model_name.clone(), seq));
}
})
} else {
self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
.await
.map(|_| ())
};
if let Err(e) = insert_result {
self.compensate_note_creation(¬e).await;
return Err(e);
}
}
}
let mut created_edges: Vec<Uuid> = Vec::with_capacity(annotates.len());
#[cfg(test)]
let annotates_iter: Vec<(usize, Uuid)> = annotates
.iter()
.enumerate()
.map(|(i, &id)| (i, id))
.collect();
#[cfg(test)]
macro_rules! next_target {
($pair:expr) => {
$pair.1
};
}
#[cfg(not(test))]
let annotates_iter: Vec<Uuid> = annotates.to_vec();
#[cfg(not(test))]
macro_rules! next_target {
($pair:expr) => {
$pair
};
}
for pair in annotates_iter {
let target_id = next_target!(pair);
#[cfg(test)]
let injected_err: Option<RuntimeError> = {
let call_idx = pair.0;
LINK_FAIL_AFTER.with(|cell| {
let n = cell.get();
if n > 0 && call_idx + 1 == n {
cell.set(0); Some(RuntimeError::Internal("injected link failure".to_string()))
} else {
None
}
})
};
#[cfg(not(test))]
let injected_err: Option<RuntimeError> = None;
let link_result = if let Some(e) = injected_err {
Err(e)
} else {
self.link(
token,
note.id,
target_id,
EdgeRelation::Annotates,
1.0,
None,
)
.await
};
match link_result {
Ok(edge) => created_edges.push(edge.id.into()),
Err(e) => {
let edge_ids = created_edges
.iter()
.map(Uuid::to_string)
.collect::<Vec<_>>()
.join(", ");
match self.compensate_note_creation_with_edges(¬e).await {
Ok(true) => return Err(e),
Ok(false) => {
return Err(RuntimeError::Internal(format!(
"create_note: annotates link failed: {e}; note {} changed before \
compensation, retaining its incident edges [{edge_ids}]",
note.id
)));
}
Err(cleanup_error) => {
return Err(RuntimeError::Internal(format!(
"create_note: annotates link failed: {e}; compensation failed \
for note {} and retained edges [{edge_ids}]: {cleanup_error}; \
note and edges remain for reconciliation",
note.id
)));
}
}
}
}
}
let created_event = khive_storage::event::Event::new(
note.namespace.clone(),
"create",
EventKind::NoteCreated,
SubstrateKind::Note,
"",
)
.with_target(note.id)
.with_payload(serde_json::json!({
"id": note.id,
"namespace": note.namespace,
"kind": note.kind,
"salience": note.salience,
}));
let event_result = match self.events(token) {
Ok(store) => store
.append_event(created_event)
.await
.map_err(RuntimeError::from),
Err(error) => Err(error),
};
let mut degradations = Vec::new();
if let Err(error) = event_result {
record_post_commit_degradation(
&mut degradations,
"create_note",
note.id,
"event_append",
error,
);
}
vector_fences.sort_by(|a, b| a.0.cmp(&b.0));
Ok((note, embedding_report, degradations, vector_fences))
}
pub async fn list_notes(
&self,
token: &NamespaceToken,
kind: Option<&str>,
limit: u32,
offset: u32,
) -> RuntimeResult<Vec<Note>> {
let visible = token.visible_namespaces();
if visible.len() == 1 {
let page = self
.notes(token)?
.query_notes_count_free(
token.namespace().as_str(),
kind,
PageRequest {
offset: offset.into(),
limit,
},
)
.await?;
return Ok(page.items);
}
use khive_storage::note::NoteFilter;
let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
let filter = NoteFilter {
kind: kind.map(|k| k.to_string()),
namespaces: ns_strs,
..Default::default()
};
let page = self
.notes(token)?
.query_notes_filtered_count_free(
token.namespace().as_str(),
&filter,
PageRequest {
offset: offset.into(),
limit,
},
)
.await?;
Ok(page.items)
}
pub async fn list_notes_after(
&self,
token: &NamespaceToken,
kind: Option<&str>,
after: Option<Uuid>,
limit: u32,
) -> RuntimeResult<(Vec<Note>, Option<Uuid>)> {
let store = self.notes(token)?;
let after = match after {
Some(id) => {
let note = self
.get_note_including_deleted(token, id)
.await?
.ok_or_else(|| RuntimeError::NotFound(format!("note cursor {id}")))?;
Self::ensure_namespace_visible(¬e.namespace, token)?;
let sequence = store.note_sequence(id).await?.ok_or_else(|| {
RuntimeError::Internal(format!(
"note cursor {id} has no insertion-sequence ledger row"
))
})?;
Some(SeekCursor { sequence, id })
}
None => None,
};
let filter = khive_storage::note::NoteFilter {
kind: kind.map(str::to_string),
namespaces: token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect(),
..Default::default()
};
let page = store
.query_notes_filtered_after(token.namespace().as_str(), &filter, after, limit)
.await?;
Ok((page.items, page.next_after.map(|cursor| cursor.id)))
}
pub async fn count_notes(
&self,
token: &NamespaceToken,
kind: Option<&str>,
) -> RuntimeResult<u64> {
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
Ok(self
.notes(token)?
.count_notes_in_namespaces(&namespaces, kind)
.await?)
}
#[allow(clippy::too_many_arguments)]
pub async fn search_notes(
&self,
token: &NamespaceToken,
query_text: &str,
query_vector: Option<Vec<f32>>,
limit: u32,
note_kind: Option<&str>,
include_superseded: bool,
tags_any: &[String],
properties_filter: Option<&serde_json::Value>,
) -> RuntimeResult<Vec<NoteSearchHit>> {
self.search_notes_with_text_mode(
token,
query_text,
query_vector,
limit,
note_kind,
include_superseded,
tags_any,
properties_filter,
TextQueryMode::Plain,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn search_notes_with_text_mode(
&self,
token: &NamespaceToken,
query_text: &str,
query_vector: Option<Vec<f32>>,
limit: u32,
note_kind: Option<&str>,
include_superseded: bool,
tags_any: &[String],
properties_filter: Option<&serde_json::Value>,
text_mode: TextQueryMode,
) -> RuntimeResult<Vec<NoteSearchHit>> {
let (hits, _vector_error) = self
.search_notes_inner(
token,
query_text,
query_vector,
limit,
note_kind,
include_superseded,
tags_any,
properties_filter,
text_mode,
false,
)
.await?;
Ok(hits)
}
#[allow(clippy::too_many_arguments)]
pub async fn search_notes_outcome(
&self,
token: &NamespaceToken,
query_text: &str,
limit: u32,
note_kind: Option<&str>,
include_superseded: bool,
tags_any: &[String],
properties_filter: Option<&serde_json::Value>,
) -> RuntimeResult<NoteSearchOutcome> {
self.search_notes_outcome_with_text_mode(
token,
query_text,
limit,
note_kind,
include_superseded,
tags_any,
properties_filter,
TextQueryMode::Plain,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn search_notes_outcome_with_text_mode(
&self,
token: &NamespaceToken,
query_text: &str,
limit: u32,
note_kind: Option<&str>,
include_superseded: bool,
tags_any: &[String],
properties_filter: Option<&serde_json::Value>,
text_mode: TextQueryMode,
) -> RuntimeResult<NoteSearchOutcome> {
let (hits, vector_error) = self
.search_notes_inner(
token,
query_text,
None,
limit,
note_kind,
include_superseded,
tags_any,
properties_filter,
text_mode,
true,
)
.await?;
Ok(NoteSearchOutcome { hits, vector_error })
}
#[allow(clippy::too_many_arguments)]
async fn search_notes_inner(
&self,
token: &NamespaceToken,
query_text: &str,
query_vector: Option<Vec<f32>>,
limit: u32,
note_kind: Option<&str>,
include_superseded: bool,
tags_any: &[String],
properties_filter: Option<&serde_json::Value>,
text_mode: TextQueryMode,
tolerate_vector_error: bool,
) -> RuntimeResult<(Vec<NoteSearchHit>, Option<String>)> {
const RRF_K: usize = 60;
let candidates = limit.saturating_mul(4).max(limit);
let visible_ns: Vec<String> = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_owned())
.collect();
#[cfg(any(test, feature = "fault-injection"))]
let fts_search_inject = {
let mut g = FTS_SEARCH_FAIL_NS.lock().unwrap();
match g.as_deref() {
Some(armed) if visible_ns.iter().any(|ns| ns == armed) => {
*g = None;
true
}
_ => false,
}
};
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_search_inject = false;
let text_store = self.text_for_notes(token)?;
let text_fut = async {
if fts_search_inject {
return Err(khive_storage::StorageError::Timeout {
operation: "fts_search".into(),
});
}
text_store
.search(TextSearchRequest {
query: query_text.to_string(),
mode: text_mode,
filter: Some(TextFilter {
namespaces: visible_ns.clone(),
record_kinds: note_kind
.map(|kind| vec![kind.to_string()])
.unwrap_or_default(),
..TextFilter::default()
}),
top_k: candidates,
snippet_chars: 200,
})
.await
};
let text_fut = crate::stage_seam::text_stage(text_fut);
let vector_fut = async {
if query_vector.is_some() || self.config().embedding_model.is_some() {
self.note_search_vector_search(token, query_vector, query_text, candidates)
.await
} else {
Ok(vec![])
}
};
let (text_search_result, vector_result) = tokio::join!(text_fut, vector_fut);
let text_hits = crate::error::fts_text_leg_or_err(
text_search_result.map_err(RuntimeError::from),
"search_notes",
query_text,
)?;
let mut vector_error: Option<String> = None;
let vector_hits = match vector_result {
Ok(hits) => hits,
Err(e) if tolerate_vector_error => {
vector_error = Some(e.to_string());
Vec::new()
}
Err(e) => return Err(e),
};
let fuse_k = text_hits.len() + vector_hits.len();
let fused = crate::fusion::rrf_fuse_k(self, text_hits, vector_hits, RRF_K, fuse_k).await?;
let candidate_ids: Vec<Uuid> = fused.iter().map(|hit| hit.entity_id).collect();
if candidate_ids.is_empty() {
return Ok((vec![], vector_error));
}
let note_store = self.notes(token)?;
let search_pool = self.backend().pool_arc();
let mailbox_view = crate::MailboxView {
actor_id: token.actor().id.clone(),
delegated: false,
};
let mut alive_notes: HashMap<Uuid, Note> = HashMap::new();
for note in note_store.get_notes_batch(&candidate_ids).await? {
search_pool.record_note_candidate_hydration_row();
if note.deleted_at.is_some() {
continue;
}
if !mailbox_view.permits_message_note(token, ¬e) {
continue;
}
if let Some(want_kind) = note_kind {
if note.kind != want_kind {
continue;
}
}
if !tags_any.is_empty() {
let note_tags: Vec<String> = note
.properties
.as_ref()
.and_then(|p| p.get("tags"))
.and_then(serde_json::Value::as_array)
.map(|arr| {
arr.iter()
.filter_map(serde_json::Value::as_str)
.map(str::to_owned)
.collect()
})
.unwrap_or_default();
if !note_tags
.iter()
.any(|t| tags_any.iter().any(|w| t.eq_ignore_ascii_case(w)))
{
continue;
}
}
if let Some(pf) = properties_filter {
if !note_props_match(note.properties.as_ref(), pf) {
continue;
}
}
alive_notes.insert(note.id, note);
}
if !include_superseded && !alive_notes.is_empty() {
let graph = self.graph(token)?;
let note_ids: Vec<Uuid> = alive_notes.keys().copied().collect();
let superseded: std::collections::HashSet<Uuid> = graph
.batch_neighbors(
¬e_ids,
NeighborQuery {
direction: Direction::In,
relations: Some(vec![EdgeRelation::Supersedes]),
limit: Some(1),
min_weight: None,
},
)
.await?
.into_iter()
.map(|(note_id, _)| note_id)
.collect();
alive_notes.retain(|id, _| !superseded.contains(id));
}
let mut hits: Vec<NoteSearchHit> = fused
.into_iter()
.filter_map(|hit| {
let note = alive_notes.get(&hit.entity_id)?;
let weighted = salience_weighted_rank(hit.score, note.salience);
Some(NoteSearchHit {
note_id: hit.entity_id,
score: weighted,
rank_score_kind: hit.rank_score_kind,
signals: hit.signals,
source: hit.source,
title: hit.title.or_else(|| note_title(note)),
snippet: hit.snippet.or_else(|| note_snippet(note)),
})
})
.collect();
hits.sort_by(|a, b| b.score.cmp(&a.score).then(a.note_id.cmp(&b.note_id)));
hits.truncate(limit as usize);
Ok((hits, vector_error))
}
pub async fn resolve_prefix(
&self,
token: &NamespaceToken,
prefix: &str,
) -> RuntimeResult<Option<Uuid>> {
let namespaces = [token.namespace().as_str().to_owned()];
self.resolve_prefix_inner(Some(&namespaces), prefix, false, false)
.await
}
pub async fn resolve_prefix_including_deleted(
&self,
token: &NamespaceToken,
prefix: &str,
) -> RuntimeResult<Option<Uuid>> {
let namespaces = [token.namespace().as_str().to_owned()];
self.resolve_prefix_inner(Some(&namespaces), prefix, true, false)
.await
}
pub async fn resolve_prefix_unfiltered(&self, prefix: &str) -> RuntimeResult<Option<Uuid>> {
self.resolve_prefix_inner(None, prefix, false, false).await
}
pub async fn resolve_prefix_unfiltered_including_deleted(
&self,
prefix: &str,
) -> RuntimeResult<Option<Uuid>> {
self.resolve_prefix_inner(None, prefix, true, false).await
}
pub(crate) async fn resolve_prefix_for_kg_read(
&self,
prefix: &str,
include_deleted: bool,
) -> RuntimeResult<Option<Uuid>> {
self.resolve_prefix_inner(None, prefix, include_deleted, true)
.await
}
async fn resolve_prefix_inner(
&self,
namespaces: Option<&[String]>,
prefix: &str,
include_deleted: bool,
require_tables: bool,
) -> RuntimeResult<Option<Uuid>> {
if !prefix.chars().all(|c| c.is_ascii_hexdigit() || c == '-') {
return Ok(None);
}
#[cfg(any(test, feature = "fault-injection"))]
if consume_fault(&PREFIX_RESOLVE_FAIL_NS, prefix) {
return Err(RuntimeError::Storage(
khive_storage::StorageError::Timeout {
operation: "resolve_prefix".into(),
},
));
}
let Some((lower, upper)) = uuid_prefix_bounds(prefix) else {
return Ok(None);
};
let tables = [
("entities", true),
("notes", true),
("events", false),
("graph_edges", false),
];
let mut matches: Vec<String> = Vec::new();
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
let mut reader = self.sql().reader().await.map_err(RuntimeError::Storage)?;
for (table, has_deleted_at) in tables {
let sql = resolve_prefix_statement(
table,
has_deleted_at,
include_deleted,
namespaces,
&lower,
&upper,
);
match reader.query_all(sql).await {
Ok(rows) => {
for row in rows {
if let Some(col) = row.columns.first() {
if let SqlValue::Text(s) = &col.value {
if seen.insert(s.clone()) {
matches.push(s.clone());
}
}
}
}
}
Err(e) => {
let msg = e.to_string();
if !require_tables && msg.contains("no such table") {
continue;
}
return Err(RuntimeError::Storage(e));
}
}
if matches.len() > 1 {
break;
}
}
if matches.len() <= 1 {
if let Some(sidecar_sql) = self.events_sidecar_sql_read_only()? {
let mut sql = resolve_prefix_statement(
"events",
false,
include_deleted,
namespaces,
&lower,
&upper,
);
sql.label = Some("resolve_prefix.events_sidecar".into());
let mut sidecar_reader =
sidecar_sql.reader().await.map_err(RuntimeError::Storage)?;
match sidecar_reader.query_all(sql).await {
Ok(rows) => {
for row in rows {
if let Some(col) = row.columns.first() {
if let SqlValue::Text(s) = &col.value {
if seen.insert(s.clone()) {
matches.push(s.clone());
}
}
}
}
}
Err(e) => {
let msg = e.to_string();
if require_tables || !msg.contains("no such table") {
return Err(RuntimeError::Storage(e));
}
}
}
}
}
match matches.len() {
0 => Ok(None),
1 => {
let uuid = Uuid::from_str(&matches[0])
.map_err(|e| RuntimeError::Internal(format!("stored UUID is invalid: {e}")))?;
Ok(Some(uuid))
}
_ => {
let uuids: Vec<uuid::Uuid> = matches
.iter()
.filter_map(|s| Uuid::from_str(s).ok())
.collect();
Err(RuntimeError::AmbiguousPrefix {
prefix: prefix.to_string(),
matches: uuids,
})
}
}
}
pub async fn resolve_by_id(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
if let Some(entity) = self.entities(token)?.get_entity(id).await? {
return Ok(Some(Resolved::Entity(entity)));
}
if let Some(note) = self.notes(token)?.get_note(id).await? {
return Ok(Some(Resolved::Note(note)));
}
Ok(None)
}
pub async fn resolve_by_id_including_deleted(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
if let Some(entity) = self
.entities(token)?
.get_entity_including_deleted(id)
.await?
{
return Ok(Some(Resolved::Entity(entity)));
}
if let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? {
return Ok(Some(Resolved::Note(note)));
}
Ok(None)
}
pub async fn resolve(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
match self.get_entity(token, id).await {
Ok(entity) => return Ok(Some(Resolved::Entity(entity))),
Err(RuntimeError::NotFound(_) | RuntimeError::NamespaceMismatch { .. }) => {}
Err(e) => return Err(e),
}
if let Some(note) = self.notes(token)?.get_note(id).await? {
if Self::ensure_namespace_visible(¬e.namespace, token).is_ok() {
return Ok(Some(Resolved::Note(note)));
}
}
if let Some(event) = self.events(token)?.get_event(id).await? {
if Self::ensure_namespace_visible(&event.namespace, token).is_ok() {
return Ok(Some(Resolved::Event(event)));
}
}
Ok(None)
}
pub async fn resolve_edge_endpoint(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
if let Some(resolved) = self.resolve_by_id(token, id).await? {
return Ok(Some(resolved));
}
if let Some(event) = self.events(token)?.get_event(id).await? {
return Ok(Some(Resolved::Event(event)));
}
Ok(None)
}
pub async fn resolve_primary(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
let ns = token.namespace().as_str();
if let Some(entity) = self.entities(token)?.get_entity(id).await? {
if Self::ensure_namespace(&entity.namespace, ns).is_ok() {
return Ok(Some(Resolved::Entity(entity)));
}
}
if let Some(note) = self.notes(token)?.get_note(id).await? {
if Self::ensure_namespace(¬e.namespace, ns).is_ok() {
return Ok(Some(Resolved::Note(note)));
}
}
if let Some(event) = self.events(token)?.get_event(id).await? {
if Self::ensure_namespace(&event.namespace, ns).is_ok() {
return Ok(Some(Resolved::Event(event)));
}
}
Ok(None)
}
pub async fn resolve_including_deleted(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<Resolved>> {
let ns = token.namespace().as_str();
if let Some(entity) = self
.entities(token)?
.get_entity_including_deleted(id)
.await?
{
if Self::ensure_namespace(&entity.namespace, ns).is_ok() {
return Ok(Some(Resolved::Entity(entity)));
}
}
if let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? {
if Self::ensure_namespace(¬e.namespace, ns).is_ok() {
return Ok(Some(Resolved::Note(note)));
}
}
if let Some(event) = self.events(token)?.get_event(id).await? {
if Self::ensure_namespace(&event.namespace, ns).is_ok() {
return Ok(Some(Resolved::Event(event)));
}
}
Ok(None)
}
async fn atomic_hard_delete_with_edge_purge(
&self,
row_statement: SqlStatement,
node_id: Uuid,
namespace: &str,
actor: &str,
substrate: SubstrateKind,
) -> RuntimeResult<bool> {
let mut statements = vec![PlanStatement {
statement: row_statement,
guard: Some(AffectedRowGuard::exactly(1)),
}];
if matches!(substrate, SubstrateKind::Entity | SubstrateKind::Note) {
statements.push(PlanStatement {
statement: khive_db::stores::attachment::delete_record_attachments_statement(
node_id,
if substrate == SubstrateKind::Entity {
AttachmentSubstrate::Entity
} else {
AttachmentSubstrate::Note
},
),
guard: None,
});
}
statements.extend(
hard_delete_lineage_warning_statements(namespace, actor, node_id, substrate)
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
}),
);
statements.push(PlanStatement {
statement: purge_incident_edges_statement(node_id),
guard: None,
});
let plan = AtomicOpPlan::Delete(DeletePlan {
target_id: node_id,
statements,
post_commit: PostCommitEffect::None,
});
match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
Ok(AtomicRunOutcome::Committed { .. }) => Ok(true),
Ok(AtomicRunOutcome::RolledBack {
failure: AtomicOpFailure::NoteConflict(conflict),
..
}) => Err(conflict.into_error().into()),
Ok(AtomicRunOutcome::RolledBack {
failure: AtomicOpFailure::EntityConflict(conflict),
..
}) => Err(conflict.into_error().into()),
Ok(AtomicRunOutcome::RolledBack {
failure: AtomicOpFailure::GuardFailed { .. },
..
}) => Ok(false),
Ok(AtomicRunOutcome::RolledBack {
failure: AtomicOpFailure::SqlError { message, .. },
..
}) => Err(RuntimeError::Internal(format!(
"hard delete + edge purge for {node_id} failed: {message}"
))),
Err(e) => Err(RuntimeError::Internal(format!(
"hard delete + edge purge for {node_id}: atomic unit seam failure: {}",
e.0
))),
}
}
pub async fn restore_entity(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<(Entity, bool)>> {
let Some(entity) = self
.entities(token)?
.get_entity_including_deleted(id)
.await?
else {
return Ok(None);
};
if entity.namespace != token.namespace().as_str() {
return Ok(None);
}
if let Some(kept_id) = entity.merged_into {
if entity.deleted_at.is_none() {
return Err(live_merged_entity_refused(id, kept_id));
}
return Err(merge_tombstone_restore_refused(id, kept_id));
}
if entity.deleted_at.is_none() {
return Ok(Some((entity, false)));
}
let updated_at =
Utc::now()
.timestamp_micros()
.max(entity.updated_at.checked_add(1).ok_or_else(|| {
RuntimeError::Internal(format!(
"entity {id} updated_at is already at i64::MAX and cannot advance"
))
})?);
let mut restored = entity;
restored.deleted_at = None;
restored.updated_at = updated_at;
restored.version = restored
.version
.checked_add(1)
.ok_or_else(|| RuntimeError::InvalidInput("entity version overflow".into()))?;
let mut statements = vec![PlanStatement {
statement: SqlStatement {
sql: "UPDATE entities SET deleted_at=NULL, updated_at=?1, version=version+1 \
WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL AND version=?4"
.into(),
params: vec![
SqlValue::Integer(updated_at),
SqlValue::Text(id.to_string()),
SqlValue::Text(token.namespace().as_str().to_owned()),
SqlValue::Integer(restored.version - 1),
],
label: Some("entity-restore".into()),
},
guard: Some(AffectedRowGuard::exactly(1)),
}];
for statement in khive_db::stores::text::delete_document_statements(
"fts_entities",
&restored.namespace,
id,
)
.into_iter()
.chain(insert_document_statements(
"fts_entities",
&entity_fts_document(&restored),
)) {
statements.push(PlanStatement {
statement,
guard: None,
});
}
let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
graph_effects: Vec::new(),
target_id: id,
statements,
post_commit: PostCommitEffect::None,
edge_natural_key: None,
idempotent_noop: false,
entity_guard: None,
note_guard: None,
note_vector_purge: None,
note_embedding_inheritance: None,
}));
match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
Ok(AtomicRunOutcome::Committed { .. }) => {
#[cfg(any(test, feature = "fault-injection"))]
if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
return Err(restore_reindex_failed(
"entity",
id,
RuntimeError::Internal("injected FTS failure".to_string()),
));
}
self.reindex_entity(token, &restored)
.await
.map_err(|e| restore_reindex_failed("entity", id, e))?;
Ok(Some((restored, true)))
}
Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
format!("entity restore rolled back: {failure:?}"),
)),
Err(error) => Err(RuntimeError::Storage(error.0)),
}
}
pub async fn restore_note(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<(Note, bool)>> {
let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? else {
return Ok(None);
};
if note.namespace != token.namespace().as_str() {
return Ok(None);
}
if note.deleted_at.is_none() {
return Ok(Some((note, false)));
}
if let Some(key) = note.key.as_deref() {
if let Some(holder) = self
.notes(token)?
.get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
.await?
.into_iter()
.find(|holder| holder.id != note.id)
{
return Err(restore_key_conflict(key, &holder));
}
}
let updated_at =
Utc::now()
.timestamp_micros()
.max(note.updated_at.checked_add(1).ok_or_else(|| {
RuntimeError::Internal(format!(
"note {id} updated_at is already at i64::MAX and cannot advance"
))
})?);
let mut params = vec![
SqlValue::Text("active".into()),
SqlValue::Integer(updated_at),
SqlValue::Text(id.to_string()),
SqlValue::Text(note.namespace.clone()),
SqlValue::Text(note.kind.clone()),
];
let key_clause = if let Some(key) = note.key.as_deref() {
params.push(SqlValue::Text(key.to_owned()));
format!(
" AND (key IS NULL OR NOT EXISTS (SELECT 1 FROM notes live \
WHERE live.namespace=?4 AND live.kind=?5 AND live.key=?{} \
AND live.deleted_at IS NULL AND live.id != notes.id))",
params.len()
)
} else {
String::new()
};
let mut restored = note.clone();
restored.status = "active".into();
restored.deleted_at = None;
restored.updated_at = updated_at;
restored.version = restored
.version
.checked_add(1)
.ok_or_else(|| RuntimeError::Internal(format!("note {id} version is exhausted")))?;
let mut statements = vec![PlanStatement {
statement: SqlStatement {
sql: format!(
"UPDATE notes SET status=?1, deleted_at=NULL, updated_at=?2 \
WHERE id=?3 AND namespace=?4 AND kind=?5 AND deleted_at IS NOT NULL{key_clause}"
),
params,
label: Some("note-restore".into()),
},
guard: Some(AffectedRowGuard::exactly(1)),
}];
for statement in
khive_db::stores::text::delete_document_statements("fts_notes", &restored.namespace, id)
.into_iter()
.chain(insert_document_statements(
"fts_notes",
¬e_fts_document(&restored),
))
{
statements.push(PlanStatement {
statement,
guard: None,
});
}
let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
graph_effects: Vec::new(),
target_id: id,
statements,
post_commit: PostCommitEffect::None,
edge_natural_key: None,
idempotent_noop: false,
entity_guard: None,
note_guard: None,
note_vector_purge: None,
note_embedding_inheritance: None,
}));
match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
Ok(AtomicRunOutcome::Committed { .. }) => {
#[cfg(any(test, feature = "fault-injection"))]
if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
return Err(restore_reindex_failed(
"note",
id,
RuntimeError::Internal("injected FTS failure".to_string()),
));
}
let reindexed = self.reindex_note_with_report(token, &restored).await;
let report = reindexed.map_err(|e| restore_reindex_failed("note", id, e))?;
let degradations = report.post_commit_degradations();
legacy_post_commit_result("restore_note", id, Some((restored, true)), degradations)
}
Ok(AtomicRunOutcome::RolledBack {
failure: AtomicOpFailure::GuardFailed { .. },
..
}) => {
if let Some(key) = note.key.as_deref() {
if let Some(holder) = self
.notes(token)?
.get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
.await?
.into_iter()
.find(|holder| holder.id != note.id)
{
return Err(restore_key_conflict(key, &holder));
}
}
Err(RuntimeError::NotFound(format!(
"note {id} is no longer a caller-owned tombstone"
)))
}
Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
format!("note restore rolled back: {failure:?}"),
)),
Err(error) => Err(RuntimeError::Storage(error.0)),
}
}
pub async fn restore_edge(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<Option<(Edge, bool)>> {
let Some(edge) = self.get_edge_including_deleted(token, id).await? else {
return Ok(None);
};
if edge.namespace != token.namespace().as_str() {
return Ok(None);
}
if edge.deleted_at.is_none() {
return Ok(Some((edge, false)));
}
let updated_at = Utc::now();
let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
graph_effects: Vec::new(),
target_id: id,
statements: vec![PlanStatement {
statement: SqlStatement {
sql: "UPDATE graph_edges SET deleted_at=NULL, updated_at=?1 \
WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL"
.into(),
params: vec![
SqlValue::Integer(updated_at.timestamp_micros()),
SqlValue::Text(id.to_string()),
SqlValue::Text(edge.namespace.clone()),
],
label: Some("edge-restore".into()),
},
guard: Some(AffectedRowGuard::exactly(1)),
}],
post_commit: PostCommitEffect::None,
edge_natural_key: None,
idempotent_noop: false,
entity_guard: None,
note_guard: None,
note_vector_purge: None,
note_embedding_inheritance: None,
}));
match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
Ok(AtomicRunOutcome::Committed { .. }) => {
let mut restored = edge;
restored.deleted_at = None;
restored.updated_at = updated_at;
Ok(Some((restored, true)))
}
Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
format!("edge restore rolled back: {failure:?}"),
)),
Err(error) => Err(RuntimeError::Storage(error.0)),
}
}
pub async fn delete_note(
&self,
token: &NamespaceToken,
id: Uuid,
hard: bool,
) -> RuntimeResult<bool> {
let (deleted, degradations) = self
.delete_note_with_post_commit_report(token, id, hard)
.await?;
legacy_post_commit_result("delete_note", id, deleted, degradations)
}
pub async fn delete_note_with_post_commit_report(
&self,
token: &NamespaceToken,
id: Uuid,
hard: bool,
) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
let note_store = self.notes(token)?;
let note = if hard {
match note_store.get_note_including_deleted(id).await? {
Some(n) => n,
None => return Ok((false, Vec::new())),
}
} else {
match note_store.get_note(id).await? {
Some(n) => n,
None => return Ok((false, Vec::new())),
}
};
if let Some(error) = self.stream_member_error(¬e).await? {
return Err(error);
}
let mode = if hard {
DeleteMode::Hard
} else {
DeleteMode::Soft
};
let record_tok = token.with_namespace(
khive_types::Namespace::parse(¬e.namespace)
.map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
);
let record_ns = note.namespace.clone();
let actor = format!("{}:{}", token.actor().kind, token.actor().id);
let deleted = if hard {
self.atomic_hard_delete_with_edge_purge(
note_hard_delete_statement(id),
id,
&record_ns,
&actor,
SubstrateKind::Note,
)
.await?
} else {
note_store.delete_note(id, mode).await?
};
let mut degradations = Vec::new();
if deleted {
let fts_result = match self.text_for_notes(&record_tok) {
Ok(store) => store
.delete_document(&record_ns, id)
.await
.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = fts_result {
record_post_commit_degradation(
&mut degradations,
"delete_note",
id,
"fts_cleanup",
error,
);
}
for model_name in self.registered_embedding_model_names() {
let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = vector_result {
record_post_commit_degradation(
&mut degradations,
"delete_note",
id,
"vector_cleanup",
format!("{model_name}: {error}"),
);
}
}
let event = khive_storage::event::Event::new(
record_ns.clone(),
"delete",
EventKind::NoteDeleted,
SubstrateKind::Note,
"",
)
.with_target(id)
.with_payload(serde_json::json!({"id": id, "namespace": record_ns, "hard": hard}));
let event_result = match self.events(&record_tok) {
Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = event_result {
record_post_commit_degradation(
&mut degradations,
"delete_note",
id,
"event_append",
error,
);
}
self.fire_note_mutation_hook(¬e.kind, id).await;
}
Ok((deleted, degradations))
}
pub async fn delete_note_row_first_for_compensation(
&self,
token: &NamespaceToken,
id: Uuid,
) -> RuntimeResult<()> {
let note_store = self.notes(token)?;
let Some(note) = note_store.get_note_including_deleted(id).await? else {
return Ok(());
};
let record_tok = NamespaceToken::for_namespace(
khive_types::Namespace::parse(¬e.namespace)
.map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
);
let record_ns = note.namespace.clone();
note_store.delete_note(id, DeleteMode::Hard).await?;
#[cfg(any(test, feature = "fault-injection"))]
{
let armed = ROLLBACK_CLEANUP_FAIL_NS.lock().unwrap().take();
if armed.as_deref() == Some(record_ns.as_str()) {
return Err(RuntimeError::Internal(
"row removed but compensation cleanup failed: injected=true".to_string(),
));
}
}
let mut cleanup_errors = Vec::new();
if let Err(e) = self.graph(&record_tok)?.purge_incident_edges(id).await {
cleanup_errors.push(format!("graph={e}"));
}
if let Err(e) = self
.text_for_notes(&record_tok)?
.delete_document(&record_ns, id)
.await
{
cleanup_errors.push(format!("fts={e}"));
}
for model_name in self.registered_embedding_model_names() {
if let Err(e) = self
.vectors_for_model(&record_tok, &model_name)?
.delete(id)
.await
{
cleanup_errors.push(format!("vector[{model_name}]={e}"));
}
}
if cleanup_errors.is_empty() {
Ok(())
} else {
Err(RuntimeError::Internal(format!(
"row removed but compensation cleanup failed: {}",
cleanup_errors.join("; ")
)))
}
}
}
#[derive(Clone, Debug, Serialize)]
pub struct QueryResult {
pub rows: Vec<SqlRow>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub warnings: Vec<String>,
pub offset: usize,
pub page_size: usize,
pub has_more: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub next_offset: Option<usize>,
pub truncated: bool,
}
#[derive(Debug)]
enum SymmetricEdgeUpdateOutcome {
Absorbed(String),
Updated,
Stale,
}
impl KhiveRuntime {
pub async fn query(&self, token: &NamespaceToken, query: &str) -> RuntimeResult<Vec<SqlRow>> {
Ok(self
.query_with_metadata(token, query, khive_query::CompileOptions::default())
.await?
.rows)
}
pub async fn query_with_metadata(
&self,
token: &NamespaceToken,
query: &str,
mut opts: khive_query::CompileOptions,
) -> RuntimeResult<QueryResult> {
use khive_query::QueryValue;
use khive_storage::types::SqlValue;
let (language, ast) = khive_query::language::parse_auto_with_language(query)?;
if opts.max_limit == 0 {
return Err(RuntimeError::InvalidInput(
"query page size must be at least 1".into(),
));
}
let offset = ast.offset;
let page_size = ast.limit.unwrap_or(opts.max_limit).min(opts.max_limit);
opts.scopes = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_string())
.collect();
let compiled = khive_query::compile(&ast, &opts)?;
let mut warnings = compiled.warnings;
let truncation_check = compiled.truncation_check;
warnings.extend(self.with_pack_edge_rules(|pack_rules| {
static_impossible_edge_pattern_warnings(language, &ast.pattern, pack_rules)
}));
let params: Vec<SqlValue> = compiled
.params
.into_iter()
.map(|qv| match qv {
QueryValue::Null => SqlValue::Null,
QueryValue::Integer(n) => SqlValue::Integer(n),
QueryValue::Float(f) => SqlValue::Float(f),
QueryValue::Text(s) => SqlValue::Text(s),
QueryValue::Blob(b) => SqlValue::Blob(b),
})
.collect();
let mut reader = self.sql().reader().await?;
let stmt = SqlStatement {
sql: compiled.sql,
params,
label: None,
};
let mut rows = reader.query_all(stmt).await?;
let mut truncated = false;
if let Some(check) = truncation_check {
if rows.len() > check.max_limit {
rows.truncate(check.max_limit);
truncated = true;
}
}
let next_offset = if truncated && language == khive_query::QueryLanguage::Gql {
let next = offset.checked_add(rows.len()).ok_or_else(|| {
RuntimeError::InvalidInput("GQL next_offset exceeds usize::MAX".into())
})?;
if next == offset {
return Err(RuntimeError::InvalidInput(
"query page did not advance; page size must be at least 1".into(),
));
}
i64::try_from(next).map_err(|_| {
RuntimeError::InvalidInput("GQL next_offset exceeds i64::MAX".into())
})?;
Some(next)
} else {
None
};
if truncated {
let Some(check) = truncation_check else {
return Err(RuntimeError::Internal(
"truncated query result is missing sentinel metadata".into(),
));
};
let bound = match check.requested_limit {
Some(requested) => {
format!("requested query LIMIT {requested} exceeds the effective page size")
}
None => "the query has no explicit LIMIT".to_string(),
};
let warning = match language {
khive_query::QueryLanguage::Gql => {
let Some(next) = next_offset else {
return Err(RuntimeError::Internal(
"truncated GQL result is missing its continuation offset".into(),
));
};
format!(
"result page capped at {} rows because {bound}; more matches exist. \
Continue the same GQL query with `SKIP {next}` (the machine-readable \
`next_offset`) and keep the same page size.",
check.max_limit
)
}
khive_query::QueryLanguage::Sparql => format!(
"result page capped at {} rows because {bound}; more matches exist. \
SPARQL OFFSET paging is not part of the supported dialect.",
check.max_limit
),
};
warnings.push(warning);
}
Ok(QueryResult {
rows,
warnings,
offset,
page_size,
has_more: truncated,
next_offset,
truncated,
})
}
pub async fn delete_entity(
&self,
token: &NamespaceToken,
id: Uuid,
hard: bool,
) -> RuntimeResult<bool> {
let (deleted, degradations) = self
.delete_entity_with_post_commit_report(token, id, hard)
.await?;
legacy_post_commit_result("delete_entity", id, deleted, degradations)
}
pub async fn delete_entity_with_post_commit_report(
&self,
token: &NamespaceToken,
id: Uuid,
hard: bool,
) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
let entity = if hard {
match self
.entities(token)?
.get_entity_including_deleted(id)
.await?
{
Some(e) => e,
None => return Ok((false, Vec::new())),
}
} else {
match self.entities(token)?.get_entity(id).await? {
Some(e) => e,
None => return Ok((false, Vec::new())),
}
};
let mode = if hard {
DeleteMode::Hard
} else {
DeleteMode::Soft
};
let record_tok = token.with_namespace(
khive_types::Namespace::parse(&entity.namespace)
.map_err(|e| RuntimeError::Internal(format!("entity namespace invalid: {e}")))?,
);
let actor = format!("{}:{}", token.actor().kind, token.actor().id);
let deleted = if hard {
self.atomic_hard_delete_with_edge_purge(
entity_hard_delete_statement(id),
id,
&entity.namespace,
&actor,
SubstrateKind::Entity,
)
.await?
} else {
self.entities(token)?.delete_entity(id, mode).await?
};
let mut degradations = Vec::new();
if deleted {
let ns = entity.namespace.clone();
let fts_result = match self.text(&record_tok) {
Ok(store) => store
.delete_document(&ns, id)
.await
.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = fts_result {
record_post_commit_degradation(
&mut degradations,
"delete_entity",
id,
"fts_cleanup",
error,
);
}
for model_name in self.registered_embedding_model_names() {
let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = vector_result {
record_post_commit_degradation(
&mut degradations,
"delete_entity",
id,
"vector_cleanup",
format!("{model_name}: {error}"),
);
}
}
let event = khive_storage::event::Event::new(
ns.clone(),
"delete",
EventKind::EntityDeleted,
SubstrateKind::Entity,
"",
)
.with_target(id)
.with_payload(serde_json::json!({"id": id, "namespace": ns, "hard": hard}));
let event_result = match self.events(&record_tok) {
Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
Err(error) => Err(error),
};
if let Err(error) = event_result {
record_post_commit_degradation(
&mut degradations,
"delete_entity",
id,
"event_append",
error,
);
}
}
Ok((deleted, degradations))
}
pub(crate) async fn delete_entity_attachments_on_core(&self, id: Uuid) -> RuntimeResult<bool> {
let core = self.core();
drop(core.attachments()?);
let statement = khive_db::stores::attachment::delete_record_attachments_statement(
id,
AttachmentSubstrate::Entity,
);
Ok(core.sql().writer().await?.execute(statement).await? > 0)
}
pub async fn count_entities(
&self,
token: &NamespaceToken,
kind: Option<&str>,
) -> RuntimeResult<u64> {
let ns_strs: Vec<String> = token
.visible_namespaces()
.iter()
.map(|ns| ns.as_str().to_owned())
.collect();
let filter = EntityFilter {
kinds: match kind {
Some(k) => vec![k.to_string()],
None => vec![],
},
namespaces: ns_strs,
..Default::default()
};
Ok(self
.entities(token)?
.count_entities(token.namespace().as_str(), filter)
.await?)
}
pub async fn entity_stats_counts(
&self,
token: &NamespaceToken,
) -> RuntimeResult<EntityStatsCounts> {
entity_stats_counts(self.entities(token)?.as_ref(), token).await
}
pub async fn get_edge(
&self,
_token: &NamespaceToken,
edge_id: Uuid,
) -> RuntimeResult<Option<Edge>> {
let mut reader = self.sql().reader().await?;
let record_ns = reader
.query_scalar(SqlStatement {
sql: "SELECT namespace FROM graph_edges \
WHERE id = ?1 AND deleted_at IS NULL LIMIT 1"
.into(),
params: vec![SqlValue::Text(edge_id.to_string())],
label: Some("get_edge_namespace".into()),
})
.await?;
let Some(SqlValue::Text(record_ns)) = record_ns else {
return Ok(None);
};
let record_tok = NamespaceToken::for_namespace(
khive_types::Namespace::parse(&record_ns)
.map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
);
Ok(self
.graph(&record_tok)?
.get_edge(LinkId::from(edge_id))
.await?)
}
pub async fn get_edges_by_id(
&self,
_token: &NamespaceToken,
ids: &[Uuid],
) -> RuntimeResult<Vec<Option<Edge>>> {
let mut edges = Vec::with_capacity(ids.len());
for chunk in ids.chunks(900) {
let window = self.prepare_edge_read_window(chunk).await?;
edges.extend(
Self::hydrate_edge_read_window(chunk, window, |record_token| {
self.graph(record_token)
})
.await?,
);
}
Ok(edges)
}
async fn prepare_edge_read_window(&self, ids: &[Uuid]) -> RuntimeResult<EdgeReadWindow> {
let placeholders = (1..=ids.len())
.map(|index| format!("?{index}"))
.collect::<Vec<_>>()
.join(",");
let mut reader = self.sql().reader().await?;
let rows = reader
.query_all(SqlStatement {
sql: format!(
"SELECT id, namespace FROM graph_edges WHERE id IN ({placeholders}) AND deleted_at IS NULL"
),
params: ids.iter().map(|id| SqlValue::Text(id.to_string())).collect(),
label: Some("get_edge_namespace".into()),
})
.await?;
let mut namespaces = HashMap::with_capacity(rows.len());
for row in rows {
let Some(SqlValue::Text(id)) = row.columns.first().map(|column| &column.value) else {
return Err(RuntimeError::Internal(
"edge namespace lookup returned an invalid id".into(),
));
};
let id = Uuid::parse_str(id).map_err(|e| {
RuntimeError::Internal(format!("edge namespace lookup returned an invalid id: {e}"))
})?;
let value = row
.columns
.get(1)
.map(|column| column.value.clone())
.unwrap_or(SqlValue::Null);
namespaces.insert(id, value);
}
let mut window = EdgeReadWindow {
outcomes: (0..ids.len()).map(|_| Some(Ok(None))).collect(),
groups: Vec::new(),
};
let mut group_indices = HashMap::new();
for (index, id) in ids.iter().enumerate() {
let Some(SqlValue::Text(record_ns)) = namespaces.get(id) else {
continue;
};
match khive_types::Namespace::parse(record_ns) {
Ok(namespace) => {
let next_group = window.groups.len();
let group = *group_indices.entry(record_ns.clone()).or_insert(next_group);
if group == next_group {
window.groups.push((namespace, Vec::new()));
}
window.groups[group].1.push(index);
window.outcomes[index] = None;
}
Err(error) => {
window.outcomes[index] = Some(Err(RuntimeError::Internal(format!(
"edge namespace invalid: {error}"
))));
}
}
}
Ok(window)
}
async fn hydrate_edge_read_window<F>(
ids: &[Uuid],
mut window: EdgeReadWindow,
mut graph: F,
) -> RuntimeResult<Vec<Option<Edge>>>
where
F: FnMut(&NamespaceToken) -> RuntimeResult<std::sync::Arc<dyn khive_storage::GraphStore>>,
{
for (namespace, indices) in window.groups {
let record_token = NamespaceToken::for_namespace(namespace);
let group_ids: Vec<LinkId> = indices
.iter()
.map(|&index| LinkId::from(ids[index]))
.collect();
let outcomes = match graph(&record_token) {
Ok(store) => store
.get_edge_read_outcomes(&group_ids)
.await
.map_err(RuntimeError::from),
Err(error) => Err(error),
};
match outcomes {
Ok(outcomes) if outcomes.len() == indices.len() => {
for (index, outcome) in indices.into_iter().zip(outcomes) {
window.outcomes[index] = Some(outcome.map_err(RuntimeError::from));
}
}
Ok(_) => {
window.outcomes[indices[0]] = Some(Err(RuntimeError::Internal(
"edge batch returned an invalid outcome count".into(),
)));
}
Err(error) => {
window.outcomes[indices[0]] = Some(Err(error));
}
}
}
window
.outcomes
.into_iter()
.map(|outcome| {
outcome.unwrap_or_else(|| {
Err(RuntimeError::Internal(
"edge batch omitted an input outcome".into(),
))
})
})
.collect()
}
pub async fn get_edge_visible(
&self,
token: &NamespaceToken,
edge_id: Uuid,
) -> RuntimeResult<Option<Edge>> {
self.get_edge(token, edge_id).await
}
pub async fn get_edge_including_deleted(
&self,
_token: &NamespaceToken,
edge_id: Uuid,
) -> RuntimeResult<Option<Edge>> {
let mut reader = self.sql().reader().await?;
let record_ns = reader
.query_scalar(SqlStatement {
sql: "SELECT namespace FROM graph_edges WHERE id = ?1 LIMIT 1".into(),
params: vec![SqlValue::Text(edge_id.to_string())],
label: Some("get_edge_including_deleted_namespace".into()),
})
.await?;
let Some(SqlValue::Text(record_ns)) = record_ns else {
return Ok(None);
};
let record_tok = NamespaceToken::for_namespace(
khive_types::Namespace::parse(&record_ns)
.map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
);
Ok(self
.graph(&record_tok)?
.get_edge_including_deleted(LinkId::from(edge_id))
.await?)
}
pub async fn get_edge_by_natural_key_including_deleted(
&self,
token: &NamespaceToken,
namespace: &str,
source_id: Uuid,
target_id: Uuid,
relation: EdgeRelation,
) -> RuntimeResult<Option<Edge>> {
Ok(self
.graph(token)?
.get_edge_by_natural_key_including_deleted(namespace, source_id, target_id, relation)
.await?)
}
pub const EDGE_LIST_MAX_LIMIT: u32 = 1000;
pub async fn list_edges(
&self,
token: &NamespaceToken,
filter: crate::curation::EdgeListFilter,
limit: u32,
offset: u32,
) -> RuntimeResult<Vec<Edge>> {
let limit = limit.min(Self::EDGE_LIST_MAX_LIMIT);
let visible = token.visible_namespaces();
if let [ns] = visible {
let temp = NamespaceToken::for_namespace(ns.clone());
let page = self
.graph(&temp)?
.query_edges(
filter.into(),
vec![SortOrder {
field: EdgeSortField::CreatedAt,
direction: khive_storage::types::SortDirection::Asc,
}],
PageRequest {
offset: offset.into(),
limit,
},
)
.await?;
return Ok(page.items);
}
let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
let sort = vec![SortOrder {
field: EdgeSortField::CreatedAt,
direction: khive_storage::types::SortDirection::Asc,
}];
let graph = self.graph(token)?;
match graph
.query_edges_in_namespaces(
&ns_strs,
filter.clone().into(),
sort.clone(),
PageRequest {
offset: offset.into(),
limit,
},
)
.await
{
Ok(page) => Ok(page.items),
Err(khive_storage::StorageError::Unsupported { operation, .. })
if operation == "query_edges_in_namespaces" =>
{
let fetch_limit = offset.saturating_add(limit);
let mut namespace_prefixes = Vec::new();
for ns in visible {
let temp = NamespaceToken::for_namespace(ns.clone());
let page = self
.graph(&temp)?
.query_edges(
filter.clone().into(),
sort.clone(),
PageRequest {
offset: 0,
limit: fetch_limit,
},
)
.await?;
namespace_prefixes.push(page.items);
}
Ok(Self::merge_paged_namespace_edges(
namespace_prefixes,
offset,
limit,
))
}
Err(error) => Err(error.into()),
}
}
fn merge_paged_namespace_edges(
namespace_prefixes: Vec<Vec<Edge>>,
offset: u32,
limit: u32,
) -> Vec<Edge> {
let mut results: Vec<Edge> = namespace_prefixes.into_iter().flatten().collect();
results.sort_by_key(|e| (e.created_at, Uuid::from(e.id)));
let start = (offset as usize).min(results.len());
let end = (start + limit as usize).min(results.len());
results[start..end].to_vec()
}
pub async fn list_edges_after(
&self,
token: &NamespaceToken,
filter: crate::curation::EdgeListFilter,
after: Option<Uuid>,
limit: u32,
) -> RuntimeResult<(Vec<Edge>, Option<Uuid>)> {
let limit = limit.clamp(1, Self::EDGE_LIST_MAX_LIMIT);
let visible = token.visible_namespaces();
let limit_usize = limit as usize;
let cursor_store = self.graph(token)?;
let after = match after {
Some(id) => {
let edge = self
.get_edge_including_deleted(token, id)
.await?
.ok_or_else(|| RuntimeError::NotFound(format!("edge cursor {id}")))?;
Self::ensure_namespace_visible(&edge.namespace, token)?;
let sequence = cursor_store.edge_sequence(id).await?.ok_or_else(|| {
RuntimeError::Internal(format!(
"edge cursor {id} has no insertion-sequence ledger row"
))
})?;
Some(SeekCursor { sequence, id })
}
None => None,
};
if let [ns] = visible {
let temp = NamespaceToken::for_namespace(ns.clone());
let page = self
.graph(&temp)?
.query_edges_sequence_after(filter.into(), after, limit)
.await?;
return Ok((page.items, page.next_after.map(|cursor| cursor.id)));
}
let probe_limit = limit.saturating_add(1);
let mut results = Vec::new();
for ns in visible {
let temp = NamespaceToken::for_namespace(ns.clone());
let page = self
.graph(&temp)?
.query_edges_sequence_after(filter.clone().into(), after, probe_limit)
.await?;
results.extend(page.items);
}
let ids = results
.iter()
.map(|edge| Uuid::from(edge.id))
.collect::<Vec<_>>();
let sequences = cursor_store
.edge_sequences(&ids)
.await?
.into_iter()
.collect::<HashMap<_, _>>();
if let Some(missing) = ids.iter().find(|id| !sequences.contains_key(id)) {
return Err(RuntimeError::Internal(format!(
"edge {missing} has no insertion-sequence ledger row"
)));
}
results.sort_by_key(|edge| {
let id = Uuid::from(edge.id);
(sequences[&id], id)
});
results.dedup_by_key(|e| Uuid::from(e.id));
let has_more = results.len() > limit_usize;
if has_more {
results.truncate(limit_usize);
}
let next_after = if has_more {
results.last().map(|e| Uuid::from(e.id))
} else {
None
};
Ok((results, next_after))
}
pub async fn count_edges_by_relation(
&self,
token: &NamespaceToken,
) -> RuntimeResult<std::collections::HashMap<String, u64>> {
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
let graph = self.graph(token)?;
let counts = match graph
.count_edges_by_relation_in_namespaces(&namespaces)
.await
{
Ok(counts) => counts,
Err(khive_storage::StorageError::Unsupported { operation, .. })
if operation == "count_edges_by_relation_in_namespaces" =>
{
let mut totals = HashMap::new();
for namespace in token.visible_namespaces() {
let scoped = NamespaceToken::for_namespace(namespace.clone());
for (relation, count) in self.graph(&scoped)?.count_edges_by_relation().await? {
*totals.entry(relation).or_insert(0) += count;
}
}
return Ok(totals
.into_iter()
.map(|(relation, count)| (relation.to_string(), count))
.collect());
}
Err(error) => return Err(error.into()),
};
Ok(counts
.into_iter()
.map(|(relation, count)| (relation.to_string(), count))
.collect())
}
pub async fn count_edges_by_endpoint_base(
&self,
token: &NamespaceToken,
) -> RuntimeResult<khive_storage::types::EdgeEndpointBaseCounts> {
use khive_storage::types::EdgeEndpointBaseCounts;
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
let graph = self.graph(token)?;
match graph
.count_edges_by_endpoint_base_in_namespaces(&namespaces)
.await
{
Ok(counts) => Ok(counts),
Err(khive_storage::StorageError::Unsupported { operation, .. })
if operation == "count_edges_by_endpoint_base_in_namespaces"
|| operation == "count_edges_by_endpoint_base" =>
{
let mut totals = EdgeEndpointBaseCounts::default();
for namespace in token.visible_namespaces() {
let scoped = NamespaceToken::for_namespace(namespace.clone());
let counts = self.graph(&scoped)?.count_edges_by_endpoint_base().await?;
totals.entity_entity =
totals.entity_entity.saturating_add(counts.entity_entity);
totals.entity_note = totals.entity_note.saturating_add(counts.entity_note);
totals.note_entity = totals.note_entity.saturating_add(counts.note_entity);
totals.note_note = totals.note_note.saturating_add(counts.note_note);
totals.unresolved = totals.unresolved.saturating_add(counts.unresolved);
}
Ok(totals)
}
Err(error) => Err(error.into()),
}
}
#[allow(clippy::too_many_arguments)]
fn update_edge_symmetric_dml(
conn: &rusqlite::Connection,
ns: &str,
edge_id_str: &str,
canon_src_str: &str,
canon_tgt_str: &str,
relation_str: &str,
weight: f64,
metadata: Option<String>,
expected_updated_at_micros: i64,
expected_deleted_at_micros: Option<i64>,
) -> Result<SymmetricEdgeUpdateOutcome, SqliteError> {
let minimum_updated_at_micros =
expected_updated_at_micros.checked_add(1).ok_or_else(|| {
SqliteError::InvalidData(format!(
"update_edge: edge {edge_id_str} updated_at is already at i64::MAX \
and cannot advance"
))
})?;
let now_ts = chrono::Utc::now()
.timestamp_micros()
.max(minimum_updated_at_micros);
let conflict_id: Option<String> = conn
.query_row(
khive_db::stores::graph::EDGE_SYMMETRIC_CONFLICT_PROBE_SQL,
rusqlite::params![
&ns,
&canon_src_str,
&canon_tgt_str,
&relation_str,
&edge_id_str
],
|row| row.get(0),
)
.optional()
.map_err(SqliteError::Rusqlite)?;
if let Some(existing_id) = conflict_id {
let affected = conn
.execute(
khive_db::stores::graph::EDGE_SYMMETRIC_DELETE_NONCANONICAL_GUARDED_SQL,
rusqlite::params![
&ns,
&edge_id_str,
expected_updated_at_micros,
expected_deleted_at_micros,
],
)
.map_err(SqliteError::Rusqlite)?;
if affected == 0 {
return Ok(SymmetricEdgeUpdateOutcome::Stale);
}
Ok(SymmetricEdgeUpdateOutcome::Absorbed(existing_id))
} else {
let affected = conn
.execute(
khive_db::stores::graph::EDGE_SYMMETRIC_UPDATE_INPLACE_SQL,
rusqlite::params![
&canon_src_str,
&canon_tgt_str,
&relation_str,
weight,
now_ts,
metadata,
&ns,
&edge_id_str,
expected_updated_at_micros,
expected_deleted_at_micros,
],
)
.map_err(SqliteError::Rusqlite)?;
if affected == 0 {
return Ok(SymmetricEdgeUpdateOutcome::Stale);
}
Ok(SymmetricEdgeUpdateOutcome::Updated)
}
}
pub async fn update_edge(
&self,
token: &NamespaceToken,
edge_id: Uuid,
patch: crate::curation::EdgePatch,
) -> RuntimeResult<Edge> {
let graph_for_fetch = self.graph(token)?;
let mut edge = graph_for_fetch
.get_edge(LinkId::from(edge_id))
.await?
.ok_or_else(|| crate::RuntimeError::NotFound(format!("edge {edge_id}")))?;
let expected_updated_at = edge.updated_at;
let expected_deleted_at = edge.deleted_at;
#[cfg(test)]
crate::curation::race_seam::pause_after_read().await;
let record_ns: String = edge.namespace.clone();
let record_tok = token.with_namespace(
khive_types::Namespace::parse(&record_ns)
.map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
);
let graph = self.graph(&record_tok)?;
let mut changed_fields: Vec<&'static str> = Vec::new();
if let Some(r) = patch.relation {
self.validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, r)
.await?;
edge.relation = r;
changed_fields.push("relation");
}
if let Some(w) = patch.weight {
if !w.is_finite() || !(0.0..=1.0).contains(&w) {
return Err(RuntimeError::InvalidInput(format!(
"edge weight must be a finite value in [0.0, 1.0]; got {w}"
)));
}
edge.weight = w;
changed_fields.push("weight");
}
if let Some(props) = patch.properties {
crate::secret_gate::reject_reserved_secret_gate_property(Some(&props))?;
edge.metadata = Some(props);
}
let (canon_src, canon_tgt) =
canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
if edge.relation.is_symmetric() {
let ns = record_ns.clone();
let edge_id_str = edge_id.to_string();
let relation_str = edge.relation.to_string();
let canon_src_str = canon_src.to_string();
let canon_tgt_str = canon_tgt.to_string();
let weight = edge.weight;
let metadata = edge
.metadata
.as_ref()
.map(|v| serde_json::to_string(v).unwrap_or_default());
let expected_updated_at_micros = expected_updated_at.timestamp_micros();
let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
let pool = self.backend().pool_arc();
let writer_task = pool
.writer_task_for_runtime_write(RuntimeWriteOperation::UpdateSymmetricEdge)
.map_err(RuntimeError::Storage)?;
let outcome: SymmetricEdgeUpdateOutcome = if let Some(writer_task) = writer_task {
writer_task
.send(move |conn| {
Self::update_edge_symmetric_dml(
conn,
&ns,
&edge_id_str,
&canon_src_str,
&canon_tgt_str,
&relation_str,
weight,
metadata,
expected_updated_at_micros,
expected_deleted_at_micros,
)
.map_err(|e| {
khive_storage::StorageError::driver(
khive_storage::StorageCapability::Graph,
"update_edge",
e,
)
})
})
.await
.map_err(RuntimeError::Storage)?
} else {
tokio::task::spawn_blocking(move || {
let guard = pool.writer()?;
guard.transaction(|conn| {
Self::update_edge_symmetric_dml(
conn,
&ns,
&edge_id_str,
&canon_src_str,
&canon_tgt_str,
&relation_str,
weight,
metadata,
expected_updated_at_micros,
expected_deleted_at_micros,
)
})
})
.await
.map_err(|e| {
RuntimeError::Internal(format!("update_edge: spawn_blocking join: {e}"))
})?
.map_err(RuntimeError::Sqlite)?
};
match outcome {
SymmetricEdgeUpdateOutcome::Absorbed(sid) => {
let surviving_uuid = Uuid::parse_str(&sid).map_err(|e| {
RuntimeError::Internal(format!(
"update_edge: surviving id parse failed: {e}"
))
})?;
edge = self
.get_edge_including_deleted(&record_tok, surviving_uuid)
.await?
.ok_or_else(|| {
RuntimeError::Internal(format!(
"update_edge: surviving canonical row {surviving_uuid} vanished after update"
))
})?;
}
SymmetricEdgeUpdateOutcome::Updated => {
edge.source_id = canon_src;
edge.target_id = canon_tgt;
}
SymmetricEdgeUpdateOutcome::Stale => {
return Err(crate::curation::stale_edge_snapshot_error(edge_id));
}
}
} else {
let minimum_updated_at_micros = expected_updated_at
.timestamp_micros()
.checked_add(1)
.ok_or_else(|| {
RuntimeError::Internal(format!(
"edge {edge_id} updated_at is already at i64::MAX and cannot advance"
))
})?;
let now_micros = chrono::Utc::now()
.timestamp_micros()
.max(minimum_updated_at_micros);
edge.updated_at =
chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
RuntimeError::Internal(format!(
"edge {edge_id}: computed updated_at {now_micros} is not a valid timestamp"
))
})?;
let persisted = graph
.replace_edge_if_unchanged(edge.clone(), expected_updated_at, expected_deleted_at)
.await?;
if !persisted {
return Err(crate::curation::stale_edge_snapshot_error(edge_id));
}
}
let event_store = self.events(&record_tok)?;
let event = khive_storage::event::Event::new(
record_ns.clone(),
"update",
EventKind::EdgeUpdated,
SubstrateKind::Entity,
"",
)
.with_target(edge_id)
.with_payload(
serde_json::json!({"id": edge_id, "namespace": record_ns, "changed_fields": changed_fields}),
);
event_store.append_event(event).await.map_err(|e| {
RuntimeError::Internal(format!("update_edge: event store write failed: {e}"))
})?;
Ok(edge)
}
pub async fn delete_edge(
&self,
token: &NamespaceToken,
edge_id: Uuid,
hard: bool,
) -> RuntimeResult<bool> {
let mode = if hard {
DeleteMode::Hard
} else {
DeleteMode::Soft
};
let edge = if hard {
self.get_edge_including_deleted(token, edge_id).await?
} else {
self.get_edge(token, edge_id).await?
};
let Some(edge) = edge else {
return Ok(false);
};
let record_ns: String = edge.namespace.clone();
let record_tok = token.with_namespace(
khive_types::Namespace::parse(&record_ns)
.map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
);
let graph = self.graph(&record_tok)?;
let actor = format!("{}:{}", token.actor().kind, token.actor().id);
let deleted = if hard {
self.atomic_hard_delete_with_edge_purge(
edge_hard_delete_statement(edge_id),
edge_id,
&record_ns,
&actor,
SubstrateKind::Entity,
)
.await?
} else {
graph.delete_edge(LinkId::from(edge_id), mode).await?
};
if deleted {
let event_store = self.events(&record_tok)?;
let event = khive_storage::event::Event::new(
record_ns.clone(),
"delete",
EventKind::EdgeDeleted,
SubstrateKind::Entity,
"",
)
.with_target(edge_id)
.with_payload(serde_json::json!({"id": edge_id, "namespace": record_ns, "hard": hard}));
event_store.append_event(event).await.map_err(|e| {
RuntimeError::Internal(format!("delete_edge: event store write failed: {e}"))
})?;
}
Ok(deleted)
}
pub async fn count_edges(
&self,
token: &NamespaceToken,
filter: crate::curation::EdgeListFilter,
) -> RuntimeResult<u64> {
let namespaces: Vec<String> = token
.visible_namespaces()
.iter()
.map(|namespace| namespace.as_str().to_owned())
.collect();
let graph = self.graph(token)?;
match graph
.count_edges_in_namespaces(&namespaces, filter.clone().into())
.await
{
Ok(count) => Ok(count),
Err(khive_storage::StorageError::Unsupported { operation, .. })
if operation == "count_edges_in_namespaces" =>
{
let mut total = 0;
for namespace in token.visible_namespaces() {
let scoped = NamespaceToken::for_namespace(namespace.clone());
total += self
.graph(&scoped)?
.count_edges(filter.clone().into())
.await?;
}
Ok(total)
}
Err(error) => Err(error.into()),
}
}
pub async fn build_edge(&self, token: &NamespaceToken, spec: &LinkSpec) -> RuntimeResult<Edge> {
self.build_edge_with_endpoint_kinds(token, spec)
.await
.map(|(edge, _)| edge)
}
async fn build_edge_with_endpoint_kinds(
&self,
token: &NamespaceToken,
spec: &LinkSpec,
) -> RuntimeResult<(Edge, (EdgeEndpointKind, EdgeEndpointKind))> {
validate_edge_metadata(spec.relation, spec.metadata.as_ref())?;
let ns_str = match &spec.namespace {
Some(s) => {
let spec_ns = crate::Namespace::parse(s)
.map_err(|e| RuntimeError::InvalidInput(format!("invalid namespace: {e}")))?;
if &spec_ns != token.namespace() {
return Err(RuntimeError::InvalidInput(
"LinkSpec namespace does not match token namespace".into(),
));
}
s.as_str()
}
None => token.namespace().as_str(),
};
let endpoint_kinds = self
.validate_edge_relation_endpoints(token, spec.source_id, spec.target_id, spec.relation)
.await?;
let (source_id, target_id) =
canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
let endpoint_kinds = canonical_edge_endpoint_kinds(
spec.source_id,
source_id,
endpoint_kinds.0,
endpoint_kinds.1,
);
let metadata = if spec.relation == EdgeRelation::DependsOn {
match (
self.resolve_edge_endpoint(token, source_id).await?,
self.resolve_edge_endpoint(token, target_id).await?,
) {
(Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
merge_dependency_kind(&src_e.kind, &tgt_e.kind, spec.metadata.clone())
}
_ => spec.metadata.clone(),
}
} else {
spec.metadata.clone()
};
validate_edge_metadata(spec.relation, metadata.as_ref())?;
let now = chrono::Utc::now();
Ok((
Edge {
id: LinkId::from(Uuid::new_v4()),
namespace: ns_str.to_string(),
source_id,
target_id,
relation: spec.relation,
weight: spec.weight,
created_at: now,
updated_at: now,
deleted_at: None,
metadata,
target_backend: None,
},
endpoint_kinds,
))
}
pub async fn link_many(
&self,
token: &NamespaceToken,
specs: Vec<LinkSpec>,
) -> RuntimeResult<Vec<Edge>> {
self.link_many_observed(token, specs)
.await
.map(|rows| rows.into_iter().map(|row| row.edge).collect())
}
pub async fn link_many_observed(
&self,
token: &NamespaceToken,
specs: Vec<LinkSpec>,
) -> RuntimeResult<Vec<EdgeUpsertResult>> {
self.link_many_guarded_observed(
token,
specs,
GraphMutationPreconditions::default(),
Vec::new(),
)
.await
.map(|(rows, _)| rows)
}
#[doc(hidden)]
pub async fn link_many_guarded_observed(
&self,
token: &NamespaceToken,
specs: Vec<LinkSpec>,
preconditions: GraphMutationPreconditions,
retirements: Vec<Edge>,
) -> RuntimeResult<(Vec<EdgeUpsertResult>, Vec<LinkId>)> {
let namespace = token.namespace().as_str();
let foreign_document = preconditions
.document
.as_ref()
.is_some_and(|guard| guard.namespace != namespace);
let foreign_edge = preconditions.edges.iter().any(|guard| {
guard.namespace != namespace
|| guard
.expected
.as_ref()
.is_some_and(|edge| edge.namespace != namespace)
});
if foreign_document
|| foreign_edge
|| retirements.iter().any(|edge| edge.namespace != namespace)
{
return Err(RuntimeError::InvalidInput(
"guarded link namespace does not match token namespace".into(),
));
}
if specs.is_empty()
&& preconditions.document.is_none()
&& preconditions.edges.is_empty()
&& retirements.is_empty()
{
return Ok((Vec::new(), Vec::new()));
}
let mut edges = Vec::with_capacity(specs.len());
let mut endpoint_kinds = Vec::with_capacity(specs.len());
for spec in &specs {
let (edge, kinds) = self.build_edge_with_endpoint_kinds(token, spec).await?;
edges.push(edge);
endpoint_kinds.push(kinds);
}
let requests = edges
.into_iter()
.zip(specs.iter())
.map(|(edge, spec)| EdgeUpsertRequest {
edge,
resurrect: spec.resurrect,
})
.collect();
let attribution = crate::EventAttribution::from_token(token);
let outcome = compose_graph_mutation_events(
self.backend(),
GraphMutationRequest::Batch {
requests,
guard_endpoints: true,
},
preconditions,
retirements,
move |outcome| {
let GraphMutationOutcome::Batch(batch) = &outcome.mutation else {
return Err(Self::link_composition_shape_error(
"expected a written batch",
));
};
if batch.rows.len() != endpoint_kinds.len() {
return Err(Self::link_composition_shape_error(
"edge result count differs from validated endpoint count",
));
}
let mut events = Vec::with_capacity(batch.rows.len() + outcome.retired.len());
for (row, (source_kind, target_kind)) in batch.rows.iter().zip(endpoint_kinds) {
events.push(Self::link_mutation_event(
&attribution,
row,
source_kind,
target_kind,
));
}
for edge in &outcome.retired {
let edge_id = Uuid::from(edge.id);
events.push(
attribution.stamp(
Event::new(
edge.namespace.clone(),
"delete",
EventKind::EdgeDeleted,
SubstrateKind::Entity,
"",
)
.with_target(edge_id)
.with_payload(serde_json::json!({
"id": edge_id, "namespace": edge.namespace, "hard": false,
})),
),
);
}
Ok(events)
},
)
.await?;
let retired = outcome.retired.into_iter().map(|edge| edge.id).collect();
let GraphMutationOutcome::Batch(outcome) = outcome.mutation else {
return Err(RuntimeError::Internal(
"link_many: unexpected composition outcome".into(),
));
};
if let Some(refusal) = outcome.refusal {
return match refusal.reason {
EdgeUpsertRefusal::MissingEndpoints(missing) => {
Err(RuntimeError::GuardedWriteFailed(guarded_link_batch_failure(
&specs[refusal.entry_index],
refusal.entry_index,
missing,
)))
}
EdgeUpsertRefusal::ResurrectionRequired { edge } => {
Err(RuntimeError::InvalidInput(format!(
"batch entry {} targets soft-deleted edge {}; pass resurrect=true for that link",
refusal.entry_index, edge.id
)))
}
};
}
Ok((outcome.rows, retired))
}
pub async fn link_commit_annotation_if_absent(
&self,
token: &NamespaceToken,
commit_id: Uuid,
project_id: Uuid,
guard: CommitAnnotationGuard,
) -> RuntimeResult<CommitAnnotationInsertOutcome> {
if !matches!(guard.expected_sha.len(), 40 | 64)
|| !guard
.expected_sha
.bytes()
.all(|byte| byte.is_ascii_hexdigit())
{
return Err(RuntimeError::InvalidInput(
"expected full commit SHA".into(),
));
}
let edge = self
.build_edge(
token,
&LinkSpec {
namespace: None,
source_id: commit_id,
target_id: project_id,
relation: EdgeRelation::Annotates,
weight: 1.0,
metadata: None,
resurrect: false,
},
)
.await?;
let attribution = crate::EventAttribution::from_token(token);
let outcome = compose_graph_mutation_events(
self.backend(),
GraphMutationRequest::CommitAnnotation { edge, guard },
GraphMutationPreconditions::default(),
Vec::new(),
move |outcome| match &outcome.mutation {
GraphMutationOutcome::CommitAnnotation(CommitAnnotationInsertOutcome::Created(
edge,
)) => Ok(vec![Self::link_mutation_event(
&attribution,
&EdgeUpsertResult {
edge: edge.clone(),
disposition: EdgeUpsertDisposition::Created,
previous: None,
},
EdgeEndpointKind::Note,
EdgeEndpointKind::Entity,
)]),
_ => Err(Self::link_composition_shape_error(
"expected a created annotation",
)),
},
)
.await?;
let GraphMutationOutcome::CommitAnnotation(result) = outcome.mutation else {
return Err(RuntimeError::Internal(
"link annotation: unexpected composition outcome".into(),
));
};
Ok(result)
}
pub async fn create_many(
&self,
token: &NamespaceToken,
specs: Vec<EntityCreateSpec>,
) -> RuntimeResult<Vec<Entity>> {
if specs.is_empty() {
return Ok(vec![]);
}
let ns = token.namespace().as_str();
let mut entities = Vec::with_capacity(specs.len());
for (index, spec) in specs.iter().enumerate() {
entities.push(self.validate_bulk_entity(ns, spec, &format!("entity[{index}]"))?);
}
#[cfg(any(test, feature = "fault-injection"))]
let fts_many_inject = consume_fault(&FTS_FAIL_MANY_NS, ns);
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_many_inject = false;
#[cfg(any(test, feature = "fault-injection"))]
let fts_many_inject_partial = consume_fault(&FTS_FAIL_MANY_PARTIAL_NS, ns);
#[cfg(not(any(test, feature = "fault-injection")))]
let fts_many_inject_partial = false;
let injected_failure_index = if fts_many_inject {
Some(0)
} else if fts_many_inject_partial {
Some(usize::from(entities.len() > 1))
} else {
None
};
let _ = self.entities(token)?;
let _ = self.text(token)?;
let plans = entities
.iter()
.enumerate()
.map(|(index, entity)| {
let mut plan = bulk_entity_plan(entity)?;
if injected_failure_index == Some(index) {
plan.statements.truncate(1);
plan.statements.push(PlanStatement {
statement: SqlStatement {
sql:
"INSERT INTO __khive_create_many_injected_failure__ DEFAULT VALUES"
.to_string(),
params: vec![],
label: Some("fts-insert-injected-failure".to_string()),
},
guard: None,
});
}
Ok(AtomicOpPlan::AddEntity(plan))
})
.collect::<RuntimeResult<Vec<_>>>()?;
match run_atomic_unit(self.sql().as_ref(), plans).await {
Ok(AtomicRunOutcome::Committed { .. }) => Ok(entities),
Ok(AtomicRunOutcome::RolledBack {
failed_op_index,
failure,
}) => Err(RuntimeError::Internal(format!(
"create_many: atomic batch rolled back at entity index {failed_op_index}: \
{failure:?}"
))),
Err(e) => Err(RuntimeError::Internal(format!(
"create_many: atomic batch failed: {}",
e.0
))),
}
}
fn validate_bulk_entity(
&self,
ns: &str,
spec: &EntityCreateSpec,
record: &str,
) -> RuntimeResult<Entity> {
self.validate_entity_kind(&spec.kind)?;
let validated_type =
self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
if spec.name.trim().is_empty() {
return Err(RuntimeError::InvalidInput("name must not be empty".into()));
}
crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
crate::secret_gate::check_at(&spec.name, record, "name")?;
if let Some(d) = &spec.description {
crate::secret_gate::check_at(d, record, "description")?;
}
if let Some(ref p) = spec.properties {
crate::secret_gate::check_json_at(p, record, "properties")?;
}
crate::secret_gate::check_tags_at(&spec.tags, record, "tags")?;
let mut entity =
Entity::new(ns, &spec.kind, &spec.name).with_entity_type(validated_type.as_deref());
if let Some(d) = &spec.description {
entity = entity.with_description(d);
}
if let Some(p) = spec.properties.clone() {
entity = entity.with_properties(p);
}
if !spec.tags.is_empty() {
entity = entity.with_tags(spec.tags.clone());
}
Ok(entity)
}
pub async fn prepare_bulk_entity_plan(
&self,
token: &NamespaceToken,
spec: EntityCreateSpec,
) -> RuntimeResult<(Entity, AtomicOpPlan)> {
let entity = self.validate_bulk_entity(token.namespace().as_str(), &spec, "entity")?;
let _ = self.entities(token)?;
let _ = self.text(token)?;
let plan = AtomicOpPlan::AddEntity(bulk_entity_plan(&entity)?);
Ok((entity, plan))
}
pub async fn prepare_bulk_note_plan(
&self,
token: &NamespaceToken,
spec: NoteCreateSpec,
) -> RuntimeResult<(Note, AtomicOpPlan)> {
let mut candidate = Note::new(token.namespace().as_str(), &spec.kind, &spec.content);
candidate.name = spec.name.clone();
candidate.properties = spec.properties.clone();
crate::note_write::validate_head(&candidate)?;
let mut prepared = crate::atomic_message::prepare_atomic_notes(
self,
vec![crate::atomic_message::AtomicNoteSpec {
token,
id: None,
kind: &spec.kind,
name: spec.name.as_deref(),
content: &spec.content,
properties: spec.properties,
}],
crate::atomic_message::AtomicNoteOptions {
salience: spec.salience,
embed: Some(false),
..Default::default()
},
)
.await?;
match (prepared.notes.pop(), prepared.plans.pop()) {
(Some(note), Some(plan)) if prepared.notes.is_empty() && prepared.plans.is_empty() => {
Ok((note, plan))
}
_ => Err(RuntimeError::Internal(
"bulk note preparation must yield exactly one note and one plan".into(),
)),
}
}
}
#[derive(Clone, Debug)]
pub struct NoteCreateSpec {
pub kind: String,
pub name: Option<String>,
pub content: String,
pub salience: Option<f64>,
pub properties: Option<serde_json::Value>,
}
fn bulk_entity_plan(entity: &Entity) -> RuntimeResult<AddEntityPlan> {
crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
let mut statements = vec![PlanStatement {
statement: entity_upsert_statement(entity),
guard: Some(AffectedRowGuard::exactly(1)),
}];
statements.extend(
insert_document_statements("fts_entities", &entity_fts_document(entity))
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
}),
);
Ok(AddEntityPlan {
entity_id: entity.id,
statements,
post_commit: PostCommitEffect::None,
})
}
fn guarded_link_batch_failure(
spec: &LinkSpec,
entry_index: usize,
missing: khive_storage::MissingEndpoints,
) -> GuardedWriteFailure {
let (source_id, target_id) =
canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
GuardedWriteFailure {
entry_index: Some(entry_index),
missing_source: missing.source.then_some(source_id),
missing_target: missing.target.then_some(target_id),
}
}
#[derive(Clone, Debug)]
pub struct LinkSpec {
pub namespace: Option<String>,
pub source_id: Uuid,
pub target_id: Uuid,
pub relation: EdgeRelation,
pub weight: f64,
pub metadata: Option<serde_json::Value>,
pub resurrect: bool,
}
#[derive(Clone, Debug)]
pub struct EntityCreateSpec {
pub kind: String,
pub entity_type: Option<String>,
pub name: String,
pub description: Option<String>,
pub properties: Option<serde_json::Value>,
pub tags: Vec<String>,
}
#[cfg(test)]
#[path = "operations_tests.rs"]
mod tests;