1use std::collections::HashMap;
7use std::str::FromStr;
8
9use chrono::Utc;
10use serde::Serialize;
11use uuid::Uuid;
12
13use khive_score::DeterministicScore;
14use khive_storage::entity::EntityTypeCounts;
15use khive_storage::graph::{CommitAnnotationGuard, CommitAnnotationInsertOutcome};
16use khive_storage::note::Note;
17use khive_storage::types::{
18 DeleteMode, DirectedNeighborHit, Direction, EdgeSortField, EdgeUpsertDisposition,
19 EdgeUpsertRefusal, EdgeUpsertRequest, EdgeUpsertResult, GraphPath, GuardedEdgeUpsertOutcome,
20 LinkId, NeighborCursor, NeighborHit, NeighborQuery, Page, PageRequest, SeekCursor, SortOrder,
21 SqlRow, SqlStatement, SqlValue, TextFilter, TextQueryMode, TextSearchRequest, TraversalRequest,
22};
23use khive_storage::{
24 Attachment, AttachmentSubstrate, Edge, EdgeRelation, Entity, EntityFilter, Event, EventFilter,
25 NewAttachment,
26};
27use khive_types::{EdgeEndpointRule, EndpointKind, EventKind, KhiveError, SubstrateKind};
28
29use khive_db::stores::entity::{entity_hard_delete_statement, entity_upsert_statement};
30use khive_db::stores::event::hard_delete_lineage_warning_statements;
31use khive_db::stores::graph::{
32 compose_graph_mutation_events, edge_hard_delete_statement, purge_incident_edges_statement,
33 GraphMutationOutcome, GraphMutationPreconditions, GraphMutationRequest,
34};
35use khive_db::stores::note::note_hard_delete_statement;
36use khive_db::stores::text::insert_document_statements;
37use khive_db::{pool::RuntimeWriteOperation, SqliteError};
38use rusqlite::OptionalExtension;
39
40#[cfg(test)]
41mod batch_edge_tests;
42
43struct EdgeReadWindow {
44 outcomes: Vec<Option<RuntimeResult<Option<Edge>>>>,
45 groups: Vec<(khive_types::Namespace, Vec<usize>)>,
46}
47
48fn restore_reindex_failed(kind: &str, id: Uuid, error: RuntimeError) -> RuntimeError {
52 RuntimeError::Internal(format!(
53 "{kind} {id} is restored and text-indexed, but its embedding rebuild failed \
54 and will be retried by the next reindex: {error}"
55 ))
56}
57
58fn merge_tombstone_restore_refused(id: Uuid, kept_id: impl std::fmt::Display) -> RuntimeError {
59 KhiveError::conflict(format!(
60 "merge_tombstone: {id} was merged into {kept_id}; a merge tombstone is not restorable, query the kept id"
61 ))
62 .with_details(khive_types::Details::new_owned([
63 ("reason", "merge_tombstone".into()),
64 ("merged_into", kept_id.to_string()),
65 ]))
66 .into()
67}
68
69#[derive(Debug, PartialEq, Eq)]
71pub struct EntityStatsCounts {
72 pub entities: u64,
73 pub entities_by_type: Option<EntityTypeCounts>,
74}
75
76pub async fn entity_stats_counts(
79 store: &dyn khive_storage::EntityStore,
80 token: &NamespaceToken,
81) -> RuntimeResult<EntityStatsCounts> {
82 let namespaces: Vec<String> = token
83 .visible_namespaces()
84 .iter()
85 .map(|namespace| namespace.as_str().to_owned())
86 .collect();
87 match store.count_entities_by_type(&namespaces).await? {
88 Some(groups) => {
89 let entities = groups.iter().try_fold(0_u64, |total, (_, count)| {
90 total.checked_add(*count).ok_or_else(|| {
91 RuntimeError::Internal(
92 "entity type counts exceed the scalar count range".into(),
93 )
94 })
95 })?;
96 Ok(EntityStatsCounts {
97 entities,
98 entities_by_type: Some(groups),
99 })
100 }
101 None => {
102 let entities = store
103 .count_entities(
104 token.namespace().as_str(),
105 EntityFilter {
106 namespaces,
107 ..EntityFilter::default()
108 },
109 )
110 .await?;
111 Ok(EntityStatsCounts {
112 entities,
113 entities_by_type: None,
114 })
115 }
116 }
117}
118
119pub struct EntityClaimSpec {
122 pub id: Uuid,
123 pub kind: String,
124 pub entity_type: Option<String>,
125 pub name: String,
126 pub description: Option<String>,
127 pub properties: Option<serde_json::Value>,
128 pub tags: Vec<String>,
129 pub identity_tag: String,
130}
131
132fn live_merged_entity_refused(id: Uuid, kept_id: impl std::fmt::Display) -> RuntimeError {
133 KhiveError::conflict(format!(
134 "live_merged_entity: {id} is live but still carries merged_into {kept_id}; a row an \
135 earlier restore left live over its merge is not restorable, re-tombstone it or query the \
136 kept id"
137 ))
138 .with_details(khive_types::Details::new_owned([
139 ("reason", "live_merged_entity".into()),
140 ("merged_into", kept_id.to_string()),
141 ]))
142 .into()
143}
144
145fn restore_key_conflict(key: &str, holder: &Note) -> RuntimeError {
146 KhiveError::conflict(format!(
147 "restore_key_conflict: key {key:?} is already held by live note {}",
148 holder.id
149 ))
150 .with_details(khive_types::Details::new_owned([
151 ("reason", "restore_key_conflict".into()),
152 ("key", key.to_owned()),
153 ("existing_id", holder.id.to_string()),
154 ]))
155 .into()
156}
157
158use crate::atomic_plan::{
159 AddEntityPlan, AffectedRowGuard, DeletePlan, PlanStatement, PostCommitEffect, UpdatePlan,
160};
161use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicOpPlan, AtomicRunOutcome};
162use crate::curation::{entity_fts_document, note_embedding_text_ref, note_fts_document};
163use crate::error::{GuardedWriteFailure, RuntimeError, RuntimeResult};
164use crate::runtime::{KhiveRuntime, NamespaceToken};
165
166#[derive(Clone, Debug, Serialize)]
169pub struct PostCommitDegradation {
170 pub stage: &'static str,
171 pub error: String,
172}
173
174impl PostCommitDegradation {
175 fn new(stage: &'static str, error: impl ToString) -> Self {
176 Self {
177 stage,
178 error: error.to_string(),
179 }
180 }
181}
182
183macro_rules! conditional_insert_stages {
185 ($($stage:ident => $label:literal),+ $(,)?) => {
186 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
192 pub enum ConditionalInsertStage {
193 $($stage,)+
194 }
195
196 impl ConditionalInsertStage {
197 pub const ALL: &'static [Self] = &[$(Self::$stage,)+];
199
200 pub const fn label(self) -> &'static str {
201 match self {
202 $(Self::$stage => $label,)+
203 }
204 }
205 }
206 };
207}
208
209conditional_insert_stages! {
210 FtsAcquisition => "fts_acquisition",
211 FtsUpsert => "fts_upsert",
212 Embedding => "embedding",
213 VectorAcquisition => "vector_acquisition",
214 VectorInsert => "vector_insert",
215}
216
217impl ConditionalInsertStage {
218 pub fn from_label(label: &str) -> Option<Self> {
219 Self::ALL
220 .iter()
221 .copied()
222 .find(|stage| stage.label() == label)
223 }
224}
225
226fn record_conditional_insert_degradation(
227 degradations: &mut Vec<PostCommitDegradation>,
228 id: Uuid,
229 stage: ConditionalInsertStage,
230 error: impl ToString,
231) {
232 record_post_commit_degradation(degradations, "try_create_note", id, stage.label(), error);
233}
234
235fn record_post_commit_degradation(
236 degradations: &mut Vec<PostCommitDegradation>,
237 operation: &'static str,
238 id: Uuid,
239 stage: &'static str,
240 error: impl ToString,
241) {
242 let degradation = PostCommitDegradation::new(stage, error);
243 tracing::warn!(%operation, %id, stage, error = %degradation.error,
244 "substrate mutation committed with post-commit degradation");
245 degradations.push(degradation);
246}
247
248fn legacy_post_commit_result<T>(
252 operation: &'static str,
253 id: Uuid,
254 value: T,
255 degradations: Vec<PostCommitDegradation>,
256) -> RuntimeResult<T> {
257 if degradations.is_empty() {
258 return Ok(value);
259 }
260 let failures = serde_json::Value::Array(
261 degradations
262 .iter()
263 .map(|failure| {
264 serde_json::json!({"stage": failure.stage, "error": failure.error.as_str()})
265 })
266 .collect(),
267 );
268 Err(KhiveError::internal(format!(
269 "{operation} committed record {id}, but post-commit work failed; do not retry the mutation; reconcile by record_id"
270 ))
271 .with_details(khive_types::Details::new_owned([
272 ("reason", "post_commit_degraded".to_string()),
273 ("operation", operation.to_string()),
274 ("record_id", id.to_string()),
275 ("committed", "true".to_string()),
276 ("retryable", "false".to_string()),
277 ("post_commit_degradations", failures.to_string()),
278 ]))
279 .into())
280}
281
282pub(crate) fn legacy_post_commit_result_with_embedding<T>(
283 operation: &'static str,
284 id: Uuid,
285 value: T,
286 embedding: crate::retrieval::EmbeddingTruncationReport,
287 degradations: Vec<PostCommitDegradation>,
288) -> RuntimeResult<T> {
289 if !embedding.any_truncated() {
290 return legacy_post_commit_result(operation, id, value, degradations);
291 }
292 let failures = serde_json::Value::Array(
293 degradations
294 .iter()
295 .map(|failure| {
296 serde_json::json!({"stage": failure.stage, "error": failure.error.as_str()})
297 })
298 .collect(),
299 );
300 Err(KhiveError::internal(format!(
301 "{operation} committed record {id}, but embedding input was truncated; do not retry the mutation; reconcile by record_id"
302 ))
303 .with_details(khive_types::Details::new_owned([
304 ("reason", "embedding_input_truncated".to_string()),
305 ("operation", operation.to_string()),
306 ("record_id", id.to_string()),
307 ("committed", "true".to_string()),
308 ("retryable", "false".to_string()),
309 (
310 "embedding_truncation_report",
311 serde_json::json!(embedding).to_string(),
312 ),
313 ("post_commit_degradations", failures.to_string()),
314 ]))
315 .into())
316}
317
318#[cfg(test)]
320std::thread_local! {
321 static LINK_FAIL_AFTER: std::cell::Cell<usize> = const { std::cell::Cell::new(0) };
322}
323
324#[cfg(any(test, feature = "fault-injection"))]
325std::thread_local! {
326 static VECTOR_FAIL_AFTER: std::cell::Cell<Option<usize>> =
327 const { std::cell::Cell::new(None) };
328}
329
330#[cfg(any(test, feature = "fault-injection"))]
333pub fn arm_vector_fail_after(n: usize) {
334 VECTOR_FAIL_AFTER.with(|cell| cell.set(Some(n)));
335}
336
337#[cfg(any(test, feature = "fault-injection"))]
340type FaultArmSet = std::sync::Mutex<std::collections::HashMap<String, std::sync::Arc<()>>>;
341#[cfg(any(test, feature = "fault-injection"))]
342const MAX_FAULT_ARMS: usize = 64;
343#[cfg(any(test, feature = "fault-injection"))]
344static FTS_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
345 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
346#[cfg(any(test, feature = "fault-injection"))]
347static VECTOR_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
348 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
349#[cfg(any(test, feature = "fault-injection"))]
351static ENTITY_COMPENSATION_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
352 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
353#[cfg(any(test, feature = "fault-injection"))]
356static FTS_FAIL_MANY_NS: std::sync::LazyLock<FaultArmSet> =
357 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
358#[cfg(any(test, feature = "fault-injection"))]
361static FTS_FAIL_MANY_PARTIAL_NS: std::sync::LazyLock<FaultArmSet> =
362 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
363#[cfg(any(test, feature = "fault-injection"))]
368static PREFIX_RESOLVE_FAIL_NS: std::sync::LazyLock<FaultArmSet> =
369 std::sync::LazyLock::new(|| std::sync::Mutex::new(std::collections::HashMap::new()));
370
371#[cfg(any(test, feature = "fault-injection"))]
373#[must_use = "the fault injection is disarmed when this guard is dropped"]
374pub struct FaultInjectionArm {
375 namespace: String,
376 token: std::sync::Arc<()>,
377 arms: &'static FaultArmSet,
378}
379
380#[cfg(any(test, feature = "fault-injection"))]
381impl Drop for FaultInjectionArm {
382 fn drop(&mut self) {
383 let mut arms = self.arms.lock().unwrap();
384 if arms
385 .get(&self.namespace)
386 .is_some_and(|token| std::sync::Arc::ptr_eq(token, &self.token))
387 {
388 arms.remove(&self.namespace);
389 }
390 }
391}
392
393#[cfg(any(test, feature = "fault-injection"))]
394fn arm_fault(arms: &'static FaultArmSet, namespace: &str, max_arms: usize) -> FaultInjectionArm {
395 let token = std::sync::Arc::new(());
396 let refusal = {
397 let mut active = arms.lock().unwrap();
398 if active.contains_key(namespace) {
399 Some("the namespace is already armed")
400 } else if active.len() >= max_arms {
401 Some("the arm set is at capacity")
402 } else {
403 active.insert(namespace.to_string(), std::sync::Arc::clone(&token));
404 None
405 }
406 };
407 if let Some(reason) = refusal {
408 panic!("cannot arm fault injection for namespace `{namespace}`: {reason}");
409 }
410 FaultInjectionArm {
411 namespace: namespace.to_string(),
412 token,
413 arms,
414 }
415}
416
417#[cfg(any(test, feature = "fault-injection"))]
418fn consume_fault(arms: &FaultArmSet, namespace: &str) -> bool {
419 arms.lock().unwrap().remove(namespace).is_some()
420}
421#[cfg(any(test, feature = "fault-injection"))]
428static FTS_SEARCH_FAIL_NS: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
429
430#[cfg(any(test, feature = "fault-injection"))]
445pub fn arm_fts_fail_scoped(ns: &str) -> FaultInjectionArm {
446 arm_fault(&FTS_FAIL_NS, ns, MAX_FAULT_ARMS)
447}
448
449#[cfg(any(test, feature = "fault-injection"))]
459pub fn arm_fts_fail_many_scoped(ns: &str) -> FaultInjectionArm {
460 arm_fault(&FTS_FAIL_MANY_NS, ns, MAX_FAULT_ARMS)
461}
462
463#[cfg(any(test, feature = "fault-injection"))]
472pub fn arm_fts_fail_many_partial_scoped(ns: &str) -> FaultInjectionArm {
473 arm_fault(&FTS_FAIL_MANY_PARTIAL_NS, ns, MAX_FAULT_ARMS)
474}
475
476#[cfg(any(test, feature = "fault-injection"))]
486pub fn arm_fts_search_fail(ns: &str) {
487 *FTS_SEARCH_FAIL_NS.lock().unwrap() = Some(ns.to_string());
488}
489
490#[cfg(any(test, feature = "fault-injection"))]
500pub fn arm_vector_fail_scoped(ns: &str) -> FaultInjectionArm {
501 arm_fault(&VECTOR_FAIL_NS, ns, MAX_FAULT_ARMS)
502}
503
504#[cfg(any(test, feature = "fault-injection"))]
507pub fn arm_entity_compensation_fail_scoped(ns: &str) -> FaultInjectionArm {
508 arm_fault(&ENTITY_COMPENSATION_FAIL_NS, ns, MAX_FAULT_ARMS)
509}
510
511#[cfg(any(test, feature = "fault-injection"))]
523pub fn arm_prefix_resolve_fail_scoped(prefix: &str) -> FaultInjectionArm {
524 arm_fault(&PREFIX_RESOLVE_FAIL_NS, prefix, MAX_FAULT_ARMS)
525}
526
527#[cfg(any(test, feature = "fault-injection"))]
532static ROLLBACK_CLEANUP_FAIL_NS: std::sync::Mutex<Option<String>> = std::sync::Mutex::new(None);
533
534#[cfg(any(test, feature = "fault-injection"))]
541pub fn arm_rollback_cleanup_fail(ns: &str) {
542 *ROLLBACK_CLEANUP_FAIL_NS.lock().unwrap() = Some(ns.to_string());
543}
544
545#[cfg(any(test, feature = "fault-injection"))]
554pub(crate) fn consume_fts_fail_fault(ns: &str) -> bool {
555 consume_fault(&FTS_FAIL_NS, ns)
556}
557#[cfg(any(test, feature = "fault-injection"))]
558pub(crate) fn consume_vector_fail_fault(ns: &str) -> bool {
559 consume_fault(&VECTOR_FAIL_NS, ns)
560}
561
562#[derive(Clone, Debug)]
564pub struct NoteSearchHit {
565 pub note_id: Uuid,
566 pub score: DeterministicScore,
567 pub rank_score_kind: crate::RankScoreKind,
568 pub signals: crate::SearchSignals,
569 pub source: crate::SearchSource,
570 pub title: Option<String>,
571 pub snippet: Option<String>,
572}
573
574fn salience_weighted_rank(score: DeterministicScore, salience: Option<f64>) -> DeterministicScore {
575 const SCALE_RAW: i128 = 1_i128 << 32;
576 let salience = DeterministicScore::from_f64(salience.unwrap_or(0.5));
577 let weight_raw = SCALE_RAW / 2 + i128::from(salience.to_raw()) / 2;
578 let weighted_raw = i128::from(score.to_raw()) * weight_raw / SCALE_RAW;
581 DeterministicScore::from_raw(weighted_raw.clamp(
582 i128::from(DeterministicScore::NEG_INF.to_raw()),
583 i128::from(DeterministicScore::MAX.to_raw()),
584 ) as i64)
585}
586
587#[derive(Clone, Debug)]
591pub struct NoteSearchOutcome {
592 pub hits: Vec<NoteSearchHit>,
593 pub vector_error: Option<String>,
594}
595
596pub fn hex_prefix_to_uuid_pattern(prefix: &str) -> String {
607 if prefix.contains('-') {
608 return prefix.to_string();
609 }
610 const BOUNDARIES: [usize; 4] = [8, 13, 18, 23]; let mut out = String::with_capacity(36);
612 for c in prefix.chars() {
613 if BOUNDARIES.contains(&out.len()) {
614 out.push('-');
615 }
616 out.push(c);
617 }
618 out
619}
620
621pub fn uuid_prefix_bounds(prefix: &str) -> Option<(String, String)> {
631 const HYPHEN_POSITIONS: [usize; 4] = [8, 13, 18, 23];
632
633 let compact = if prefix.contains('-') {
634 if prefix.len() > 36 {
635 return None;
636 }
637 let mut compact = String::with_capacity(32);
638 for (index, byte) in prefix.bytes().enumerate() {
639 if HYPHEN_POSITIONS.contains(&index) {
640 if byte != b'-' {
641 return None;
642 }
643 } else if byte.is_ascii_hexdigit() {
644 compact.push(char::from(byte.to_ascii_lowercase()));
645 } else {
646 return None;
647 }
648 }
649 compact
650 } else {
651 if prefix.is_empty()
652 || prefix.len() > 32
653 || !prefix.bytes().all(|byte| byte.is_ascii_hexdigit())
654 {
655 return None;
656 }
657 prefix.to_ascii_lowercase()
658 };
659
660 if compact.is_empty() || compact.len() > 32 {
661 return None;
662 }
663
664 let lower = hex_prefix_to_uuid_pattern(&compact);
665 let mut successor = compact.into_bytes();
666 let mut carried_past_start = true;
667 for index in (0..successor.len()).rev() {
668 let next = match successor[index] {
669 b'0'..=b'8' | b'a'..=b'e' => Some(successor[index] + 1),
670 b'9' => Some(b'a'),
671 b'f' => None,
672 _ => return None,
673 };
674 if let Some(next) = next {
675 successor[index] = next;
676 successor.truncate(index + 1);
677 carried_past_start = false;
678 break;
679 }
680 }
681
682 let upper = if carried_past_start {
683 "g".to_string()
684 } else {
685 let compact_upper = String::from_utf8(successor).ok()?;
686 hex_prefix_to_uuid_pattern(&compact_upper)
687 };
688 Some((lower, upper))
689}
690
691fn resolve_prefix_statement(
692 table: &str,
693 has_deleted_at: bool,
694 include_deleted: bool,
695 namespaces: Option<&[String]>,
696 lower: &str,
697 upper: &str,
698) -> SqlStatement {
699 let namespace_clause = namespaces.map(|namespaces| {
700 let placeholders: Vec<String> = (0..namespaces.len())
701 .map(|index| format!("?{}", index + 3))
702 .collect();
703 format!(" AND namespace IN ({})", placeholders.join(", "))
704 });
705 let deleted_filter = if has_deleted_at && !include_deleted {
706 " AND deleted_at IS NULL"
707 } else {
708 ""
709 };
710 let mut params = vec![
711 SqlValue::Text(lower.to_owned()),
712 SqlValue::Text(upper.to_owned()),
713 ];
714 if let Some(namespaces) = namespaces {
715 params.extend(
716 namespaces
717 .iter()
718 .map(|namespace| SqlValue::Text(namespace.clone())),
719 );
720 }
721
722 SqlStatement {
723 sql: format!(
724 "SELECT id FROM {table} \
725 WHERE id >= ?1 AND id < ?2{namespace_clause}{deleted_filter} ORDER BY id LIMIT 2",
726 namespace_clause = namespace_clause.as_deref().unwrap_or("")
727 ),
728 params,
729 label: Some("resolve_prefix".into()),
730 }
731}
732
733fn text_preview(text: &str, max_chars: usize) -> Option<String> {
734 let trimmed = text.trim();
735 if trimmed.is_empty() {
736 None
737 } else {
738 Some(trimmed.chars().take(max_chars).collect())
739 }
740}
741
742fn normalize_symmetric_direction(
748 direction: Direction,
749 relations: Option<&[EdgeRelation]>,
750) -> Direction {
751 let Some(rels) = relations else {
752 return direction;
753 };
754 if rels.is_empty() {
755 return direction;
756 }
757 let all_symmetric = rels
758 .iter()
759 .all(|r| matches!(r, EdgeRelation::CompetesWith | EdgeRelation::ComposedWith));
760 if all_symmetric {
761 Direction::Both
762 } else {
763 direction
764 }
765}
766
767fn direction_sort_rank(direction: &Direction) -> u8 {
774 match direction {
775 Direction::Out => 0,
776 Direction::In => 1,
777 Direction::Both => 2,
778 }
779}
780
781fn note_title(note: &Note) -> Option<String> {
782 note.name
783 .clone()
784 .filter(|s| !s.trim().is_empty())
785 .or_else(|| Some(format!("[{}]", note.kind.as_str())))
786}
787
788fn note_snippet(note: &Note) -> Option<String> {
789 text_preview(¬e.content, 200)
790}
791
792#[derive(Clone, Debug)]
794pub enum Resolved {
795 Entity(Entity),
796 Note(Note),
797 Event(Event),
798 PackRecord {
805 pack: String,
806 kind: String,
807 data: serde_json::Value,
808 },
809}
810
811#[derive(Clone, Copy, Debug, Eq, PartialEq)]
818pub enum EdgeEndpointKind {
819 Entity,
820 Note,
821 Event,
822 Edge,
823}
824
825impl EdgeEndpointKind {
826 pub const fn name(self) -> &'static str {
828 match self {
829 Self::Entity => "entity",
830 Self::Note => "note",
831 Self::Event => "event",
832 Self::Edge => "edge",
833 }
834 }
835}
836
837fn resolved_pair(r: Option<&Resolved>) -> Option<(&'static str, &str, Option<&str>)> {
843 match r? {
844 Resolved::Entity(e) => Some(("entity", e.kind.as_str(), e.entity_type.as_deref())),
845 Resolved::Note(n) => Some(("note", n.kind.as_str(), None)),
846 Resolved::Event(_) => None,
847 Resolved::PackRecord { .. } => None,
848 }
849}
850
851pub fn endpoint_matches(
859 spec: &EndpointKind,
860 substrate: &str,
861 kind: &str,
862 entity_type: Option<&str>,
863) -> bool {
864 match spec {
865 EndpointKind::EntityOfKind(k) => substrate == "entity" && *k == kind,
866 EndpointKind::NoteOfKind(k) => substrate == "note" && *k == kind,
867 EndpointKind::EntityOfType {
868 kind: k,
869 entity_type: t,
870 } => substrate == "entity" && *k == kind && entity_type == Some(*t),
871 }
872}
873
874fn pattern_endpoint_matches(
888 spec: &EndpointKind,
889 substrate: &str,
890 kind: &str,
891 entity_type: Option<&str>,
892) -> bool {
893 match spec {
894 EndpointKind::EntityOfType {
895 kind: k,
896 entity_type: t,
897 } => substrate == "entity" && *k == kind && entity_type.is_none_or(|et| et == *t),
898 _ => endpoint_matches(spec, substrate, kind, entity_type),
899 }
900}
901
902pub fn accepted_pack_relations_for_entities(
916 rules: &[EdgeEndpointRule],
917 src_kind: &str,
918 src_entity_type: Option<&str>,
919 tgt_kind: &str,
920 tgt_entity_type: Option<&str>,
921) -> Vec<EdgeRelation> {
922 let mut relations: Vec<EdgeRelation> = rules
923 .iter()
924 .filter(|r| {
925 endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
926 && endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
927 })
928 .map(|r| r.relation)
929 .collect();
930 relations.sort_by_key(|r| r.as_str());
931 relations.dedup();
932 relations
933}
934
935pub fn accepted_entity_relations_for_entities(
946 rules: &[EdgeEndpointRule],
947 src_kind: &str,
948 src_entity_type: Option<&str>,
949 tgt_kind: &str,
950 tgt_entity_type: Option<&str>,
951) -> Vec<EdgeRelation> {
952 let mut relations: Vec<EdgeRelation> = BASE_ENTITY_ENDPOINT_RULES
953 .iter()
954 .filter(|(src, _relation, tgt)| (*src == "*" || *src == src_kind) && *tgt == tgt_kind)
955 .map(|(_src, relation, _tgt)| *relation)
956 .collect();
957 relations.extend(
958 accepted_pack_relations_for_entities(
959 rules,
960 src_kind,
961 src_entity_type,
962 tgt_kind,
963 tgt_entity_type,
964 )
965 .into_iter()
966 .filter(|relation| {
967 *relation != EdgeRelation::Annotates && !crate::pack::is_special_relation(*relation)
968 }),
969 );
970 relations.sort_by_key(|relation| relation.as_str());
971 relations.dedup();
972 relations
973}
974
975fn accepted_entity_relations_description(
976 rules: &[EdgeEndpointRule],
977 src_kind: &str,
978 src_entity_type: Option<&str>,
979 tgt_kind: &str,
980 tgt_entity_type: Option<&str>,
981) -> String {
982 let relations = accepted_entity_relations_for_entities(
983 rules,
984 src_kind,
985 src_entity_type,
986 tgt_kind,
987 tgt_entity_type,
988 );
989 if relations.is_empty() {
990 "none".to_string()
991 } else {
992 relations
993 .iter()
994 .map(EdgeRelation::as_str)
995 .collect::<Vec<_>>()
996 .join(", ")
997 }
998}
999
1000fn accepted_pack_relations_for_pattern_entities(
1006 rules: &[EdgeEndpointRule],
1007 src_kind: &str,
1008 src_entity_type: Option<&str>,
1009 tgt_kind: &str,
1010 tgt_entity_type: Option<&str>,
1011) -> Vec<EdgeRelation> {
1012 let mut relations: Vec<EdgeRelation> = rules
1013 .iter()
1014 .filter(|r| {
1015 pattern_endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
1016 && pattern_endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
1017 })
1018 .map(|r| r.relation)
1019 .collect();
1020 relations.sort_by_key(|r| r.as_str());
1021 relations.dedup();
1022 relations
1023}
1024
1025fn accepted_entity_kind_pairs_for_relation(
1036 pack_rules: &[EdgeEndpointRule],
1037 relation: EdgeRelation,
1038) -> Vec<(&'static str, &'static str)> {
1039 let mut pairs = Vec::new();
1040 for src in khive_types::EntityKind::ALL {
1041 for tgt in khive_types::EntityKind::ALL {
1042 let allowed = base_entity_rule_allows(src.name(), relation, tgt.name())
1043 || (!crate::pack::is_special_relation(relation)
1044 && accepted_pack_relations_for_pattern_entities(
1045 pack_rules,
1046 src.name(),
1047 None,
1048 tgt.name(),
1049 None,
1050 )
1051 .contains(&relation));
1052 if allowed {
1053 pairs.push((src.name(), tgt.name()));
1054 }
1055 }
1056 }
1057 pairs
1058}
1059
1060fn static_impossible_edge_pattern_warnings(
1075 language: khive_query::QueryLanguage,
1076 pattern: &khive_query::ast::MatchPattern,
1077 pack_rules: &[EdgeEndpointRule],
1078) -> Vec<String> {
1079 use khive_query::ast::{EdgeDirection, PatternElement};
1080
1081 if language != khive_query::QueryLanguage::Gql {
1082 return Vec::new();
1083 }
1084
1085 let elements = &pattern.elements;
1086 let mut warnings = Vec::new();
1087
1088 for (i, el) in elements.iter().enumerate() {
1089 let PatternElement::Edge(edge) = el else {
1090 continue;
1091 };
1092 if edge.relations.len() != 1 || edge.min_hops != 1 || edge.max_hops != 1 {
1093 continue;
1094 }
1095 let (left, right) = match (elements.get(i.wrapping_sub(1)), elements.get(i + 1)) {
1096 (Some(PatternElement::Node(l)), Some(PatternElement::Node(r))) => (l, r),
1097 _ => continue,
1098 };
1099 let (src_node, tgt_node) = match edge.direction {
1100 EdgeDirection::Out => (left, right),
1101 EdgeDirection::In => (right, left),
1102 EdgeDirection::Both => continue,
1103 };
1104 let (Some(src_raw), Some(tgt_raw)) = (src_node.kind.as_deref(), tgt_node.kind.as_deref())
1105 else {
1106 continue;
1107 };
1108 let (Ok(src_kind), Ok(tgt_kind)) = (
1109 src_raw.parse::<khive_types::EntityKind>(),
1110 tgt_raw.parse::<khive_types::EntityKind>(),
1111 ) else {
1112 continue;
1113 };
1114 let Ok(relation) = edge.relations[0].parse::<EdgeRelation>() else {
1115 continue;
1116 };
1117
1118 let possible = base_entity_rule_allows(src_kind.name(), relation, tgt_kind.name())
1119 || (!crate::pack::is_special_relation(relation)
1120 && accepted_pack_relations_for_pattern_entities(
1121 pack_rules,
1122 src_kind.name(),
1123 src_node.entity_type.as_deref(),
1124 tgt_kind.name(),
1125 tgt_node.entity_type.as_deref(),
1126 )
1127 .contains(&relation));
1128 if possible {
1129 continue;
1130 }
1131
1132 let accepted = accepted_entity_kind_pairs_for_relation(pack_rules, relation);
1133 let accepted_str = if accepted.is_empty() {
1134 "none".to_string()
1135 } else {
1136 accepted
1137 .iter()
1138 .map(|(s, t)| format!("{s}->{t}"))
1139 .collect::<Vec<_>>()
1140 .join(", ")
1141 };
1142 warnings.push(format!(
1143 "pattern ({src})-[:{relation}]->({tgt}) can never match: '{relation}' does not accept \
1144 {src}->{tgt} endpoints; accepted source->target kinds for '{relation}': {accepted_str}",
1145 src = src_kind.name(),
1146 tgt = tgt_kind.name(),
1147 ));
1148 }
1149
1150 warnings
1151}
1152
1153fn pack_rule_allows(
1156 rules: &[EdgeEndpointRule],
1157 relation: EdgeRelation,
1158 src: Option<&Resolved>,
1159 tgt: Option<&Resolved>,
1160) -> bool {
1161 let Some((src_sub, src_kind, src_type)) = resolved_pair(src) else {
1162 return false;
1163 };
1164 let Some((tgt_sub, tgt_kind, tgt_type)) = resolved_pair(tgt) else {
1165 return false;
1166 };
1167 rules.iter().any(|r| {
1168 r.relation == relation
1169 && endpoint_matches(&r.source, src_sub, src_kind, src_type)
1170 && endpoint_matches(&r.target, tgt_sub, tgt_kind, tgt_type)
1171 })
1172}
1173
1174pub const BASE_ENTITY_ENDPOINT_RULES: &[(&str, EdgeRelation, &str)] = &[
1184 ("concept", EdgeRelation::Contains, "concept"),
1186 ("project", EdgeRelation::Contains, "project"),
1187 ("project", EdgeRelation::Contains, "artifact"),
1188 ("org", EdgeRelation::Contains, "project"),
1189 ("org", EdgeRelation::Contains, "service"),
1190 ("concept", EdgeRelation::PartOf, "concept"),
1191 ("project", EdgeRelation::PartOf, "project"),
1192 ("project", EdgeRelation::PartOf, "org"),
1193 ("*", EdgeRelation::InstanceOf, "concept"),
1194 ("service", EdgeRelation::InstanceOf, "project"),
1195 ("document", EdgeRelation::LinksTo, "document"),
1200 ("concept", EdgeRelation::LocatedIn, "concept"),
1203 ("org", EdgeRelation::LocatedIn, "concept"),
1204 ("person", EdgeRelation::Owns, "org"),
1206 ("org", EdgeRelation::Owns, "org"),
1207 ("concept", EdgeRelation::Extends, "concept"),
1209 ("concept", EdgeRelation::VariantOf, "concept"),
1210 ("artifact", EdgeRelation::VariantOf, "artifact"),
1211 ("concept", EdgeRelation::IntroducedBy, "document"),
1212 ("concept", EdgeRelation::IntroducedBy, "person"),
1213 ("artifact", EdgeRelation::IntroducedBy, "document"),
1214 ("project", EdgeRelation::IntroducedBy, "document"),
1215 ("service", EdgeRelation::IntroducedBy, "document"),
1218 ("document", EdgeRelation::IntroducedBy, "person"),
1219 ("document", EdgeRelation::IntroducedBy, "org"),
1220 ("concept", EdgeRelation::IntroducedBy, "org"),
1221 ("artifact", EdgeRelation::DerivedFrom, "dataset"),
1223 ("artifact", EdgeRelation::DerivedFrom, "document"),
1224 ("artifact", EdgeRelation::DerivedFrom, "project"),
1225 ("artifact", EdgeRelation::DerivedFrom, "artifact"),
1226 ("document", EdgeRelation::DerivedFrom, "document"),
1229 ("document", EdgeRelation::Precedes, "document"),
1231 ("dataset", EdgeRelation::Precedes, "dataset"),
1232 ("artifact", EdgeRelation::Precedes, "artifact"),
1233 ("service", EdgeRelation::Precedes, "service"),
1234 ("project", EdgeRelation::Precedes, "project"),
1235 ("project", EdgeRelation::DependsOn, "project"),
1237 ("service", EdgeRelation::DependsOn, "project"),
1238 ("service", EdgeRelation::DependsOn, "service"),
1239 ("service", EdgeRelation::DependsOn, "artifact"),
1240 ("service", EdgeRelation::DependsOn, "dataset"),
1241 ("artifact", EdgeRelation::DependsOn, "project"),
1242 ("artifact", EdgeRelation::DependsOn, "service"),
1243 ("document", EdgeRelation::DependsOn, "document"),
1244 ("concept", EdgeRelation::Enables, "concept"),
1245 ("service", EdgeRelation::Enables, "concept"),
1246 ("dataset", EdgeRelation::Enables, "concept"),
1247 ("project", EdgeRelation::Implements, "concept"),
1249 ("service", EdgeRelation::Implements, "concept"),
1250 ("concept", EdgeRelation::CompetesWith, "concept"),
1252 ("project", EdgeRelation::CompetesWith, "project"),
1253 ("service", EdgeRelation::CompetesWith, "service"),
1254 ("org", EdgeRelation::CompetesWith, "org"),
1255 ("concept", EdgeRelation::ComposedWith, "concept"),
1256 ("project", EdgeRelation::ComposedWith, "project"),
1257 ("concept", EdgeRelation::Supersedes, "concept"),
1259 ("document", EdgeRelation::Supersedes, "document"),
1260 ("artifact", EdgeRelation::Supersedes, "artifact"),
1261 ("service", EdgeRelation::Supersedes, "service"),
1262 ("dataset", EdgeRelation::Supersedes, "dataset"),
1263 ("concept", EdgeRelation::Supports, "concept"),
1265 ("document", EdgeRelation::Supports, "concept"),
1266 ("dataset", EdgeRelation::Supports, "concept"),
1267 ("artifact", EdgeRelation::Supports, "concept"),
1268 ("concept", EdgeRelation::Refutes, "concept"),
1269 ("document", EdgeRelation::Refutes, "concept"),
1270 ("dataset", EdgeRelation::Refutes, "concept"),
1271 ("artifact", EdgeRelation::Refutes, "concept"),
1272];
1273
1274pub fn base_entity_endpoint_rules() -> &'static [(&'static str, EdgeRelation, &'static str)] {
1280 BASE_ENTITY_ENDPOINT_RULES
1281}
1282
1283pub fn base_entity_rule_allows(src_kind: &str, relation: EdgeRelation, tgt_kind: &str) -> bool {
1289 BASE_ENTITY_ENDPOINT_RULES.iter().any(|(src, rel, tgt)| {
1290 *rel == relation && (*src == "*" || *src == src_kind) && *tgt == tgt_kind
1291 })
1292}
1293
1294pub(crate) fn canonical_edge_endpoints(
1300 relation: EdgeRelation,
1301 source_id: Uuid,
1302 target_id: Uuid,
1303) -> (Uuid, Uuid) {
1304 relation.canonical_endpoints(source_id, target_id)
1305}
1306
1307pub(crate) fn canonical_edge_endpoint_kinds(
1309 requested_source_id: Uuid,
1310 canonical_source_id: Uuid,
1311 source_kind: EdgeEndpointKind,
1312 target_kind: EdgeEndpointKind,
1313) -> (EdgeEndpointKind, EdgeEndpointKind) {
1314 if requested_source_id == canonical_source_id {
1315 (source_kind, target_kind)
1316 } else {
1317 (target_kind, source_kind)
1318 }
1319}
1320
1321pub(crate) fn infer_dependency_kind(src_kind: &str, tgt_kind: &str) -> Option<&'static str> {
1327 match (src_kind, tgt_kind) {
1328 ("project", "project") => Some("build"),
1329 ("service", "service") => Some("runtime"),
1330 ("service", "dataset") => Some("data"),
1331 ("service", "artifact") => Some("artifact"),
1332 ("artifact", "project") | ("artifact", "service") => Some("tooling"),
1333 ("document", "document") => Some("normative"),
1334 _ => None,
1335 }
1336}
1337
1338pub(crate) fn merge_dependency_kind(
1348 src_kind: &str,
1349 tgt_kind: &str,
1350 metadata: Option<serde_json::Value>,
1351) -> Option<serde_json::Value> {
1352 let metadata = metadata.filter(|value| !value.is_null());
1355 if let Some(ref m) = metadata {
1356 if m.get("dependency_kind").is_some() {
1357 return metadata;
1358 }
1359 }
1360 let Some(inferred) = infer_dependency_kind(src_kind, tgt_kind) else {
1361 return metadata;
1362 };
1363 let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
1364 if let Some(o) = obj.as_object_mut() {
1365 o.insert("dependency_kind".to_string(), serde_json::json!(inferred));
1366 }
1367 Some(obj)
1368}
1369
1370pub fn merge_entry_metadata(
1384 metadata: Option<serde_json::Value>,
1385 dependency_kind: Option<String>,
1386) -> RuntimeResult<Option<serde_json::Value>> {
1387 validate_metadata_shape(metadata.as_ref())?;
1388 let metadata = metadata.filter(|value| !value.is_null());
1389 let Some(dk) = dependency_kind else {
1390 return Ok(metadata);
1391 };
1392 let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
1393 let map = obj
1394 .as_object_mut()
1395 .ok_or_else(|| RuntimeError::InvalidInput("metadata must be a JSON object".into()))?;
1396 map.entry("dependency_kind".to_string())
1397 .or_insert_with(|| serde_json::json!(dk));
1398 Ok(Some(obj))
1399}
1400
1401const VALID_DEPENDENCY_KINDS: &[&str] = &[
1403 "build",
1404 "runtime",
1405 "data",
1406 "artifact",
1407 "tooling",
1408 "normative",
1409];
1410
1411pub(crate) fn validate_edge_weight(weight: f64) -> RuntimeResult<()> {
1417 if !weight.is_finite() || !(0.0..=1.0).contains(&weight) {
1418 return Err(RuntimeError::InvalidInput(format!(
1419 "edge weight must be finite and in [0.0, 1.0], got {weight}"
1420 )));
1421 }
1422 Ok(())
1423}
1424
1425fn validate_metadata_shape(metadata: Option<&serde_json::Value>) -> RuntimeResult<()> {
1426 if metadata.is_some_and(|value| !value.is_null() && !value.is_object()) {
1427 return Err(RuntimeError::InvalidInput(
1428 "metadata must be a JSON object".into(),
1429 ));
1430 }
1431 Ok(())
1432}
1433
1434pub(crate) fn validate_edge_metadata(
1439 relation: EdgeRelation,
1440 metadata: Option<&serde_json::Value>,
1441) -> RuntimeResult<()> {
1442 validate_metadata_shape(metadata)?;
1443 let Some(meta) = metadata.filter(|value| !value.is_null()) else {
1444 return Ok(());
1445 };
1446 let object = meta.as_object().expect("validated metadata object");
1447 if object
1448 .get("optional")
1449 .is_some_and(|value| !value.is_boolean())
1450 {
1451 return Err(RuntimeError::InvalidInput(
1452 "metadata.optional must be a boolean".into(),
1453 ));
1454 }
1455 if let Some(dk) = meta.get("dependency_kind") {
1456 if relation != EdgeRelation::DependsOn {
1457 return Err(RuntimeError::InvalidInput(format!(
1458 "dependency_kind is only valid on depends_on edges (got {})",
1459 relation.as_str()
1460 )));
1461 }
1462 let dk_str = dk
1463 .as_str()
1464 .ok_or_else(|| RuntimeError::InvalidInput("dependency_kind must be a string".into()))?;
1465 if !VALID_DEPENDENCY_KINDS.contains(&dk_str) {
1466 return Err(RuntimeError::InvalidInput(format!(
1467 "unknown dependency_kind {dk_str:?}; valid: {}",
1468 VALID_DEPENDENCY_KINDS.join(" | ")
1469 )));
1470 }
1471 }
1472 Ok(())
1473}
1474
1475fn note_props_match(note_props: Option<&serde_json::Value>, filter: &serde_json::Value) -> bool {
1480 let required = match filter.as_object() {
1481 Some(obj) if !obj.is_empty() => obj,
1482 _ => return true,
1483 };
1484 let actual = match note_props.and_then(serde_json::Value::as_object) {
1485 Some(obj) => obj,
1486 None => return false,
1487 };
1488 required
1489 .iter()
1490 .all(|(k, v)| actual.get(k).is_some_and(|av| av == v))
1491}
1492
1493fn note_graph_name(note: &Note) -> String {
1494 note.name
1495 .as_deref()
1496 .filter(|name| !name.trim().is_empty())
1497 .map(str::to_owned)
1498 .unwrap_or_else(|| format!("[{}]", note.kind))
1499}
1500
1501fn merge_traversal_paths_by_root(paths: Vec<GraphPath>, limit: Option<u32>) -> Vec<GraphPath> {
1506 let mut order: Vec<Uuid> = Vec::new();
1507 let mut merged: HashMap<Uuid, GraphPath> = HashMap::new();
1508 let mut node_index: HashMap<Uuid, HashMap<Uuid, usize>> = HashMap::new();
1512
1513 for path in paths {
1514 let existing = merged.entry(path.root_id).or_insert_with(|| {
1515 order.push(path.root_id);
1516 GraphPath {
1517 root_id: path.root_id,
1518 nodes: Vec::new(),
1519 total_weight: 0.0,
1520 }
1521 });
1522 let index = node_index.entry(path.root_id).or_default();
1523 for node in path.nodes {
1524 match index.get(&node.node_id) {
1525 Some(&i) => {
1526 if node.depth < existing.nodes[i].depth {
1527 existing.nodes[i] = node;
1528 }
1529 }
1530 None => {
1531 index.insert(node.node_id, existing.nodes.len());
1532 existing.nodes.push(node);
1533 }
1534 }
1535 }
1536 }
1537
1538 order
1539 .into_iter()
1540 .filter_map(|root_id| merged.remove(&root_id))
1541 .map(|mut path| {
1542 path.nodes.sort_by_key(|n| n.depth);
1544 if let Some(lim) = limit {
1545 let lim = lim as usize;
1546 let mut non_root_kept = 0usize;
1547 path.nodes.retain(|n| {
1548 if n.depth == 0 {
1549 return true;
1550 }
1551 if non_root_kept < lim {
1552 non_root_kept += 1;
1553 true
1554 } else {
1555 false
1556 }
1557 });
1558 }
1559 recompute_total_weight(&mut path);
1560 path
1561 })
1562 .collect()
1563}
1564
1565fn recompute_total_weight(path: &mut GraphPath) {
1574 path.total_weight = path.nodes.iter().map(|n| n.weight).fold(0.0_f64, f64::max);
1575}
1576
1577async fn drain_embed_join_set<T: Send + 'static>(
1591 mut join_set: tokio::task::JoinSet<(usize, RuntimeResult<T>)>,
1592 model_count: usize,
1593) -> RuntimeResult<Vec<T>> {
1594 let mut vectors: Vec<Option<T>> = (0..model_count).map(|_| None).collect();
1595
1596 while let Some(joined) = join_set.join_next().await {
1597 match joined {
1598 Ok((idx, Ok(vector))) => vectors[idx] = Some(vector),
1599 Ok((_idx, Err(e))) => {
1600 join_set.abort_all();
1601 return Err(e);
1602 }
1603 Err(join_err) => {
1604 join_set.abort_all();
1605 return Err(RuntimeError::Internal(format!(
1606 "embed task panicked: {join_err}"
1607 )));
1608 }
1609 }
1610 }
1611
1612 Ok(vectors
1613 .into_iter()
1614 .map(|v| v.expect("every model index observed exactly once by join_set drain"))
1615 .collect())
1616}
1617
1618impl KhiveRuntime {
1619 async fn compensate_entity_create(
1622 &self,
1623 token: &NamespaceToken,
1624 entity_id: Uuid,
1625 namespace: &str,
1626 vector_models: &[String],
1627 ) -> Vec<String> {
1628 let mut cleanup_errors = Vec::new();
1629
1630 #[cfg(any(test, feature = "fault-injection"))]
1631 let entity_delete_injected = consume_fault(&ENTITY_COMPENSATION_FAIL_NS, namespace);
1632 #[cfg(not(any(test, feature = "fault-injection")))]
1633 let entity_delete_injected = false;
1634
1635 if entity_delete_injected {
1636 cleanup_errors.push("entity row delete: injected compensation failure".to_string());
1637 } else {
1638 match self.entities(token) {
1639 Ok(store) => {
1640 if let Err(error) = store.delete_entity(entity_id, DeleteMode::Hard).await {
1641 cleanup_errors.push(format!("entity row delete: {error}"));
1642 }
1643 }
1644 Err(error) => cleanup_errors.push(format!("entity store access: {error}")),
1645 }
1646 }
1647
1648 match self.text(token) {
1649 Ok(fts) => {
1650 if let Err(error) = fts.delete_document(namespace, entity_id).await {
1651 cleanup_errors.push(format!("FTS document delete: {error}"));
1652 }
1653 }
1654 Err(error) => cleanup_errors.push(format!("FTS store access: {error}")),
1655 }
1656
1657 for model_name in vector_models {
1658 match self.vectors_for_model(token, model_name) {
1659 Ok(vectors) => {
1660 if let Err(error) = vectors.delete(entity_id).await {
1661 cleanup_errors
1662 .push(format!("vector delete for model {model_name}: {error}"));
1663 }
1664 }
1665 Err(error) => cleanup_errors.push(format!(
1666 "vector store access for model {model_name}: {error}"
1667 )),
1668 }
1669 }
1670
1671 cleanup_errors
1672 }
1673
1674 fn entity_create_failure(
1675 entity_id: Uuid,
1676 primary: RuntimeError,
1677 cleanup_errors: Vec<String>,
1678 ) -> RuntimeError {
1679 if cleanup_errors.is_empty() {
1680 primary
1681 } else {
1682 RuntimeError::Khive(KhiveError::internal(format!(
1683 "create_entity indexing failed for record {entity_id}; primary failure: \
1684 {primary}; compensation failure(s): {}; partial persistence is possible; \
1685 inspect and reconcile this record before retrying",
1686 cleanup_errors.join("; ")
1687 )))
1688 }
1689 }
1690
1691 pub async fn claim_entity_if_absent(
1694 &self,
1695 token: &NamespaceToken,
1696 spec: EntityClaimSpec,
1697 ) -> RuntimeResult<(Entity, bool)> {
1698 self.validate_entity_kind(&spec.kind)?;
1699 let entity_type =
1700 self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
1701 crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
1702 crate::secret_gate::check_at(&spec.name, "entity", "name")?;
1703 if let Some(description) = &spec.description {
1704 crate::secret_gate::check_at(description, "entity", "description")?;
1705 }
1706 if let Some(properties) = &spec.properties {
1707 crate::secret_gate::check_json_at(properties, "entity", "properties")?;
1708 }
1709 crate::secret_gate::check_tags_at(&spec.tags, "entity", "tags")?;
1710
1711 let mut proposed = Entity::new(token.namespace().as_str(), &spec.kind, &spec.name);
1712 proposed.id = spec.id;
1713 proposed.entity_type = entity_type.clone();
1714 proposed.description = spec.description;
1715 proposed.properties = spec.properties;
1716 proposed.tags = spec.tags;
1717
1718 let store = self.entities(token)?;
1719 let inserted = store.insert_entity_if_absent(proposed.clone()).await?;
1720 let entity = if inserted {
1721 proposed
1722 } else {
1723 store
1724 .get_entity_including_deleted(spec.id)
1725 .await?
1726 .ok_or_else(|| {
1727 RuntimeError::Internal(format!(
1728 "entity claim {} lost but the winning row is missing",
1729 spec.id
1730 ))
1731 })?
1732 };
1733 if entity.deleted_at.is_some() {
1734 return Err(RuntimeError::InvalidInput(format!(
1735 "entity claim {} is soft-deleted; restore it explicitly",
1736 entity.id
1737 )));
1738 }
1739 if entity.namespace != token.namespace().as_str()
1740 || entity.kind != spec.kind
1741 || entity.entity_type.as_deref() != entity_type.as_deref()
1742 || !entity.name.eq_ignore_ascii_case(&spec.name)
1743 || !entity
1744 .tags
1745 .iter()
1746 .any(|tag| tag.eq_ignore_ascii_case(&spec.identity_tag))
1747 {
1748 return Err(RuntimeError::InvalidInput(format!(
1749 "entity claim {} belongs to a different record",
1750 entity.id
1751 )));
1752 }
1753
1754 self.ensure_claimed_entity_create_event(token, &entity)
1755 .await?;
1756 self.reindex_claimed_entity(token, &entity).await?;
1757 Ok((entity, inserted))
1758 }
1759
1760 pub async fn ensure_claimed_entity_create_event(
1763 &self,
1764 token: &NamespaceToken,
1765 entity: &Entity,
1766 ) -> RuntimeResult<()> {
1767 if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
1768 return Err(RuntimeError::InvalidInput(format!(
1769 "entity {} is not a live row in the write namespace",
1770 entity.id
1771 )));
1772 }
1773 let events = self.events(token).map_err(|error| {
1774 RuntimeError::Internal(format!(
1775 "entity {} persists but its create event store is unavailable: {error}",
1776 entity.id
1777 ))
1778 })?;
1779 let filter = EventFilter {
1780 target_id: Some(entity.id),
1781 kinds: vec![EventKind::EntityCreated],
1782 verbs: vec!["create".into()],
1783 substrates: vec![SubstrateKind::Entity],
1784 after: Some(entity.created_at.saturating_sub(1)),
1785 ..EventFilter::default()
1786 };
1787 let page = PageRequest {
1788 offset: 0,
1789 limit: 1,
1790 };
1791 if !events
1792 .query_events(filter.clone(), page.clone())
1793 .await?
1794 .items
1795 .is_empty()
1796 {
1797 return Ok(());
1798 }
1799
1800 let mut event = Event::new(
1801 entity.namespace.clone(),
1802 "create",
1803 EventKind::EntityCreated,
1804 SubstrateKind::Entity,
1805 "",
1806 )
1807 .with_target(entity.id)
1808 .with_payload(serde_json::json!({
1809 "id": entity.id,
1810 "namespace": &entity.namespace,
1811 "kind": &entity.kind,
1812 }));
1813 let event_seed = Uuid::new_v5(&Uuid::NAMESPACE_URL, b"khive:claimed-entity-create:v1");
1814 let mut event_key = Vec::with_capacity(24);
1815 event_key.extend_from_slice(entity.id.as_bytes());
1816 event_key.extend_from_slice(&entity.created_at.to_be_bytes());
1817 event.id = Uuid::new_v5(&event_seed, &event_key);
1818 if let Err(error) = events.append_event(event).await {
1819 if events.query_events(filter, page).await?.items.is_empty() {
1820 return Err(RuntimeError::Internal(format!(
1821 "entity {} persists but its create event failed: {error}",
1822 entity.id
1823 )));
1824 }
1825 }
1826 Ok(())
1827 }
1828
1829 pub async fn reindex_claimed_entity(
1832 &self,
1833 token: &NamespaceToken,
1834 entity: &Entity,
1835 ) -> RuntimeResult<()> {
1836 if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
1837 return Err(RuntimeError::InvalidInput(format!(
1838 "entity {} is not a live row in the write namespace",
1839 entity.id
1840 )));
1841 }
1842 let doc = entity_fts_document(entity);
1843 let embed_body = doc.body.clone();
1844 #[cfg(any(test, feature = "fault-injection"))]
1845 let fts_inject = consume_fault(&FTS_FAIL_NS, &entity.namespace);
1846 #[cfg(not(any(test, feature = "fault-injection")))]
1847 let fts_inject = false;
1848 let fts_result = if fts_inject {
1849 Err(RuntimeError::Internal("injected FTS failure".into()))
1850 } else {
1851 match self.text(token) {
1852 Ok(text) => text.upsert_document(doc).await.map_err(Into::into),
1853 Err(error) => Err(error),
1854 }
1855 };
1856 fts_result.map_err(|error| {
1857 RuntimeError::Internal(format!(
1858 "entity {} persists but its text index failed: {error}",
1859 entity.id
1860 ))
1861 })?;
1862
1863 for model_name in self.registered_embedding_model_names() {
1864 let outcome = self
1865 .embed_document_with_model_outcome_for_token(token, &model_name, &embed_body)
1866 .await
1867 .map_err(|error| {
1868 RuntimeError::Internal(format!(
1869 "entity {} persists but model {model_name} embedding failed: {error}",
1870 entity.id
1871 ))
1872 })?;
1873 #[cfg(any(test, feature = "fault-injection"))]
1874 let vector_inject = consume_fault(&VECTOR_FAIL_NS, &entity.namespace);
1875 #[cfg(not(any(test, feature = "fault-injection")))]
1876 let vector_inject = false;
1877 if vector_inject {
1878 return Err(RuntimeError::Internal(format!(
1879 "entity {} persists but model {model_name} vector indexing failed: injected vector failure",
1880 entity.id
1881 )));
1882 }
1883 self.vectors_for_model(token, &model_name)
1884 .map_err(|error| {
1885 RuntimeError::Internal(format!(
1886 "entity {} persists but model {model_name} vector store is unavailable: {error}",
1887 entity.id
1888 ))
1889 })?
1890 .insert(
1891 entity.id,
1892 SubstrateKind::Entity,
1893 &entity.namespace,
1894 "entity.body",
1895 vec![outcome.vector],
1896 )
1897 .await
1898 .map_err(|error| {
1899 RuntimeError::Internal(format!(
1900 "entity {} persists but model {model_name} vector indexing failed: {error}",
1901 entity.id
1902 ))
1903 })?;
1904 }
1905 Ok(())
1906 }
1907
1908 #[allow(clippy::too_many_arguments)]
1919 #[cfg(test)]
1920 pub(crate) async fn create_entity(
1921 &self,
1922 token: &NamespaceToken,
1923 kind: &str,
1924 entity_type: Option<&str>,
1925 name: &str,
1926 description: Option<&str>,
1927 properties: Option<serde_json::Value>,
1928 tags: Vec<String>,
1929 ) -> RuntimeResult<Entity> {
1930 let (entity, _, degradations) = self
1931 .create_entity_with_embedding_report_inner(
1932 token,
1933 kind,
1934 entity_type,
1935 name,
1936 description,
1937 properties,
1938 tags,
1939 Vec::new(),
1940 )
1941 .await?;
1942 legacy_post_commit_result("create_entity", entity.id, entity, degradations)
1943 }
1944
1945 #[allow(clippy::too_many_arguments)]
1958 pub async fn create_entity_with_attachments(
1959 &self,
1960 token: &NamespaceToken,
1961 kind: &str,
1962 entity_type: Option<&str>,
1963 name: &str,
1964 description: Option<&str>,
1965 properties: Option<serde_json::Value>,
1966 tags: Vec<String>,
1967 attachments: Vec<NewAttachment>,
1968 ) -> RuntimeResult<Entity> {
1969 let (entity, embedding, degradations) = self
1970 .create_entity_with_attachments_inner(
1971 token,
1972 kind,
1973 entity_type,
1974 name,
1975 description,
1976 properties,
1977 tags,
1978 attachments,
1979 )
1980 .await?;
1981 legacy_post_commit_result_with_embedding(
1982 "create_entity_with_attachments",
1983 entity.id,
1984 entity,
1985 embedding,
1986 degradations,
1987 )
1988 }
1989
1990 #[allow(clippy::too_many_arguments)]
1992 pub async fn create_entity_with_attachments_and_report(
1993 &self,
1994 token: &NamespaceToken,
1995 kind: &str,
1996 entity_type: Option<&str>,
1997 name: &str,
1998 description: Option<&str>,
1999 properties: Option<serde_json::Value>,
2000 tags: Vec<String>,
2001 attachments: Vec<NewAttachment>,
2002 ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
2003 let (entity, embedding, degradations) = self
2004 .create_entity_with_attachments_inner(
2005 token,
2006 kind,
2007 entity_type,
2008 name,
2009 description,
2010 properties,
2011 tags,
2012 attachments,
2013 )
2014 .await?;
2015 legacy_post_commit_result(
2016 "create_entity_with_attachments_and_report",
2017 entity.id,
2018 (entity, embedding),
2019 degradations,
2020 )
2021 }
2022
2023 #[allow(clippy::too_many_arguments)]
2024 async fn create_entity_with_attachments_inner(
2025 &self,
2026 token: &NamespaceToken,
2027 kind: &str,
2028 entity_type: Option<&str>,
2029 name: &str,
2030 description: Option<&str>,
2031 properties: Option<serde_json::Value>,
2032 tags: Vec<String>,
2033 attachments: Vec<NewAttachment>,
2034 ) -> RuntimeResult<(
2035 Entity,
2036 crate::retrieval::EmbeddingTruncationReport,
2037 Vec<PostCommitDegradation>,
2038 )> {
2039 drop(self.attachments()?);
2043 let blob_store = self.blob_store().ok_or_else(|| {
2044 RuntimeError::Unconfigured(
2045 "create_entity_with_attachments requires an installed BlobStore".to_string(),
2046 )
2047 })?;
2048 let mut roles = std::collections::HashSet::with_capacity(attachments.len());
2049 for attachment in &attachments {
2050 attachment.validate()?;
2051 if !roles.insert(attachment.role.as_str()) {
2052 return Err(RuntimeError::InvalidInput(format!(
2053 "duplicate attachment role {:?}",
2054 attachment.role
2055 )));
2056 }
2057 }
2058 for attachment in &attachments {
2059 if !blob_store.exists(&attachment.content_ref).await? {
2060 return Err(RuntimeError::InvalidInput(format!(
2061 "create_entity_with_attachments requires a published blob; no object exists for {}",
2062 attachment.content_ref
2063 )));
2064 }
2065 }
2066 let validated_type = self.validate_entity_type_for_kind(kind, entity_type)?;
2067 let (entity, embedding, degradations) = self
2068 .create_entity_with_embedding_report_inner(
2069 token,
2070 kind,
2071 validated_type.as_deref(),
2072 name,
2073 description,
2074 properties,
2075 tags,
2076 attachments,
2077 )
2078 .await?;
2079 Ok((entity, embedding, degradations))
2080 }
2081
2082 #[allow(clippy::too_many_arguments)]
2083 pub async fn create_entity_with_embedding_report(
2084 &self,
2085 token: &NamespaceToken,
2086 kind: &str,
2087 entity_type: Option<&str>,
2088 name: &str,
2089 description: Option<&str>,
2090 properties: Option<serde_json::Value>,
2091 tags: Vec<String>,
2092 ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
2093 let (entity, embedding, degradations) = self
2094 .create_entity_with_embedding_report_inner(
2095 token,
2096 kind,
2097 entity_type,
2098 name,
2099 description,
2100 properties,
2101 tags,
2102 Vec::new(),
2103 )
2104 .await?;
2105 legacy_post_commit_result(
2106 "create_entity_with_embedding_report",
2107 entity.id,
2108 (entity, embedding),
2109 degradations,
2110 )
2111 }
2112
2113 #[allow(clippy::too_many_arguments)]
2115 pub async fn create_entity_with_post_commit_report(
2116 &self,
2117 token: &NamespaceToken,
2118 kind: &str,
2119 entity_type: Option<&str>,
2120 name: &str,
2121 description: Option<&str>,
2122 properties: Option<serde_json::Value>,
2123 tags: Vec<String>,
2124 ) -> RuntimeResult<(
2125 Entity,
2126 crate::retrieval::EmbeddingTruncationReport,
2127 Vec<PostCommitDegradation>,
2128 )> {
2129 self.create_entity_with_embedding_report_inner(
2130 token,
2131 kind,
2132 entity_type,
2133 name,
2134 description,
2135 properties,
2136 tags,
2137 Vec::new(),
2138 )
2139 .await
2140 }
2141
2142 #[allow(clippy::too_many_arguments)]
2143 async fn create_entity_with_embedding_report_inner(
2144 &self,
2145 token: &NamespaceToken,
2146 kind: &str,
2147 entity_type: Option<&str>,
2148 name: &str,
2149 description: Option<&str>,
2150 properties: Option<serde_json::Value>,
2151 tags: Vec<String>,
2152 attachments: Vec<NewAttachment>,
2153 ) -> RuntimeResult<(
2154 Entity,
2155 crate::retrieval::EmbeddingTruncationReport,
2156 Vec<PostCommitDegradation>,
2157 )> {
2158 self.validate_entity_kind(kind)?;
2159 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
2160 crate::secret_gate::check_at(name, "entity", "name")?;
2162 if let Some(d) = description {
2163 crate::secret_gate::check_at(d, "entity", "description")?;
2164 }
2165 if let Some(ref p) = properties {
2166 crate::secret_gate::check_json_at(p, "entity", "properties")?;
2167 }
2168 crate::secret_gate::check_tags_at(&tags, "entity", "tags")?;
2169 let ns = token.namespace().as_str();
2170 let mut entity = Entity::new(ns, kind, name).with_entity_type(entity_type);
2171 if let Some(d) = description {
2172 entity = entity.with_description(d);
2173 }
2174 if let Some(p) = properties {
2175 entity = entity.with_properties(p);
2176 }
2177 if !tags.is_empty() {
2178 entity = entity.with_tags(tags);
2179 }
2180 let projected_content_ref = attachments
2181 .iter()
2182 .find(|attachment| attachment.role == "content")
2183 .map(|attachment| attachment.content_ref.to_string());
2184 let attachment_rows = attachments
2185 .into_iter()
2186 .map(|attachment| {
2187 Attachment::from_new(
2188 entity.id,
2189 AttachmentSubstrate::Entity,
2190 attachment,
2191 entity.created_at,
2192 )
2193 })
2194 .collect();
2195 self.entities(token)?
2196 .upsert_entity_with_attachments(entity.clone(), attachment_rows)
2197 .await?;
2198 entity.content_ref = projected_content_ref;
2199
2200 let doc = entity_fts_document(&entity);
2201 let embed_body = doc.body.clone();
2202
2203 {
2205 #[cfg(any(test, feature = "fault-injection"))]
2206 let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
2207 #[cfg(not(any(test, feature = "fault-injection")))]
2208 let fts_inject = false;
2209 let fts_result: RuntimeResult<()> = if fts_inject {
2210 Err(RuntimeError::Internal("injected FTS failure".to_string()))
2211 } else {
2212 match self.text(token) {
2213 Ok(fts) => fts.upsert_document(doc).await.map_err(RuntimeError::from),
2214 Err(e) => Err(e),
2215 }
2216 };
2217 if let Err(e) = fts_result {
2218 let cleanup_errors = self
2219 .compensate_entity_create(token, entity.id, ns, &[])
2220 .await;
2221 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2222 }
2223 }
2224
2225 let embed_model_names = {
2228 let names = self.registered_embedding_model_names();
2229 if names.is_empty() {
2230 vec![]
2231 } else {
2232 names
2233 }
2234 };
2235
2236 let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
2237 if embed_model_names.len() == 1 {
2238 let model_name = &embed_model_names[0];
2239 let vec_result = self
2240 .embed_document_with_model_outcome_for_token(token, model_name, &embed_body)
2241 .await;
2242
2243 #[cfg(any(test, feature = "fault-injection"))]
2244 let vec_inject = consume_fault(&VECTOR_FAIL_NS, ns);
2245 #[cfg(not(any(test, feature = "fault-injection")))]
2246 let vec_inject = false;
2247 let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
2248 if vec_inject {
2249 Err(RuntimeError::Internal(
2250 "injected vector failure".to_string(),
2251 ))
2252 } else {
2253 vec_result
2254 };
2255
2256 let single_result: RuntimeResult<()> = match vec_result {
2257 Ok(outcome) => {
2258 embedding_report.observe(&outcome);
2259 match self.vectors_for_model(token, model_name) {
2260 Ok(vs) => vs
2261 .insert(
2262 entity.id,
2263 SubstrateKind::Entity,
2264 ns,
2265 "entity.body",
2266 vec![outcome.vector],
2267 )
2268 .await
2269 .map_err(RuntimeError::from),
2270 Err(e) => Err(e),
2271 }
2272 }
2273 Err(e) => Err(e),
2274 };
2275 if let Err(e) = single_result {
2276 let cleanup_errors = self
2277 .compensate_entity_create(
2278 token,
2279 entity.id,
2280 ns,
2281 std::slice::from_ref(model_name),
2282 )
2283 .await;
2284 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2285 }
2286 } else if !embed_model_names.is_empty() {
2287 let rt_clone = self.clone();
2290 let body_owned = embed_body.clone();
2291 let usage_ctx = crate::usage::current();
2292 let mut join_set = tokio::task::JoinSet::new();
2293 for (idx, model_name) in embed_model_names.iter().enumerate() {
2294 let rt = rt_clone.clone();
2295 let text = body_owned.clone();
2296 let name = model_name.clone();
2297 let ctx = usage_ctx.clone();
2298 let token = (*token).clone();
2299 join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
2300 let fut = rt.embed_document_with_model_outcome_for_token(&token, &name, &text);
2301 let result = match ctx {
2302 Some(ctx) => crate::usage::scope(ctx, fut).await,
2303 None => fut.await,
2304 };
2305 (idx, result)
2306 }));
2307 }
2308 let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
2312 Ok(outcomes) => outcomes,
2313 Err(e) => {
2314 let cleanup_errors = self
2315 .compensate_entity_create(token, entity.id, ns, &[])
2316 .await;
2317 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2318 }
2319 };
2320 let mut inserted_models: Vec<String> = Vec::with_capacity(embed_model_names.len());
2322 for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
2323 embedding_report.observe(&outcome);
2324 #[cfg(any(test, feature = "fault-injection"))]
2326 let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
2327 Some(0) => {
2328 cell.set(None);
2329 true
2330 }
2331 Some(n) => {
2332 cell.set(Some(n - 1));
2333 false
2334 }
2335 None => false,
2336 });
2337 #[cfg(not(any(test, feature = "fault-injection")))]
2338 let count_inject = false;
2339
2340 let insert_result = if count_inject {
2341 Err(RuntimeError::Internal(
2342 "injected vector insert failure".to_string(),
2343 ))
2344 } else {
2345 match self.vectors_for_model(token, model_name) {
2346 Ok(vs) => vs
2347 .insert(
2348 entity.id,
2349 SubstrateKind::Entity,
2350 ns,
2351 "entity.body",
2352 vec![outcome.vector],
2353 )
2354 .await
2355 .map_err(RuntimeError::from),
2356 Err(e) => Err(e),
2357 }
2358 };
2359 if let Err(e) = insert_result {
2360 let mut cleanup_models = inserted_models.clone();
2363 cleanup_models.push(model_name.clone());
2364 let cleanup_errors = self
2365 .compensate_entity_create(token, entity.id, ns, &cleanup_models)
2366 .await;
2367 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2368 }
2369 inserted_models.push(model_name.clone());
2370 }
2371 }
2372
2373 let created_event = khive_storage::event::Event::new(
2380 entity.namespace.clone(),
2381 "create",
2382 EventKind::EntityCreated,
2383 SubstrateKind::Entity,
2384 "",
2385 )
2386 .with_target(entity.id)
2387 .with_payload(serde_json::json!({
2388 "id": entity.id,
2389 "namespace": entity.namespace,
2390 "kind": entity.kind,
2391 }));
2392 let event_result = match self.events(token) {
2393 Ok(store) => store
2394 .append_event(created_event)
2395 .await
2396 .map_err(RuntimeError::from),
2397 Err(error) => Err(error),
2398 };
2399 let mut degradations = Vec::new();
2400 if let Err(error) = event_result {
2401 record_post_commit_degradation(
2402 &mut degradations,
2403 "create_entity",
2404 entity.id,
2405 "event_append",
2406 error,
2407 );
2408 }
2409
2410 Ok((entity, embedding_report, degradations))
2411 }
2412
2413 pub async fn get_entity(&self, token: &NamespaceToken, id: Uuid) -> RuntimeResult<Entity> {
2426 let store = self.entities(token)?;
2427 if let Some(entity) = store.get_entity(id).await? {
2428 return Ok(entity);
2429 }
2430 if let Some(tombstone) = store.get_entity_including_deleted(id).await? {
2431 if let Some(kept_id) = tombstone.merged_into {
2432 return Err(RuntimeError::NotFound(format!(
2433 "{id} was merged into {kept_id}; query the kept id"
2434 )));
2435 }
2436 }
2437 Err(RuntimeError::NotFound(format!("entity {id}")))
2438 }
2439
2440 pub async fn get_entity_including_deleted(
2444 &self,
2445 token: &NamespaceToken,
2446 id: Uuid,
2447 ) -> RuntimeResult<Option<Entity>> {
2448 self.entities(token)?
2449 .get_entity_including_deleted(id)
2450 .await
2451 .map_err(Into::into)
2452 }
2453
2454 pub async fn get_note_including_deleted(
2458 &self,
2459 token: &NamespaceToken,
2460 id: Uuid,
2461 ) -> RuntimeResult<Option<khive_storage::note::Note>> {
2462 self.notes(token)?
2463 .get_note_including_deleted(id)
2464 .await
2465 .map_err(Into::into)
2466 }
2467
2468 pub async fn get_entities_by_ids(
2472 &self,
2473 token: &NamespaceToken,
2474 ids: &[Uuid],
2475 ) -> RuntimeResult<Vec<Entity>> {
2476 if ids.is_empty() {
2477 return Ok(vec![]);
2478 }
2479 let filter = EntityFilter {
2480 ids: ids.to_vec(),
2481 ..Default::default()
2482 };
2483 let page = self
2484 .entities(token)?
2485 .query_entities(
2486 token.namespace().as_str(),
2487 filter,
2488 PageRequest {
2489 offset: 0,
2490 limit: ids.len() as u32,
2491 },
2492 )
2493 .await?;
2494 Ok(page.items)
2495 }
2496
2497 async fn get_entities_by_ids_visible(
2506 &self,
2507 token: &NamespaceToken,
2508 ids: &[Uuid],
2509 ) -> RuntimeResult<Vec<Entity>> {
2510 if ids.is_empty() {
2511 return Ok(vec![]);
2512 }
2513 let namespaces: Vec<String> = token
2514 .visible_namespaces()
2515 .iter()
2516 .map(|ns| ns.as_str().to_owned())
2517 .collect();
2518 let filter = EntityFilter {
2519 ids: ids.to_vec(),
2520 namespaces,
2521 ..Default::default()
2522 };
2523 let page = self
2524 .entities(token)?
2525 .query_entities(
2526 token.namespace().as_str(),
2527 filter,
2528 PageRequest {
2529 offset: 0,
2530 limit: ids.len() as u32,
2531 },
2532 )
2533 .await?;
2534 Ok(page.items)
2535 }
2536
2537 pub(crate) fn ensure_namespace(record_ns: &str, caller_primary_ns: &str) -> RuntimeResult<()> {
2546 if record_ns == caller_primary_ns {
2547 return Ok(());
2548 }
2549 Err(RuntimeError::NotFound("not found in this namespace".into()))
2550 }
2551
2552 pub(crate) fn ensure_namespace_visible(
2558 record_ns: &str,
2559 token: &NamespaceToken,
2560 ) -> RuntimeResult<()> {
2561 for ns in token.visible_namespaces() {
2562 if record_ns == ns.as_str() {
2563 return Ok(());
2564 }
2565 }
2566 Err(RuntimeError::NotFound("not found in this namespace".into()))
2567 }
2568
2569 pub async fn list_entities(
2576 &self,
2577 token: &NamespaceToken,
2578 kind: Option<&str>,
2579 entity_type: Option<&str>,
2580 limit: u32,
2581 offset: u32,
2582 ) -> RuntimeResult<Vec<Entity>> {
2583 let filter = EntityFilter {
2584 kinds: kind
2585 .map(|value| vec![value.to_string()])
2586 .unwrap_or_default(),
2587 entity_types: entity_type
2588 .map(|value| vec![value.to_string()])
2589 .unwrap_or_default(),
2590 legacy_entity_type_fallback: true,
2591 ..Default::default()
2592 };
2593 self.list_entities_filtered(token, filter, limit, offset)
2594 .await
2595 }
2596
2597 pub async fn list_entities_filtered(
2600 &self,
2601 token: &NamespaceToken,
2602 mut filter: EntityFilter,
2603 limit: u32,
2604 offset: u32,
2605 ) -> RuntimeResult<Vec<Entity>> {
2606 filter.namespaces = token
2607 .visible_namespaces()
2608 .iter()
2609 .map(|namespace| namespace.as_str().to_owned())
2610 .collect();
2611 let page = self
2612 .entities(token)?
2613 .query_entities_count_free(
2614 token.namespace().as_str(),
2615 filter,
2616 PageRequest {
2617 offset: offset.into(),
2618 limit,
2619 },
2620 )
2621 .await?;
2622 Ok(page.items)
2623 }
2624
2625 pub async fn list_entities_after(
2632 &self,
2633 token: &NamespaceToken,
2634 kind: Option<&str>,
2635 entity_type: Option<&str>,
2636 tags_any: &[String],
2637 after: Option<Uuid>,
2638 limit: u32,
2639 ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
2640 let filter = EntityFilter {
2641 kinds: kind
2642 .map(|value| vec![value.to_string()])
2643 .unwrap_or_default(),
2644 entity_types: entity_type
2645 .map(|value| vec![value.to_string()])
2646 .unwrap_or_default(),
2647 legacy_entity_type_fallback: true,
2648 tags_any: tags_any.to_vec(),
2649 ..Default::default()
2650 };
2651 self.list_entities_after_filtered(token, filter, after, limit)
2652 .await
2653 }
2654
2655 pub async fn list_entities_after_filtered(
2658 &self,
2659 token: &NamespaceToken,
2660 mut filter: EntityFilter,
2661 after: Option<Uuid>,
2662 limit: u32,
2663 ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
2664 let store = self.entities(token)?;
2665 let after = match after {
2666 Some(id) => {
2667 let entity = self
2668 .get_entity_including_deleted(token, id)
2669 .await?
2670 .ok_or_else(|| RuntimeError::NotFound(format!("entity cursor {id}")))?;
2671 Self::ensure_namespace_visible(&entity.namespace, token)?;
2672 let sequence = store.entity_sequence(id).await?.ok_or_else(|| {
2673 RuntimeError::Internal(format!(
2674 "entity cursor {id} has no insertion-sequence ledger row"
2675 ))
2676 })?;
2677 Some(SeekCursor { sequence, id })
2678 }
2679 None => None,
2680 };
2681 filter.namespaces = token
2682 .visible_namespaces()
2683 .iter()
2684 .map(|namespace| namespace.as_str().to_owned())
2685 .collect();
2686 let page = store
2687 .query_entities_after(token.namespace().as_str(), filter, after, limit)
2688 .await?;
2689 Ok((page.items, page.next_after.map(|cursor| cursor.id)))
2690 }
2691
2692 pub async fn list_entities_tagged(
2699 &self,
2700 token: &NamespaceToken,
2701 kind: Option<&str>,
2702 domain_tag: Option<&str>,
2703 limit: u32,
2704 offset: u32,
2705 ) -> RuntimeResult<Vec<Entity>> {
2706 let ns_strs: Vec<String> = token
2707 .visible_namespaces()
2708 .iter()
2709 .map(|ns| ns.as_str().to_owned())
2710 .collect();
2711 let filter = EntityFilter {
2712 kinds: match kind {
2713 Some(k) => vec![k.to_string()],
2714 None => vec![],
2715 },
2716 tags_any: match domain_tag {
2717 Some(t) if !t.is_empty() => vec![t.to_string()],
2718 _ => vec![],
2719 },
2720 namespaces: ns_strs,
2721 ..Default::default()
2722 };
2723 let page = self
2724 .entities(token)?
2725 .query_entities_count_free(
2726 token.namespace().as_str(),
2727 filter,
2728 PageRequest {
2729 offset: offset.into(),
2730 limit,
2731 },
2732 )
2733 .await?;
2734 Ok(page.items)
2735 }
2736
2737 pub async fn count_entities_tagged(
2742 &self,
2743 token: &NamespaceToken,
2744 kind: Option<&str>,
2745 domain_tag: Option<&str>,
2746 ) -> RuntimeResult<u64> {
2747 let ns_strs: Vec<String> = token
2748 .visible_namespaces()
2749 .iter()
2750 .map(|ns| ns.as_str().to_owned())
2751 .collect();
2752 let filter = EntityFilter {
2753 kinds: match kind {
2754 Some(k) => vec![k.to_string()],
2755 None => vec![],
2756 },
2757 tags_any: match domain_tag {
2758 Some(t) if !t.is_empty() => vec![t.to_string()],
2759 _ => vec![],
2760 },
2761 namespaces: ns_strs,
2762 ..Default::default()
2763 };
2764 Ok(self
2765 .entities(token)?
2766 .count_entities(token.namespace().as_str(), filter)
2767 .await?)
2768 }
2769
2770 pub async fn list_events(
2772 &self,
2773 token: &NamespaceToken,
2774 filter: EventFilter,
2775 page: PageRequest,
2776 ) -> RuntimeResult<Page<Event>> {
2777 self.events(token)?
2778 .query_events(filter, page)
2779 .await
2780 .map_err(Into::into)
2781 }
2782
2783 pub(crate) async fn validate_edge_relation_endpoints(
2801 &self,
2802 token: &NamespaceToken,
2803 source_id: Uuid,
2804 target_id: Uuid,
2805 relation: EdgeRelation,
2806 ) -> RuntimeResult<(EdgeEndpointKind, EdgeEndpointKind)> {
2807 if source_id == target_id {
2808 return Err(RuntimeError::InvalidInput(
2809 "self-loop edges are not allowed: source_id and target_id must be different".into(),
2810 ));
2811 }
2812 if relation == EdgeRelation::Annotates {
2813 match self.resolve_edge_endpoint(token, source_id).await? {
2817 Some(Resolved::Note(_)) => {}
2818 Some(_) => {
2819 return Err(RuntimeError::InvalidInput(format!(
2820 "annotates source {source_id} must be a note"
2821 )));
2822 }
2823 None => {
2824 if self.get_edge(token, source_id).await?.is_some() {
2826 return Err(RuntimeError::InvalidInput(format!(
2827 "annotates source {source_id} must be a note"
2828 )));
2829 }
2830 return Err(RuntimeError::NotFound(format!(
2831 "link source {source_id} not found"
2832 )));
2833 }
2834 }
2835 let target_kind = match self.resolve_edge_endpoint(token, target_id).await? {
2837 Some(Resolved::Entity(_)) => EdgeEndpointKind::Entity,
2838 Some(Resolved::Note(_)) => EdgeEndpointKind::Note,
2839 Some(Resolved::Event(_)) => EdgeEndpointKind::Event,
2840 Some(Resolved::PackRecord { .. }) => {
2841 return Err(RuntimeError::InvalidInput(
2842 "pack-private record is not a valid edge endpoint for annotates".into(),
2843 ));
2844 }
2845 None => match self.get_edge(token, target_id).await {
2846 Ok(Some(_)) => EdgeEndpointKind::Edge,
2847 Ok(None) | Err(RuntimeError::NotFound(_)) => {
2848 return Err(RuntimeError::NotFound(format!(
2849 "link target {target_id} not found"
2850 )));
2851 }
2852 Err(error) => return Err(error),
2853 },
2854 };
2855 return Ok((EdgeEndpointKind::Note, target_kind));
2856 } else if crate::pack::is_special_relation(relation) {
2857 let rel_name = relation.as_str();
2861 let src = match self.resolve_edge_endpoint(token, source_id).await? {
2862 Some(r) => r,
2863 None => {
2864 if self.get_edge(token, source_id).await?.is_some() {
2865 return Err(RuntimeError::InvalidInput(format!(
2866 "{rel_name} source {source_id} must be a note or entity (got edge)"
2867 )));
2868 }
2869 return Err(RuntimeError::NotFound(format!(
2870 "link source {source_id} not found"
2871 )));
2872 }
2873 };
2874 let tgt = match self.resolve_edge_endpoint(token, target_id).await? {
2875 Some(r) => r,
2876 None => {
2877 if self.get_edge(token, target_id).await?.is_some() {
2878 return Err(RuntimeError::InvalidInput(format!(
2879 "{rel_name} target {target_id} must be a note or entity (got edge)"
2880 )));
2881 }
2882 return Err(RuntimeError::NotFound(format!(
2883 "link target {target_id} not found"
2884 )));
2885 }
2886 };
2887 return match (&src, &tgt) {
2888 (Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
2889 if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
2890 let legal_relations = accepted_entity_relations_description(
2891 &self.pack_edge_rules(),
2892 &src_e.kind,
2893 src_e.entity_type.as_deref(),
2894 &tgt_e.kind,
2895 tgt_e.entity_type.as_deref(),
2896 );
2897 let rule_hint = match relation {
2898 EdgeRelation::Supports | EdgeRelation::Refutes => {
2899 "requires concept|document|dataset|artifact -> concept \
2900 (or same-substrate note -> note)"
2901 }
2902 _ => "requires same-kind entity endpoints",
2903 };
2904 return Err(RuntimeError::InvalidInput(format!(
2905 "({}) -[{rel_name}]-> ({}) is not in the base endpoint \
2906 allowlist; {rel_name} {rule_hint}; currently legal relations for \
2907 {} -> {} under the loaded endpoint rules: {legal_relations}",
2908 src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
2909 )));
2910 }
2911 Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
2912 }
2913 (Resolved::Note(_), Resolved::Note(_)) => {
2914 Ok((EdgeEndpointKind::Note, EdgeEndpointKind::Note))
2915 }
2916 (Resolved::Event(_), _) => {
2917 return Err(RuntimeError::InvalidInput(format!(
2918 "{rel_name} does not apply to events; source {source_id} is an event"
2919 )));
2920 }
2921 (_, Resolved::Event(_)) => {
2922 return Err(RuntimeError::InvalidInput(format!(
2923 "{rel_name} does not apply to events; target {target_id} is an event"
2924 )));
2925 }
2926 (Resolved::Entity(_), Resolved::Note(_)) => {
2927 return Err(RuntimeError::InvalidInput(format!(
2928 "{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
2929 got source={source_id} (entity) target={target_id} (note)"
2930 )));
2931 }
2932 (Resolved::Note(_), Resolved::Entity(_)) => {
2933 return Err(RuntimeError::InvalidInput(format!(
2934 "{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
2935 got source={source_id} (note) target={target_id} (entity)"
2936 )));
2937 }
2938 (Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
2939 return Err(RuntimeError::InvalidInput(format!(
2940 "pack-private record is not a valid edge endpoint for {rel_name}"
2941 )));
2942 }
2943 };
2944 } else {
2945 let src_res = self.resolve_edge_endpoint(token, source_id).await?;
2952 let tgt_res = self.resolve_edge_endpoint(token, target_id).await?;
2953 let pack_rules = self.pack_edge_rules();
2954
2955 if pack_rule_allows(&pack_rules, relation, src_res.as_ref(), tgt_res.as_ref()) {
2956 let kind = |resolved: Option<&Resolved>| match resolved {
2957 Some(Resolved::Entity(_)) => Some(EdgeEndpointKind::Entity),
2958 Some(Resolved::Note(_)) => Some(EdgeEndpointKind::Note),
2959 _ => None,
2960 };
2961 return match (kind(src_res.as_ref()), kind(tgt_res.as_ref())) {
2962 (Some(source_kind), Some(target_kind)) => Ok((source_kind, target_kind)),
2963 _ => Err(RuntimeError::Internal(
2964 "pack endpoint rule admitted an unsupported substrate".into(),
2965 )),
2966 };
2967 }
2968
2969 let (src_kind, src_entity_type) = match src_res.as_ref() {
2971 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2972 Some(_) => {
2973 return Err(RuntimeError::InvalidInput(format!(
2974 "link source {source_id} must be an entity for relation {relation:?} \
2975 (only `annotates` crosses substrates)"
2976 )));
2977 }
2978 None => {
2979 if self.get_edge(token, source_id).await?.is_some() {
2980 return Err(RuntimeError::InvalidInput(format!(
2981 "link source {source_id} must be an entity for relation {relation:?} \
2982 (only `annotates` crosses substrates)"
2983 )));
2984 }
2985 return Err(RuntimeError::NotFound(format!(
2986 "link source {source_id} not found"
2987 )));
2988 }
2989 };
2990 let (tgt_kind, tgt_entity_type) = match tgt_res.as_ref() {
2991 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2992 Some(_) => {
2993 return Err(RuntimeError::InvalidInput(format!(
2994 "link target {target_id} must be an entity for relation {relation:?} \
2995 (only `annotates` crosses substrates)"
2996 )));
2997 }
2998 None => {
2999 if self.get_edge(token, target_id).await?.is_some() {
3000 return Err(RuntimeError::InvalidInput(format!(
3001 "link target {target_id} must be an entity for relation {relation:?} \
3002 (only `annotates` crosses substrates)"
3003 )));
3004 }
3005 return Err(RuntimeError::NotFound(format!(
3006 "link target {target_id} not found"
3007 )));
3008 }
3009 };
3010 if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
3011 let legal_relations = accepted_entity_relations_description(
3012 &pack_rules,
3013 src_kind,
3014 src_entity_type,
3015 tgt_kind,
3016 tgt_entity_type,
3017 );
3018 return Err(RuntimeError::InvalidInput(format!(
3019 "({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
3020 allowlist; use pack EDGE_RULES to extend the allowlist; currently legal \
3021 relations for {src_kind} -> {tgt_kind} under the loaded endpoint rules: \
3022 {legal_relations}",
3023 relation.as_str()
3024 )));
3025 }
3026 }
3027 Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
3028 }
3029
3030 pub async fn validate_link_endpoints(
3035 &self,
3036 token: &NamespaceToken,
3037 source_id: Uuid,
3038 target_id: Uuid,
3039 relation: EdgeRelation,
3040 ) -> RuntimeResult<()> {
3041 self.validate_edge_relation_endpoints(token, source_id, target_id, relation)
3042 .await
3043 .map(|_| ())
3044 }
3045
3046 pub fn validate_link_endpoints_by_resolved(
3057 &self,
3058 source_id: Uuid,
3059 target_id: Uuid,
3060 relation: EdgeRelation,
3061 src: Option<&Resolved>,
3062 tgt: Option<&Resolved>,
3063 ) -> RuntimeResult<()> {
3064 if source_id == target_id {
3065 return Err(RuntimeError::InvalidInput(
3066 "self-loop edges are not allowed: source_id and target_id must be different".into(),
3067 ));
3068 }
3069
3070 if relation == EdgeRelation::Annotates {
3071 match src {
3072 Some(Resolved::Note(_)) => {}
3073 Some(_) => {
3074 return Err(RuntimeError::InvalidInput(format!(
3075 "annotates source {source_id} must be a note"
3076 )));
3077 }
3078 None => {
3079 return Err(RuntimeError::NotFound(format!(
3080 "link source {source_id} not found"
3081 )));
3082 }
3083 }
3084 if tgt.is_none() {
3085 return Err(RuntimeError::NotFound(format!(
3086 "link target {target_id} not found"
3087 )));
3088 }
3089 return Ok(());
3090 }
3091
3092 if crate::pack::is_special_relation(relation) {
3093 let rel_name = relation.as_str();
3094 let src = src.ok_or_else(|| {
3095 RuntimeError::NotFound(format!("link source {source_id} not found"))
3096 })?;
3097 let tgt = tgt.ok_or_else(|| {
3098 RuntimeError::NotFound(format!("link target {target_id} not found"))
3099 })?;
3100 match (src, tgt) {
3101 (Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
3102 if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
3103 let legal_relations = accepted_entity_relations_description(
3104 &self.pack_edge_rules(),
3105 &src_e.kind,
3106 src_e.entity_type.as_deref(),
3107 &tgt_e.kind,
3108 tgt_e.entity_type.as_deref(),
3109 );
3110 let rule_hint = match relation {
3111 EdgeRelation::Supports | EdgeRelation::Refutes => {
3112 "requires concept|document|dataset|artifact -> concept \
3113 (or same-substrate note -> note)"
3114 }
3115 _ => "requires same-kind entity endpoints",
3116 };
3117 return Err(RuntimeError::InvalidInput(format!(
3118 "({}) -[{rel_name}]-> ({}) is not in the base endpoint \
3119 allowlist; {rel_name} {rule_hint}; currently legal relations for \
3120 {} -> {} under the loaded endpoint rules: {legal_relations}",
3121 src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
3122 )));
3123 }
3124 }
3125 (Resolved::Note(_), Resolved::Note(_)) => {}
3126 (Resolved::Entity(_), Resolved::Note(_)) => {
3127 return Err(RuntimeError::InvalidInput(format!(
3128 "{rel_name} endpoints must be the same substrate \
3129 (note→note or entity→entity); got source={source_id} (entity) \
3130 target={target_id} (note)"
3131 )));
3132 }
3133 (Resolved::Note(_), Resolved::Entity(_)) => {
3134 return Err(RuntimeError::InvalidInput(format!(
3135 "{rel_name} endpoints must be the same substrate \
3136 (note→note or entity→entity); got source={source_id} (note) \
3137 target={target_id} (entity)"
3138 )));
3139 }
3140 (Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
3141 return Err(RuntimeError::InvalidInput(format!(
3142 "pack-private record is not a valid edge endpoint for {rel_name}"
3143 )));
3144 }
3145 _ => {
3146 return Err(RuntimeError::InvalidInput(format!(
3147 "{rel_name} endpoints must be notes or entities (not events)"
3148 )));
3149 }
3150 }
3151 return Ok(());
3152 }
3153
3154 let pack_rules = self.pack_edge_rules();
3157 if pack_rule_allows(&pack_rules, relation, src, tgt) {
3158 return Ok(());
3159 }
3160
3161 let (src_kind, src_entity_type) = match src {
3162 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
3163 Some(_) => {
3164 return Err(RuntimeError::InvalidInput(format!(
3165 "link source {source_id} must be an entity for relation {relation:?} \
3166 (only `annotates` crosses substrates)"
3167 )));
3168 }
3169 None => {
3170 return Err(RuntimeError::NotFound(format!(
3171 "link source {source_id} not found"
3172 )));
3173 }
3174 };
3175 let (tgt_kind, tgt_entity_type) = match tgt {
3176 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
3177 Some(_) => {
3178 return Err(RuntimeError::InvalidInput(format!(
3179 "link target {target_id} must be an entity for relation {relation:?} \
3180 (only `annotates` crosses substrates)"
3181 )));
3182 }
3183 None => {
3184 return Err(RuntimeError::NotFound(format!(
3185 "link target {target_id} not found"
3186 )));
3187 }
3188 };
3189
3190 if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
3191 let legal_relations = accepted_entity_relations_description(
3192 &pack_rules,
3193 src_kind,
3194 src_entity_type,
3195 tgt_kind,
3196 tgt_entity_type,
3197 );
3198 return Err(RuntimeError::InvalidInput(format!(
3199 "({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
3200 allowlist; use pack EDGE_RULES to extend the allowlist; currently legal relations \
3201 for {src_kind} -> {tgt_kind} under the loaded endpoint rules: {legal_relations}",
3202 relation.as_str()
3203 )));
3204 }
3205
3206 Ok(())
3207 }
3208
3209 pub fn validate_annotates_endpoint_kinds(
3221 &self,
3222 source_id: Uuid,
3223 target_id: Uuid,
3224 source: Option<EdgeEndpointKind>,
3225 target: Option<EdgeEndpointKind>,
3226 ) -> RuntimeResult<()> {
3227 if source_id == target_id {
3228 return Err(RuntimeError::InvalidInput(
3229 "self-loop edges are not allowed: source_id and target_id must be different".into(),
3230 ));
3231 }
3232 match source {
3233 Some(EdgeEndpointKind::Note) => {}
3234 Some(_) => {
3235 return Err(RuntimeError::InvalidInput(format!(
3236 "annotates source {source_id} must be a note"
3237 )));
3238 }
3239 None => {
3240 return Err(RuntimeError::NotFound(format!(
3241 "link source {source_id} not found"
3242 )));
3243 }
3244 }
3245 if target.is_none() {
3246 return Err(RuntimeError::NotFound(format!(
3247 "link target {target_id} not found"
3248 )));
3249 }
3250 Ok(())
3251 }
3252
3253 pub async fn link(
3273 &self,
3274 token: &NamespaceToken,
3275 source_id: Uuid,
3276 target_id: Uuid,
3277 relation: EdgeRelation,
3278 weight: f64,
3279 metadata: Option<serde_json::Value>,
3280 ) -> RuntimeResult<Edge> {
3281 self.link_observed(
3282 token, source_id, target_id, relation, weight, metadata, false,
3283 )
3284 .await
3285 .map(|result| result.edge)
3286 }
3287
3288 #[allow(clippy::too_many_arguments)]
3293 pub async fn link_observed(
3294 &self,
3295 token: &NamespaceToken,
3296 source_id: Uuid,
3297 target_id: Uuid,
3298 relation: EdgeRelation,
3299 weight: f64,
3300 metadata: Option<serde_json::Value>,
3301 resurrect: bool,
3302 ) -> RuntimeResult<EdgeUpsertResult> {
3303 validate_edge_weight(weight)?;
3304 validate_edge_metadata(relation, metadata.as_ref())?;
3305 let (source_kind, target_kind) = self
3306 .validate_edge_relation_endpoints(token, source_id, target_id, relation)
3307 .await?;
3308 let (canonical_source, canonical_target) =
3309 canonical_edge_endpoints(relation, source_id, target_id);
3310 let (source_kind, target_kind) =
3311 canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
3312 let (source_id, target_id) = (canonical_source, canonical_target);
3313 let metadata = if relation == EdgeRelation::DependsOn {
3314 match (
3319 self.resolve_edge_endpoint(token, source_id).await?,
3320 self.resolve_edge_endpoint(token, target_id).await?,
3321 ) {
3322 (Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
3323 merge_dependency_kind(&src_e.kind, &tgt_e.kind, metadata)
3324 }
3325 _ => metadata,
3326 }
3327 } else {
3328 metadata
3329 };
3330 validate_edge_metadata(relation, metadata.as_ref())?;
3331 let now = chrono::Utc::now();
3332 let ns = token.namespace().as_str();
3333 let edge = Edge {
3334 id: LinkId::from(Uuid::new_v4()),
3335 namespace: ns.to_string(),
3336 source_id,
3337 target_id,
3338 relation,
3339 weight,
3340 created_at: now,
3341 updated_at: now,
3342 deleted_at: None,
3343 metadata,
3344 target_backend: None,
3345 };
3346 let attribution = crate::EventAttribution::from_token(token);
3356 let outcome = compose_graph_mutation_events(
3357 self.backend(),
3358 GraphMutationRequest::Single {
3359 request: EdgeUpsertRequest { edge, resurrect },
3360 guard_endpoints: true,
3361 },
3362 GraphMutationPreconditions::default(),
3363 Vec::new(),
3364 move |outcome| match &outcome.mutation {
3365 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
3366 Ok(vec![Self::link_mutation_event(
3367 &attribution,
3368 result,
3369 source_kind,
3370 target_kind,
3371 )])
3372 }
3373 _ => Err(Self::link_composition_shape_error(
3374 "expected a written singleton",
3375 )),
3376 },
3377 )
3378 .await?;
3379 let GraphMutationOutcome::Single(outcome) = outcome.mutation else {
3380 return Err(RuntimeError::Internal(
3381 "link: unexpected composition outcome".into(),
3382 ));
3383 };
3384 let result = match outcome {
3385 GuardedEdgeUpsertOutcome::Written(result) => result,
3386 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::MissingEndpoints(missing)) => {
3387 return Err(RuntimeError::GuardedWriteFailed(GuardedWriteFailure {
3388 entry_index: None,
3389 missing_source: missing.source.then_some(source_id),
3390 missing_target: missing.target.then_some(target_id),
3391 }));
3392 }
3393 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::ResurrectionRequired { edge }) => {
3394 return Err(RuntimeError::InvalidInput(format!(
3395 "edge {} is soft-deleted; pass resurrect=true to link explicitly",
3396 edge.id
3397 )))
3398 }
3399 };
3400 Ok(result)
3401 }
3402
3403 #[allow(clippy::too_many_arguments)]
3411 pub async fn link_with_target_backend(
3412 &self,
3413 token: &NamespaceToken,
3414 source_id: Uuid,
3415 target_id: Uuid,
3416 source_kind: EdgeEndpointKind,
3417 target_kind: EdgeEndpointKind,
3418 relation: EdgeRelation,
3419 weight: f64,
3420 metadata: Option<serde_json::Value>,
3421 target_backend: Option<String>,
3422 ) -> RuntimeResult<Edge> {
3423 self.link_with_target_backend_observed(
3424 token,
3425 source_id,
3426 target_id,
3427 source_kind,
3428 target_kind,
3429 relation,
3430 weight,
3431 metadata,
3432 target_backend,
3433 false,
3434 )
3435 .await
3436 .map(|result| result.edge)
3437 }
3438
3439 #[allow(clippy::too_many_arguments)]
3443 pub async fn link_with_target_backend_observed(
3444 &self,
3445 token: &NamespaceToken,
3446 source_id: Uuid,
3447 target_id: Uuid,
3448 source_kind: EdgeEndpointKind,
3449 target_kind: EdgeEndpointKind,
3450 relation: EdgeRelation,
3451 weight: f64,
3452 metadata: Option<serde_json::Value>,
3453 target_backend: Option<String>,
3454 resurrect: bool,
3455 ) -> RuntimeResult<EdgeUpsertResult> {
3456 validate_edge_weight(weight)?;
3457 let (canonical_source, canonical_target) =
3458 canonical_edge_endpoints(relation, source_id, target_id);
3459 let (source_kind, target_kind) =
3460 canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
3461 let (source_id, target_id) = (canonical_source, canonical_target);
3462 validate_edge_metadata(relation, metadata.as_ref())?;
3463 let now = chrono::Utc::now();
3464 let ns = token.namespace().as_str();
3465 let edge = Edge {
3466 id: LinkId::from(Uuid::new_v4()),
3467 namespace: ns.to_string(),
3468 source_id,
3469 target_id,
3470 relation,
3471 weight,
3472 created_at: now,
3473 updated_at: now,
3474 deleted_at: None,
3475 metadata,
3476 target_backend,
3477 };
3478 let attribution = crate::EventAttribution::from_token(token);
3479 let outcome = compose_graph_mutation_events(
3480 self.backend(),
3481 GraphMutationRequest::Single {
3482 request: EdgeUpsertRequest { edge, resurrect },
3483 guard_endpoints: false,
3484 },
3485 GraphMutationPreconditions::default(),
3486 Vec::new(),
3487 move |outcome| match &outcome.mutation {
3488 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
3489 Ok(vec![Self::link_mutation_event(
3490 &attribution,
3491 result,
3492 source_kind,
3493 target_kind,
3494 )])
3495 }
3496 _ => Err(Self::link_composition_shape_error(
3497 "expected a written singleton",
3498 )),
3499 },
3500 )
3501 .await?;
3502 match outcome.mutation {
3503 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => Ok(result),
3504 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Refused(
3505 EdgeUpsertRefusal::ResurrectionRequired { edge },
3506 )) => {
3507 let error = khive_storage::StorageError::Conflict {
3508 capability: khive_storage::StorageCapability::Graph,
3509 operation: "upsert_edge_observed".into(),
3510 message: format!(
3511 "edge {} is soft-deleted; explicit resurrection is required",
3512 edge.id,
3513 ),
3514 };
3515 Err(RuntimeError::InvalidInput(format!(
3516 "edge natural key is soft-deleted; pass resurrect=true to link explicitly: {error}"
3517 )))
3518 }
3519 _ => Err(RuntimeError::Internal(
3520 "link: unexpected composition outcome".into(),
3521 )),
3522 }
3523 }
3524
3525 fn link_mutation_event(
3526 attribution: &crate::EventAttribution,
3527 result: &EdgeUpsertResult,
3528 source_kind: EdgeEndpointKind,
3529 target_kind: EdgeEndpointKind,
3530 ) -> Event {
3531 let kind = match result.disposition {
3532 EdgeUpsertDisposition::Created => EventKind::LinkCreated,
3533 EdgeUpsertDisposition::Updated | EdgeUpsertDisposition::Resurrected => {
3534 EventKind::EdgeUpdated
3535 }
3536 };
3537 let edge_id = Uuid::from(result.edge.id);
3538 let mut payload = serde_json::json!({
3539 "id": edge_id,
3540 "namespace": result.edge.namespace,
3541 "mutation": result.disposition.name(),
3542 "source_id": result.edge.source_id,
3543 "target_id": result.edge.target_id,
3544 "relation": result.edge.relation,
3545 "weight": result.edge.weight,
3546 "metadata": result.edge.metadata,
3547 "previous": result.previous,
3548 });
3549 if kind == EventKind::LinkCreated {
3550 payload["source_kind"] = serde_json::json!(source_kind.name());
3551 payload["target_kind"] = serde_json::json!(target_kind.name());
3552 }
3553 attribution.stamp(
3554 Event::new(
3555 result.edge.namespace.clone(),
3556 "link",
3557 kind,
3558 SubstrateKind::Entity,
3559 "",
3560 )
3561 .with_target(edge_id)
3562 .with_payload(payload),
3563 )
3564 }
3565
3566 fn link_composition_shape_error(message: &'static str) -> khive_storage::StorageError {
3567 khive_storage::StorageError::InvalidInput {
3568 capability: khive_storage::StorageCapability::Graph,
3569 operation: "link_mutation_event".into(),
3570 message: message.into(),
3571 }
3572 }
3573
3574 pub(crate) async fn substrate_exists_in_ns(
3581 &self,
3582 token: &NamespaceToken,
3583 id: Uuid,
3584 ) -> RuntimeResult<bool> {
3585 if self.resolve(token, id).await?.is_some() {
3586 return Ok(true);
3587 }
3588 match self.get_edge_visible(token, id).await {
3589 Ok(Some(_)) => Ok(true),
3590 Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
3591 Err(err) => Err(err),
3592 }
3593 }
3594
3595 pub(crate) async fn substrate_exists_by_id(
3602 &self,
3603 token: &NamespaceToken,
3604 id: Uuid,
3605 ) -> RuntimeResult<bool> {
3606 if self.resolve_edge_endpoint(token, id).await?.is_some() {
3607 return Ok(true);
3608 }
3609 match self.get_edge(token, id).await {
3610 Ok(Some(_)) => Ok(true),
3611 Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
3612 Err(err) => Err(err),
3613 }
3614 }
3615
3616 pub async fn latest_annotating_note(
3621 &self,
3622 token: &NamespaceToken,
3623 node_id: Uuid,
3624 kind: &str,
3625 tag: &str,
3626 ) -> RuntimeResult<Option<Uuid>> {
3627 self.latest_annotating_note_inner(token, node_id, kind, tag, None)
3628 .await
3629 }
3630
3631 pub async fn latest_annotating_note_with_property(
3634 &self,
3635 token: &NamespaceToken,
3636 node_id: Uuid,
3637 kind: &str,
3638 tag: &str,
3639 property_key: &str,
3640 property_value: &str,
3641 ) -> RuntimeResult<Option<Uuid>> {
3642 self.latest_annotating_note_inner(
3643 token,
3644 node_id,
3645 kind,
3646 tag,
3647 Some((property_key, property_value)),
3648 )
3649 .await
3650 }
3651
3652 async fn latest_annotating_note_inner(
3653 &self,
3654 token: &NamespaceToken,
3655 node_id: Uuid,
3656 kind: &str,
3657 tag: &str,
3658 required_property: Option<(&str, &str)>,
3659 ) -> RuntimeResult<Option<Uuid>> {
3660 if !self.substrate_exists_in_ns(token, node_id).await? {
3661 return Ok(None);
3662 }
3663 let mut latest: Option<(Uuid, i64)> = None;
3664 for namespace in token.visible_namespaces() {
3665 let scoped = NamespaceToken::for_namespace(namespace.clone());
3666 let graph = self.graph(&scoped)?;
3667 let candidate = match required_property {
3668 Some((key, value)) => {
3669 graph
3670 .latest_annotating_note_with_property(node_id, kind, tag, key, value)
3671 .await?
3672 }
3673 None => graph.latest_annotating_note(node_id, kind, tag).await?,
3674 };
3675 if let Some(candidate) = candidate {
3676 if latest.is_none_or(|(id, created_at)| {
3677 candidate.1 > created_at || (candidate.1 == created_at && candidate.0 < id)
3678 }) {
3679 latest = Some(candidate);
3680 }
3681 }
3682 }
3683 Ok(latest.map(|(id, _)| id))
3684 }
3685
3686 pub async fn neighbors(
3695 &self,
3696 token: &NamespaceToken,
3697 node_id: Uuid,
3698 direction: Direction,
3699 limit: Option<u32>,
3700 relations: Option<Vec<EdgeRelation>>,
3701 ) -> RuntimeResult<Vec<NeighborHit>> {
3702 self.neighbors_with_query(
3703 token,
3704 node_id,
3705 NeighborQuery {
3706 direction,
3707 relations,
3708 limit,
3709 min_weight: None,
3710 },
3711 )
3712 .await
3713 }
3714
3715 pub async fn neighbors_with_query(
3725 &self,
3726 token: &NamespaceToken,
3727 node_id: Uuid,
3728 query: NeighborQuery,
3729 ) -> RuntimeResult<Vec<NeighborHit>> {
3730 self.neighbors_with_query_page(token, node_id, query, None, None, true)
3731 .await
3732 }
3733
3734 pub async fn neighbors_with_query_page(
3738 &self,
3739 token: &NamespaceToken,
3740 node_id: Uuid,
3741 query: NeighborQuery,
3742 after: Option<NeighborCursor>,
3743 neighbor_kinds: Option<Vec<String>>,
3744 enrich: bool,
3745 ) -> RuntimeResult<Vec<NeighborHit>> {
3746 if !self.substrate_exists_by_id(token, node_id).await? {
3749 return Err(RuntimeError::NotFound(format!(
3750 "neighbor anchor {node_id} not found"
3751 )));
3752 }
3753
3754 self.neighbors_for_resolved_kg_read(
3755 token,
3756 node_id,
3757 crate::KgNeighborRead {
3758 query,
3759 after,
3760 neighbor_kinds,
3761 enrich,
3762 namespace: None,
3763 },
3764 )
3765 .await
3766 }
3767
3768 pub async fn neighbors_for_resolved_kg_read(
3776 &self,
3777 token: &NamespaceToken,
3778 node_id: Uuid,
3779 options: crate::KgNeighborRead,
3780 ) -> RuntimeResult<Vec<NeighborHit>> {
3781 self.neighbors_for_resolved_kg_read_inner(token, node_id, options, false)
3782 .await
3783 .map(|(hits, _)| hits)
3784 }
3785
3786 pub async fn neighbors_for_resolved_kg_read_with_entity_kinds(
3791 &self,
3792 token: &NamespaceToken,
3793 node_id: Uuid,
3794 options: crate::KgNeighborRead,
3795 ) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
3796 self.neighbors_for_resolved_kg_read_inner(token, node_id, options, true)
3797 .await
3798 }
3799
3800 async fn neighbors_for_resolved_kg_read_inner(
3801 &self,
3802 token: &NamespaceToken,
3803 node_id: Uuid,
3804 options: crate::KgNeighborRead,
3805 with_entity_kinds: bool,
3806 ) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
3807 let crate::KgNeighborRead {
3808 mut query,
3809 after,
3810 neighbor_kinds,
3811 enrich,
3812 namespace,
3813 } = options;
3814 let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace.as_ref())?;
3815 query.direction =
3816 normalize_symmetric_direction(query.direction, query.relations.as_deref());
3817 let mut hits = Vec::new();
3818 for ns in namespaces {
3819 let temp = NamespaceToken::for_namespace(ns.clone());
3820 let mut ns_hits = self
3821 .graph(&temp)?
3822 .neighbors_page(node_id, query.clone(), after, neighbor_kinds.clone())
3823 .await?;
3824 hits.append(&mut ns_hits);
3825 }
3826 hits.sort_by_key(|h| (h.node_id, h.edge_id));
3827 hits.dedup_by_key(|h| (h.node_id, h.edge_id));
3828 if enrich {
3829 self.enrich_neighbor_hits(token, &mut hits).await;
3830 }
3831 let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.node_id).collect();
3833 let (deleted, entity_kinds) = self
3834 .neighbor_node_screen(
3835 candidate_ids,
3836 (with_entity_kinds && !enrich).then_some(token),
3837 )
3838 .await?;
3839 if !deleted.is_empty() {
3840 hits.retain(|h| !deleted.contains(&h.node_id));
3841 }
3842 hits.sort_by(|a, b| {
3849 b.weight
3850 .partial_cmp(&a.weight)
3851 .unwrap_or(std::cmp::Ordering::Equal)
3852 .then(a.node_id.cmp(&b.node_id))
3853 .then(a.edge_id.cmp(&b.edge_id))
3854 });
3855 Ok((hits, entity_kinds))
3856 }
3857
3858 pub async fn annotation_neighbors_by_target_id(
3866 &self,
3867 target_id: Uuid,
3868 ) -> RuntimeResult<Vec<NeighborHit>> {
3869 let mut reader = self.sql().reader().await?;
3870 let rows = reader
3871 .query_all(SqlStatement {
3872 sql: "SELECT source_id, id, weight FROM graph_edges \
3873 WHERE target_id = ?1 AND relation = ?2 AND deleted_at IS NULL \
3874 ORDER BY weight DESC, source_id ASC"
3875 .to_string(),
3876 params: vec![
3877 SqlValue::Text(target_id.to_string()),
3878 SqlValue::Text(EdgeRelation::Annotates.to_string()),
3879 ],
3880 label: Some("annotations.by_target_id_unfiltered".into()),
3881 })
3882 .await?;
3883
3884 rows.into_iter()
3885 .map(|row| {
3886 let parse_uuid = |name: &str| match row.get(name) {
3887 Some(SqlValue::Text(value)) => Uuid::from_str(value).map_err(|error| {
3888 RuntimeError::Internal(format!("graph_edges.{name} is not a UUID: {error}"))
3889 }),
3890 Some(value) => Err(RuntimeError::Internal(format!(
3891 "graph_edges.{name} has unexpected SQL value {value:?}"
3892 ))),
3893 None => Err(RuntimeError::Internal(format!(
3894 "graph_edges row missing {name}"
3895 ))),
3896 };
3897 let weight = match row.get("weight") {
3898 Some(SqlValue::Float(value)) => Ok(*value),
3899 Some(value) => Err(RuntimeError::Internal(format!(
3900 "graph_edges.weight has unexpected SQL value {value:?}"
3901 ))),
3902 None => Err(RuntimeError::Internal(
3903 "graph_edges row missing weight".into(),
3904 )),
3905 }?;
3906
3907 Ok(NeighborHit {
3908 node_id: parse_uuid("source_id")?,
3909 edge_id: parse_uuid("id")?,
3910 relation: EdgeRelation::Annotates,
3911 weight,
3912 name: None,
3913 kind: None,
3914 entity_type: None,
3915 })
3916 })
3917 .collect()
3918 }
3919
3920 pub async fn neighbors_with_query_directed(
3929 &self,
3930 token: &NamespaceToken,
3931 node_id: Uuid,
3932 query: NeighborQuery,
3933 ) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
3934 if !self.substrate_exists_by_id(token, node_id).await? {
3935 return Err(RuntimeError::NotFound(format!(
3936 "neighbor anchor {node_id} not found"
3937 )));
3938 }
3939
3940 self.directed_neighbors_for_resolved_kg_read(token, node_id, query, None)
3941 .await
3942 }
3943
3944 pub(crate) async fn directed_neighbors_for_resolved_kg_read(
3945 &self,
3946 token: &NamespaceToken,
3947 node_id: Uuid,
3948 query: NeighborQuery,
3949 namespace: Option<&crate::Namespace>,
3950 ) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
3951 let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace)?;
3952 let mut hits: Vec<DirectedNeighborHit> = Vec::new();
3953 for ns in namespaces {
3954 let temp = NamespaceToken::for_namespace(ns.clone());
3955 let mut ns_hits = self
3956 .graph(&temp)?
3957 .neighbors_both_directions(node_id, query.clone())
3958 .await?;
3959 hits.append(&mut ns_hits);
3960 }
3961 hits.sort_by_key(|h| {
3965 (
3966 h.hit.node_id,
3967 h.hit.edge_id,
3968 direction_sort_rank(&h.direction),
3969 )
3970 });
3971 hits.dedup_by_key(|h| {
3972 (
3973 h.hit.node_id,
3974 h.hit.edge_id,
3975 direction_sort_rank(&h.direction),
3976 )
3977 });
3978
3979 let mut plain_hits: Vec<NeighborHit> = hits.iter().map(|h| h.hit.clone()).collect();
3980 self.enrich_neighbor_hits(token, &mut plain_hits).await;
3981 for (dh, enriched) in hits.iter_mut().zip(plain_hits) {
3982 dh.hit = enriched;
3983 }
3984
3985 let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.hit.node_id).collect();
3987 let deleted = self.deleted_entity_ids(candidate_ids).await?;
3988 if !deleted.is_empty() {
3989 hits.retain(|h| !deleted.contains(&h.hit.node_id));
3990 }
3991 hits.sort_by(|a, b| {
3995 b.hit
3996 .weight
3997 .partial_cmp(&a.hit.weight)
3998 .unwrap_or(std::cmp::Ordering::Equal)
3999 .then(a.hit.node_id.cmp(&b.hit.node_id))
4000 .then(a.hit.edge_id.cmp(&b.hit.edge_id))
4001 });
4002 Ok(hits.into_iter().map(|h| (h.hit, h.direction)).collect())
4003 }
4004
4005 pub async fn traverse(
4011 &self,
4012 token: &NamespaceToken,
4013 request: TraversalRequest,
4014 ) -> RuntimeResult<Vec<GraphPath>> {
4015 let mut request = request;
4016 request.validate().map_err(RuntimeError::InvalidInput)?;
4017 let mut roots = Vec::with_capacity(request.roots.len());
4018 let mut seen_roots = std::collections::HashSet::with_capacity(request.roots.len());
4019 for root in request.roots.drain(..) {
4020 if seen_roots.insert(root) {
4021 if !self.substrate_exists_by_id(token, root).await? {
4022 return Err(RuntimeError::NotFound(format!(
4023 "traverse root {root} not found"
4024 )));
4025 }
4026 roots.push(root);
4027 }
4028 }
4029 request.roots = roots;
4030 if request.roots.is_empty() {
4031 return Ok(Vec::new());
4032 }
4033
4034 let mut paths = Vec::new();
4035 for ns in token.visible_namespaces() {
4036 let temp = NamespaceToken::for_namespace(ns.clone());
4037 let mut ns_paths = self.graph(&temp)?.traverse(request.clone()).await?;
4038 paths.append(&mut ns_paths);
4039 }
4040 let mut paths =
4044 merge_traversal_paths_by_root(paths, Some(request.options.effective_limit()));
4045 self.enrich_path_nodes(token, &mut paths, request.include_properties)
4046 .await;
4047 let all_node_ids: Vec<Uuid> = paths
4049 .iter()
4050 .flat_map(|p| p.nodes.iter().map(|n| n.node_id))
4051 .collect();
4052 let deleted = self.deleted_entity_ids(all_node_ids).await?;
4053 if !deleted.is_empty() {
4054 for path in paths.iter_mut() {
4055 path.nodes.retain(|n| !deleted.contains(&n.node_id));
4056 recompute_total_weight(path);
4057 }
4058 paths.retain(|p| !p.nodes.is_empty());
4059 }
4060 Ok(paths)
4061 }
4062
4063 async fn deleted_entity_ids(
4083 &self,
4084 ids: Vec<Uuid>,
4085 ) -> RuntimeResult<std::collections::HashSet<Uuid>> {
4086 self.neighbor_node_screen(ids, None)
4087 .await
4088 .map(|(deleted, _)| deleted)
4089 }
4090
4091 async fn neighbor_node_screen(
4095 &self,
4096 ids: Vec<Uuid>,
4097 kind_token: Option<&NamespaceToken>,
4098 ) -> RuntimeResult<(std::collections::HashSet<Uuid>, HashMap<Uuid, String>)> {
4099 if ids.is_empty() {
4100 return Ok((std::collections::HashSet::new(), HashMap::new()));
4101 }
4102 let id_strs: Vec<String> = ids.iter().map(|u| u.to_string()).collect();
4103 let n = id_strs.len();
4104 let entities_placeholders = (0..n)
4111 .map(|i| format!("?{}", i + 1))
4112 .collect::<Vec<_>>()
4113 .join(",");
4114 let notes_placeholders = (0..n)
4115 .map(|i| format!("?{}", n + i + 1))
4116 .collect::<Vec<_>>()
4117 .join(",");
4118 let sql_str = if kind_token.is_some() {
4119 format!(
4120 "SELECT id, kind, namespace, deleted_at IS NOT NULL AS is_deleted \
4121 FROM entities WHERE id IN ({entities_placeholders}) \
4122 UNION ALL \
4123 SELECT id, NULL, NULL, 1 FROM notes \
4124 WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
4125 )
4126 } else {
4127 format!(
4128 "SELECT id FROM entities WHERE id IN ({entities_placeholders}) AND deleted_at IS NOT NULL \
4129 UNION \
4130 SELECT id FROM notes WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
4131 )
4132 };
4133 let params: Vec<SqlValue> = id_strs
4135 .iter()
4136 .chain(id_strs.iter())
4137 .cloned()
4138 .map(SqlValue::Text)
4139 .collect();
4140 let stmt = SqlStatement {
4141 sql: sql_str,
4142 params,
4143 label: Some("deleted_entity_ids".into()),
4144 };
4145 let mut out = std::collections::HashSet::new();
4146 let mut entity_kinds = HashMap::new();
4147 let sql = self.sql();
4148 let mut reader = sql.reader().await?;
4149 let rows = reader.query_all(stmt).await?;
4150 for row in rows {
4151 if let Some(col) = row.columns.first() {
4152 if let SqlValue::Text(s) = &col.value {
4153 if let Ok(u) = s.parse::<Uuid>() {
4154 if kind_token.is_none()
4155 || matches!(
4156 row.columns.get(3).map(|col| &col.value),
4157 Some(SqlValue::Integer(1))
4158 )
4159 {
4160 out.insert(u);
4161 } else if let (
4162 Some(token),
4163 Some(SqlValue::Text(kind)),
4164 Some(SqlValue::Text(namespace)),
4165 Some(SqlValue::Integer(0)),
4166 ) = (
4167 kind_token,
4168 row.columns.get(1).map(|col| &col.value),
4169 row.columns.get(2).map(|col| &col.value),
4170 row.columns.get(3).map(|col| &col.value),
4171 ) {
4172 if token
4173 .visible_namespaces()
4174 .iter()
4175 .any(|ns| ns.as_str() == namespace.as_str())
4176 {
4177 entity_kinds.insert(u, kind.clone());
4178 }
4179 }
4180 }
4181 }
4182 }
4183 }
4184 Ok((out, entity_kinds))
4185 }
4186
4187 async fn enrich_neighbor_hits(&self, token: &NamespaceToken, hits: &mut [NeighborHit]) {
4196 if hits.is_empty() {
4197 return;
4198 }
4199
4200 let unique_ids: Vec<Uuid> = {
4202 let mut seen = std::collections::HashSet::new();
4203 hits.iter()
4204 .filter_map(|h| {
4205 if seen.insert(h.node_id) {
4206 Some(h.node_id)
4207 } else {
4208 None
4209 }
4210 })
4211 .collect()
4212 };
4213
4214 let entity_map: HashMap<Uuid, Entity> = self
4215 .get_entities_by_ids_visible(token, &unique_ids)
4216 .await
4217 .unwrap_or_default()
4218 .into_iter()
4219 .map(|e| (e.id, e))
4220 .collect();
4221
4222 let residual_ids: Vec<Uuid> = unique_ids
4224 .iter()
4225 .filter(|id| !entity_map.contains_key(id))
4226 .copied()
4227 .collect();
4228
4229 let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
4230 if let Ok(store) = self.notes(token) {
4231 store
4232 .get_notes_batch(&residual_ids)
4233 .await
4234 .unwrap_or_default()
4235 .into_iter()
4236 .map(|n| (n.id, n))
4237 .collect()
4238 } else {
4239 HashMap::new()
4240 }
4241 } else {
4242 HashMap::new()
4243 };
4244
4245 for hit in hits.iter_mut() {
4246 if let Some(entity) = entity_map.get(&hit.node_id) {
4247 hit.name = Some(entity.name.clone());
4248 hit.kind = Some(entity.kind.clone());
4249 hit.entity_type = entity.entity_type.clone();
4250 } else if let Some(note) = note_map.get(&hit.node_id) {
4251 hit.name = Some(note_graph_name(note));
4252 hit.kind = Some(note.kind.clone());
4253 }
4254 }
4255 }
4256
4257 async fn enrich_path_nodes(
4268 &self,
4269 token: &NamespaceToken,
4270 paths: &mut [GraphPath],
4271 include_properties: bool,
4272 ) {
4273 if paths.is_empty() {
4274 return;
4275 }
4276
4277 let unique_ids: Vec<Uuid> = {
4279 let mut seen = std::collections::HashSet::new();
4280 paths
4281 .iter()
4282 .flat_map(|p| p.nodes.iter())
4283 .filter_map(|n| {
4284 if seen.insert(n.node_id) {
4285 Some(n.node_id)
4286 } else {
4287 None
4288 }
4289 })
4290 .collect()
4291 };
4292
4293 let entity_map: HashMap<Uuid, Entity> = self
4294 .get_entities_by_ids_visible(token, &unique_ids)
4295 .await
4296 .unwrap_or_default()
4297 .into_iter()
4298 .map(|e| (e.id, e))
4299 .collect();
4300
4301 let residual_ids: Vec<Uuid> = unique_ids
4302 .iter()
4303 .filter(|id| !entity_map.contains_key(id))
4304 .copied()
4305 .collect();
4306
4307 let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
4308 if let Ok(store) = self.notes(token) {
4309 store
4310 .get_notes_batch(&residual_ids)
4311 .await
4312 .unwrap_or_default()
4313 .into_iter()
4314 .map(|n| (n.id, n))
4315 .collect()
4316 } else {
4317 HashMap::new()
4318 }
4319 } else {
4320 HashMap::new()
4321 };
4322
4323 for path in paths.iter_mut() {
4324 for node in path.nodes.iter_mut() {
4325 if let Some(entity) = entity_map.get(&node.node_id) {
4326 node.name = Some(entity.name.clone());
4327 node.kind = Some(entity.kind.clone());
4328 if include_properties {
4329 node.properties = entity.properties.clone();
4330 }
4331 } else if let Some(note) = note_map.get(&node.node_id) {
4332 node.name = Some(note_graph_name(note));
4333 node.kind = Some(note.kind.clone());
4334 }
4335 }
4336 }
4337 }
4338
4339 #[allow(clippy::too_many_arguments)]
4357 pub async fn create_note(
4358 &self,
4359 token: &NamespaceToken,
4360 kind: &str,
4361 name: Option<&str>,
4362 content: &str,
4363 salience: Option<f64>,
4364 properties: Option<serde_json::Value>,
4365 annotates: Vec<Uuid>,
4366 ) -> RuntimeResult<Note> {
4367 let (note, embedding, degradations, _) = self
4368 .create_note_inner(
4369 token, kind, name, content, None, salience, None, properties, annotates, None,
4370 false, false,
4371 )
4372 .await?;
4373 legacy_post_commit_result_with_embedding(
4374 "create_note",
4375 note.id,
4376 note,
4377 embedding,
4378 degradations,
4379 )
4380 }
4381
4382 pub async fn create_web_receipt_note(
4387 &self,
4388 token: &NamespaceToken,
4389 summary: &str,
4390 request: serde_json::Value,
4391 annotates: Vec<Uuid>,
4392 ) -> RuntimeResult<Note> {
4393 let properties = serde_json::json!({
4394 "tags": ["web.receipt"],
4395 "request": request,
4396 });
4397 let (note, _, degradations, _) = self
4398 .create_note_inner(
4399 token,
4400 "observation",
4401 None,
4402 summary,
4403 None,
4404 None,
4405 None,
4406 Some(properties),
4407 annotates,
4408 None,
4409 false,
4410 true,
4411 )
4412 .await?;
4413 legacy_post_commit_result("create_web_receipt_note", note.id, note, degradations)
4414 }
4415
4416 #[allow(clippy::too_many_arguments)]
4427 pub async fn create_note_with_embedding_content(
4428 &self,
4429 token: &NamespaceToken,
4430 kind: &str,
4431 name: Option<&str>,
4432 content: &str,
4433 embedding_content: Option<&str>,
4434 salience: Option<f64>,
4435 properties: Option<serde_json::Value>,
4436 annotates: Vec<Uuid>,
4437 ) -> RuntimeResult<Note> {
4438 let (note, embedding, degradations, _) = self
4439 .create_note_inner(
4440 token,
4441 kind,
4442 name,
4443 content,
4444 embedding_content,
4445 salience,
4446 None,
4447 properties,
4448 annotates,
4449 None,
4450 false,
4451 false,
4452 )
4453 .await?;
4454 legacy_post_commit_result_with_embedding(
4455 "create_note_with_embedding_content",
4456 note.id,
4457 note,
4458 embedding,
4459 degradations,
4460 )
4461 }
4462
4463 #[allow(clippy::too_many_arguments)]
4464 pub async fn create_note_with_embedding_content_and_report(
4465 &self,
4466 token: &NamespaceToken,
4467 kind: &str,
4468 name: Option<&str>,
4469 content: &str,
4470 embedding_content: Option<&str>,
4471 salience: Option<f64>,
4472 properties: Option<serde_json::Value>,
4473 annotates: Vec<Uuid>,
4474 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4475 let (note, embedding, degradations, _) = self
4476 .create_note_inner(
4477 token,
4478 kind,
4479 name,
4480 content,
4481 embedding_content,
4482 salience,
4483 None,
4484 properties,
4485 annotates,
4486 None,
4487 false,
4488 false,
4489 )
4490 .await?;
4491 legacy_post_commit_result(
4492 "create_note_with_embedding_content_and_report",
4493 note.id,
4494 (note, embedding),
4495 degradations,
4496 )
4497 }
4498
4499 #[allow(clippy::too_many_arguments)]
4501 pub async fn create_note_with_embedding_content_and_post_commit_report(
4502 &self,
4503 token: &NamespaceToken,
4504 kind: &str,
4505 name: Option<&str>,
4506 content: &str,
4507 embedding_content: Option<&str>,
4508 salience: Option<f64>,
4509 properties: Option<serde_json::Value>,
4510 annotates: Vec<Uuid>,
4511 ) -> RuntimeResult<(
4512 Note,
4513 crate::retrieval::EmbeddingTruncationReport,
4514 Vec<PostCommitDegradation>,
4515 )> {
4516 let (note, embedding, degradations, _) = self
4517 .create_note_inner(
4518 token,
4519 kind,
4520 name,
4521 content,
4522 embedding_content,
4523 salience,
4524 None,
4525 properties,
4526 annotates,
4527 None,
4528 false,
4529 false,
4530 )
4531 .await?;
4532 Ok((note, embedding, degradations))
4533 }
4534
4535 #[allow(clippy::too_many_arguments)]
4539 pub async fn create_note_with_decay(
4540 &self,
4541 token: &NamespaceToken,
4542 kind: &str,
4543 name: Option<&str>,
4544 content: &str,
4545 salience: Option<f64>,
4546 decay_factor: f64,
4547 properties: Option<serde_json::Value>,
4548 annotates: Vec<Uuid>,
4549 ) -> RuntimeResult<Note> {
4550 self.create_note_with_decay_for_embedding_model(
4551 token,
4552 kind,
4553 name,
4554 content,
4555 salience,
4556 decay_factor,
4557 properties,
4558 annotates,
4559 None,
4560 )
4561 .await
4562 }
4563
4564 #[allow(clippy::too_many_arguments)]
4566 pub async fn create_note_with_decay_and_report(
4567 &self,
4568 token: &NamespaceToken,
4569 kind: &str,
4570 name: Option<&str>,
4571 content: &str,
4572 salience: Option<f64>,
4573 decay_factor: f64,
4574 properties: Option<serde_json::Value>,
4575 annotates: Vec<Uuid>,
4576 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4577 self.create_note_with_decay_for_embedding_model_and_report(
4578 token,
4579 kind,
4580 name,
4581 content,
4582 salience,
4583 decay_factor,
4584 properties,
4585 annotates,
4586 None,
4587 )
4588 .await
4589 }
4590
4591 #[allow(clippy::too_many_arguments)]
4596 pub async fn create_note_with_decay_for_embedding_model(
4597 &self,
4598 token: &NamespaceToken,
4599 kind: &str,
4600 name: Option<&str>,
4601 content: &str,
4602 salience: Option<f64>,
4603 decay_factor: f64,
4604 properties: Option<serde_json::Value>,
4605 annotates: Vec<Uuid>,
4606 embedding_model: Option<&str>,
4607 ) -> RuntimeResult<Note> {
4608 let (note, embedding, degradations, _) = self
4609 .create_note_inner(
4610 token,
4611 kind,
4612 name,
4613 content,
4614 None,
4615 salience,
4616 Some(decay_factor),
4617 properties,
4618 annotates,
4619 embedding_model,
4620 false,
4621 false,
4622 )
4623 .await?;
4624 legacy_post_commit_result_with_embedding(
4625 "create_note_with_decay_for_embedding_model",
4626 note.id,
4627 note,
4628 embedding,
4629 degradations,
4630 )
4631 }
4632
4633 #[allow(clippy::too_many_arguments)]
4635 pub async fn create_note_with_decay_for_embedding_model_and_report(
4636 &self,
4637 token: &NamespaceToken,
4638 kind: &str,
4639 name: Option<&str>,
4640 content: &str,
4641 salience: Option<f64>,
4642 decay_factor: f64,
4643 properties: Option<serde_json::Value>,
4644 annotates: Vec<Uuid>,
4645 embedding_model: Option<&str>,
4646 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4647 let (note, embedding, degradations, _) = self
4648 .create_note_inner(
4649 token,
4650 kind,
4651 name,
4652 content,
4653 None,
4654 salience,
4655 Some(decay_factor),
4656 properties,
4657 annotates,
4658 embedding_model,
4659 false,
4660 false,
4661 )
4662 .await?;
4663 legacy_post_commit_result(
4664 "create_note_with_decay_for_embedding_model_and_report",
4665 note.id,
4666 (note, embedding),
4667 degradations,
4668 )
4669 }
4670
4671 #[allow(clippy::too_many_arguments)]
4675 pub async fn create_note_with_decay_for_embedding_model_with_visibility(
4676 &self,
4677 token: &NamespaceToken,
4678 kind: &str,
4679 name: Option<&str>,
4680 content: &str,
4681 salience: Option<f64>,
4682 decay_factor: f64,
4683 properties: Option<serde_json::Value>,
4684 annotates: Vec<Uuid>,
4685 embedding_model: Option<&str>,
4686 ) -> RuntimeResult<(Note, Vec<(String, u64)>)> {
4687 let (note, _, degradations, fences) = self
4688 .create_note_inner(
4689 token,
4690 kind,
4691 name,
4692 content,
4693 None,
4694 salience,
4695 Some(decay_factor),
4696 properties,
4697 annotates,
4698 embedding_model,
4699 true,
4700 false,
4701 )
4702 .await?;
4703 legacy_post_commit_result(
4704 "create_note_with_decay_for_embedding_model_with_visibility",
4705 note.id,
4706 (note, fences),
4707 degradations,
4708 )
4709 }
4710
4711 #[allow(clippy::too_many_arguments)]
4715 pub async fn create_note_with_decay_for_embedding_model_with_visibility_and_report(
4716 &self,
4717 token: &NamespaceToken,
4718 kind: &str,
4719 name: Option<&str>,
4720 content: &str,
4721 salience: Option<f64>,
4722 decay_factor: f64,
4723 properties: Option<serde_json::Value>,
4724 annotates: Vec<Uuid>,
4725 embedding_model: Option<&str>,
4726 ) -> RuntimeResult<(
4727 Note,
4728 Vec<(String, u64)>,
4729 crate::retrieval::EmbeddingTruncationReport,
4730 )> {
4731 let (note, embedding, degradations, fences) = self
4732 .create_note_inner(
4733 token,
4734 kind,
4735 name,
4736 content,
4737 None,
4738 salience,
4739 Some(decay_factor),
4740 properties,
4741 annotates,
4742 embedding_model,
4743 true,
4744 false,
4745 )
4746 .await?;
4747 legacy_post_commit_result(
4748 "create_note_with_decay_for_embedding_model_with_visibility_and_report",
4749 note.id,
4750 (note, fences, embedding),
4751 degradations,
4752 )
4753 }
4754
4755 pub async fn try_create_note(
4778 &self,
4779 token: &NamespaceToken,
4780 kind: &str,
4781 name: Option<&str>,
4782 content: &str,
4783 properties: Option<serde_json::Value>,
4784 ) -> RuntimeResult<Option<Note>> {
4785 self.try_create_note_impl(token, kind, name, content, properties, false, None, None)
4786 .await
4787 }
4788
4789 #[allow(clippy::too_many_arguments)]
4804 pub async fn try_create_note_as_trusted_ingest(
4805 &self,
4806 _capability: &crate::pack::ChannelIngestCapability,
4807 token: &NamespaceToken,
4808 kind: &str,
4809 name: Option<&str>,
4810 content: &str,
4811 properties: Option<serde_json::Value>,
4812 expires_after: Option<std::time::Duration>,
4813 ) -> RuntimeResult<Option<Note>> {
4814 self.try_create_note_impl(
4815 token,
4816 kind,
4817 name,
4818 content,
4819 properties,
4820 true,
4821 None,
4822 expires_after,
4823 )
4824 .await
4825 }
4826
4827 #[allow(clippy::too_many_arguments)]
4831 pub async fn try_create_note_as_trusted_ingest_with_attachment(
4832 &self,
4833 _capability: &crate::pack::ChannelIngestCapability,
4834 token: &NamespaceToken,
4835 kind: &str,
4836 name: Option<&str>,
4837 content: &str,
4838 properties: Option<serde_json::Value>,
4839 attachment: NewAttachment,
4840 expires_after: Option<std::time::Duration>,
4841 ) -> RuntimeResult<Option<Note>> {
4842 self.try_create_note_impl(
4843 token,
4844 kind,
4845 name,
4846 content,
4847 properties,
4848 true,
4849 Some(attachment),
4850 expires_after,
4851 )
4852 .await
4853 }
4854
4855 #[allow(clippy::too_many_arguments)]
4856 async fn try_create_note_impl(
4857 &self,
4858 token: &NamespaceToken,
4859 kind: &str,
4860 name: Option<&str>,
4861 content: &str,
4862 properties: Option<serde_json::Value>,
4863 allow_transport_owned_message_properties: bool,
4864 attachment: Option<NewAttachment>,
4865 expires_after: Option<std::time::Duration>,
4866 ) -> RuntimeResult<Option<Note>> {
4867 self.validate_note_kind(kind)?;
4868 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
4869 crate::secret_gate::check_at(content, "note", "content")?;
4870 if let Some(n) = name {
4871 crate::secret_gate::check_at(n, "note", "name")?;
4872 }
4873 if let Some(ref p) = properties {
4874 crate::secret_gate::check_json_at(p, "note", "properties")?;
4875 }
4876 if let Some(ref attachment) = attachment {
4877 drop(self.attachments()?);
4880 attachment.validate()?;
4881 let blob_store = self.blob_store().ok_or_else(|| {
4882 RuntimeError::Unconfigured(
4883 "trusted ingest attachment requires an installed BlobStore".to_string(),
4884 )
4885 })?;
4886 if !blob_store.exists(&attachment.content_ref).await? {
4887 return Err(RuntimeError::InvalidInput(format!(
4888 "trusted ingest attachment refers to an unpublished blob: {}",
4889 attachment.content_ref
4890 )));
4891 }
4892 }
4893 if !allow_transport_owned_message_properties && kind == "message" {
4894 if let Some(key) = properties
4895 .as_ref()
4896 .and_then(serde_json::Value::as_object)
4897 .and_then(|supplied| {
4898 crate::curation::kind_owned_properties("message")
4899 .iter()
4900 .copied()
4901 .find(|key| supplied.contains_key(*key))
4902 })
4903 {
4904 return Err(RuntimeError::InvalidInput(format!(
4905 "`{key}` is transport-owned on a `message` note and cannot be supplied \
4906 through `try_create_note`; only the trusted channel-ingest path may \
4907 establish quarantine disposition and channel provenance"
4908 )));
4909 }
4910 }
4911
4912 let ns = token.namespace().as_str();
4913 let mut note = Note::new(ns, kind, content);
4914 if let Some(retention) = expires_after {
4915 let duration_us = i64::try_from(retention.as_micros()).map_err(|_| {
4916 RuntimeError::InvalidInput(
4917 "trusted ingest retention exceeds i64 microseconds".into(),
4918 )
4919 })?;
4920 note.expires_at = Some(note.created_at.checked_add(duration_us).ok_or_else(|| {
4921 RuntimeError::InvalidInput("trusted ingest expiry exceeds i64 microseconds".into())
4922 })?);
4923 }
4924 if let Some(n) = name {
4925 note = note.with_name(n);
4926 }
4927 if let Some(p) = properties {
4928 note = note.with_properties(p);
4929 }
4930
4931 let inserted = if let Some(attachment) = attachment {
4938 self.raw_notes(token)?
4939 .try_insert_note_with_attachments(
4940 note.clone(),
4941 vec![Attachment::from_new(
4942 note.id,
4943 AttachmentSubstrate::Note,
4944 attachment,
4945 note.created_at,
4946 )],
4947 )
4948 .await?
4949 } else {
4950 self.raw_notes(token)?.try_insert_note(note.clone()).await?
4951 };
4952 if !inserted {
4953 return Ok(None);
4954 }
4955
4956 let mut degradations = Vec::new();
4957 match self.text_for_notes(token) {
4958 Ok(fts) => {
4959 if let Err(error) = fts.upsert_document(note_fts_document(¬e)).await {
4960 record_conditional_insert_degradation(
4961 &mut degradations,
4962 note.id,
4963 ConditionalInsertStage::FtsUpsert,
4964 error,
4965 );
4966 }
4967 }
4968 Err(error) => record_conditional_insert_degradation(
4969 &mut degradations,
4970 note.id,
4971 ConditionalInsertStage::FtsAcquisition,
4972 error,
4973 ),
4974 }
4975
4976 let embed_model_names = self.embedding_models_for_note_kind(kind);
4977 for model_name in &embed_model_names {
4978 match self
4979 .embed_document_with_model_outcome_for_token(
4980 token,
4981 model_name,
4982 note_embedding_text_ref(¬e),
4983 )
4984 .await
4985 {
4986 Ok(outcome) => {
4987 if outcome.truncated {
4988 tracing::warn!(
4989 note_id = %note.id,
4990 model = %outcome.model_name,
4991 source_bytes = outcome.source_bytes,
4992 embedded_bytes = outcome.embedded_bytes,
4993 "try_create_note: embedding input truncated; full content stored unchanged"
4994 );
4995 }
4996 match self.vectors_for_model(token, model_name) {
4997 Ok(vs) => {
4998 if let Err(error) = vs
4999 .insert(
5000 note.id,
5001 SubstrateKind::Note,
5002 ns,
5003 "note.content",
5004 vec![outcome.vector],
5005 )
5006 .await
5007 {
5008 record_conditional_insert_degradation(
5009 &mut degradations,
5010 note.id,
5011 ConditionalInsertStage::VectorInsert,
5012 format!("model {model_name}: {error}"),
5013 );
5014 }
5015 }
5016 Err(error) => record_conditional_insert_degradation(
5017 &mut degradations,
5018 note.id,
5019 ConditionalInsertStage::VectorAcquisition,
5020 format!("model {model_name}: {error}"),
5021 ),
5022 }
5023 }
5024 Err(error) => record_conditional_insert_degradation(
5025 &mut degradations,
5026 note.id,
5027 ConditionalInsertStage::Embedding,
5028 format!("model {model_name}: {error}"),
5029 ),
5030 }
5031 }
5032
5033 legacy_post_commit_result("try_create_note", note.id, Some(note), degradations)
5034 }
5035
5036 #[allow(clippy::too_many_arguments)]
5040 async fn create_note_inner(
5041 &self,
5042 token: &NamespaceToken,
5043 kind: &str,
5044 name: Option<&str>,
5045 content: &str,
5046 embedding_content: Option<&str>,
5047 salience: Option<f64>,
5048 decay_factor: Option<f64>,
5049 properties: Option<serde_json::Value>,
5050 annotates: Vec<Uuid>,
5051 embedding_model: Option<&str>,
5052 capture_visibility: bool,
5053 web_receipt: bool,
5054 ) -> RuntimeResult<(
5055 Note,
5056 crate::retrieval::EmbeddingTruncationReport,
5057 Vec<PostCommitDegradation>,
5058 Vec<(String, u64)>,
5059 )> {
5060 self.validate_note_kind(kind)?;
5061 let mut properties = self.derive_note_write_properties(kind, token, properties)?;
5067 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
5068 if web_receipt {
5069 let map = properties
5070 .as_mut()
5071 .and_then(serde_json::Value::as_object_mut)
5072 .expect("web receipt properties are constructed as an object");
5073 map.insert(
5074 crate::secret_gate::RESERVED_WEB_RECEIPT_KEY.to_string(),
5075 serde_json::Value::String(
5076 crate::secret_gate::WEB_RECEIPT_PROVENANCE_VALUE.to_string(),
5077 ),
5078 );
5079 }
5080 crate::secret_gate::check_at(content, "note", "content")?;
5082 if let Some(n) = name {
5083 crate::secret_gate::check_at(n, "note", "name")?;
5084 }
5085 if let Some(ref p) = properties {
5086 crate::secret_gate::check_json_at(p, "note", "properties")?;
5087 }
5088 if let Some(ec) = embedding_content {
5094 if ec.is_empty() {
5095 return Err(RuntimeError::InvalidInput(
5096 "embedding_content must not be empty".into(),
5097 ));
5098 }
5099 if ec.len() >= content.len() || !content.starts_with(ec) {
5100 return Err(RuntimeError::InvalidInput(
5101 "embedding_content must be a proper prefix of content".into(),
5102 ));
5103 }
5104 crate::secret_gate::check_at(ec, "note", "embedding_content")?;
5105 }
5106 let ns = token.namespace().as_str();
5107
5108 for &target_id in &annotates {
5111 if !self.substrate_exists_by_id(token, target_id).await? {
5112 return Err(RuntimeError::NotFound(format!(
5113 "create_note annotates target {target_id} not found"
5114 )));
5115 }
5116 }
5117
5118 if let Some(s) = salience {
5121 if !s.is_finite() || !(0.0..=1.0).contains(&s) {
5122 return Err(RuntimeError::InvalidInput(format!(
5123 "salience must be a finite value in [0.0, 1.0]; got {s}"
5124 )));
5125 }
5126 }
5127 if let Some(d) = decay_factor {
5128 if !d.is_finite() || d < 0.0 {
5129 return Err(RuntimeError::InvalidInput(format!(
5130 "decay_factor must be a finite value >= 0.0; got {d}"
5131 )));
5132 }
5133 }
5134
5135 if let Some(model_name) = embedding_model {
5139 self.resolve_embedding_model(Some(model_name))?;
5140 }
5141
5142 let mut note = Note::new(ns, kind, content);
5143 if let Some(s) = salience {
5144 note = note.with_salience(s);
5145 }
5146 if let Some(df) = decay_factor {
5147 note = note.with_decay(df);
5148 }
5149 if let Some(n) = name {
5150 note = note.with_name(n);
5151 }
5152 if let Some(p) = properties {
5153 note = note.with_properties(p);
5154 }
5155 let notes = if web_receipt {
5156 self.raw_notes(token)?
5157 } else {
5158 self.notes(token)?
5159 };
5160 notes.upsert_note(note.clone()).await?;
5161
5162 let embed_model_names: Vec<String> = if let Some(m) = embedding_model {
5168 vec![m.to_string()]
5169 } else {
5170 self.embedding_models_for_note_kind(kind)
5171 };
5172
5173 {
5175 #[cfg(any(test, feature = "fault-injection"))]
5181 let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
5182 #[cfg(not(any(test, feature = "fault-injection")))]
5183 let fts_inject = false;
5184 let fts_result: RuntimeResult<()> = if fts_inject {
5185 Err(RuntimeError::Internal("injected FTS failure".to_string()))
5186 } else {
5187 let statements =
5188 khive_db::stores::text::delete_document_statements("fts_notes", ns, note.id)
5189 .into_iter()
5190 .chain(khive_db::stores::text::insert_document_statements(
5191 "fts_notes",
5192 ¬e_fts_document(¬e),
5193 ))
5194 .collect();
5195 self.apply_note_index_revision(¬e, statements)
5196 .await
5197 .map(|_| ())
5198 };
5199
5200 if let Err(e) = fts_result {
5201 self.compensate_note_creation(¬e).await;
5202 return Err(e);
5203 }
5204 }
5205
5206 let canonical_embed_text = note_embedding_text_ref(¬e);
5216 let embed_text = embedding_content.unwrap_or(canonical_embed_text);
5217
5218 let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
5219 let mut vector_fences = Vec::with_capacity(embed_model_names.len());
5220 if embed_model_names.len() == 1 {
5221 let model_name = &embed_model_names[0];
5223 let vec_result = self
5224 .embed_document_with_model_outcome_for_token(token, model_name, embed_text)
5225 .await;
5226
5227 #[cfg(any(test, feature = "fault-injection"))]
5237 let vec_inject = {
5238 let ns_inject = consume_fault(&VECTOR_FAIL_NS, ns);
5239 let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
5240 Some(0) => {
5241 cell.set(None);
5242 true
5243 }
5244 Some(n) => {
5245 cell.set(Some(n - 1));
5246 false
5247 }
5248 None => false,
5249 });
5250 ns_inject || count_inject
5251 };
5252 #[cfg(not(any(test, feature = "fault-injection")))]
5253 let vec_inject = false;
5254 let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
5255 if vec_inject {
5256 Err(RuntimeError::Internal(
5257 "injected vector failure".to_string(),
5258 ))
5259 } else {
5260 vec_result
5261 };
5262
5263 let single_model_result: RuntimeResult<()> = match vec_result {
5264 Ok(outcome) => {
5265 embedding_report.observe(&outcome);
5266 if capture_visibility {
5267 match self
5268 .publish_note_vector_revision_with_seq(
5269 token,
5270 ¬e,
5271 model_name,
5272 &outcome.vector,
5273 )
5274 .await
5275 {
5276 Ok(Some(seq)) => {
5277 vector_fences.push((model_name.clone(), seq));
5278 Ok(())
5279 }
5280 Ok(None) => Ok(()),
5281 Err(error) => Err(error),
5282 }
5283 } else {
5284 self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
5285 .await
5286 .map(|_| ())
5287 }
5288 }
5289 Err(e) => Err(e),
5290 };
5291 if let Err(e) = single_model_result {
5292 self.compensate_note_creation(¬e).await;
5293 return Err(e);
5294 }
5295 } else if !embed_model_names.is_empty() {
5296 let rt_clone = self.clone();
5299 let content_owned: std::sync::Arc<str> = std::sync::Arc::from(embed_text);
5302 let usage_ctx = crate::usage::current();
5303 let mut join_set = tokio::task::JoinSet::new();
5304 for (idx, model_name) in embed_model_names.iter().enumerate() {
5305 let rt = rt_clone.clone();
5306 let text = std::sync::Arc::clone(&content_owned);
5307 let name = model_name.clone();
5308 let ctx = usage_ctx.clone();
5309 let token = (*token).clone();
5310 join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
5311 let fut = rt.embed_document_with_model_outcome_for_token(
5312 &token,
5313 &name,
5314 text.as_ref(),
5315 );
5316 let result = match ctx {
5317 Some(ctx) => crate::usage::scope(ctx, fut).await,
5318 None => fut.await,
5319 };
5320 (idx, result)
5321 }));
5322 }
5323 let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
5327 Ok(outcomes) => outcomes,
5328 Err(e) => {
5329 self.compensate_note_creation(¬e).await;
5330 return Err(e);
5331 }
5332 };
5333 for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
5335 embedding_report.observe(&outcome);
5336 let insert_result = if capture_visibility {
5337 self.publish_note_vector_revision_with_seq(
5338 token,
5339 ¬e,
5340 model_name,
5341 &outcome.vector,
5342 )
5343 .await
5344 .map(|seq| {
5345 if let Some(seq) = seq {
5346 vector_fences.push((model_name.clone(), seq));
5347 }
5348 })
5349 } else {
5350 self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
5351 .await
5352 .map(|_| ())
5353 };
5354 if let Err(e) = insert_result {
5355 self.compensate_note_creation(¬e).await;
5356 return Err(e);
5357 }
5358 }
5359 }
5360
5361 let mut created_edges: Vec<Uuid> = Vec::with_capacity(annotates.len());
5366
5367 #[cfg(test)]
5370 let annotates_iter: Vec<(usize, Uuid)> = annotates
5371 .iter()
5372 .enumerate()
5373 .map(|(i, &id)| (i, id))
5374 .collect();
5375 #[cfg(test)]
5376 macro_rules! next_target {
5377 ($pair:expr) => {
5378 $pair.1
5379 };
5380 }
5381 #[cfg(not(test))]
5382 let annotates_iter: Vec<Uuid> = annotates.to_vec();
5383 #[cfg(not(test))]
5384 macro_rules! next_target {
5385 ($pair:expr) => {
5386 $pair
5387 };
5388 }
5389
5390 for pair in annotates_iter {
5391 let target_id = next_target!(pair);
5392
5393 #[cfg(test)]
5395 let injected_err: Option<RuntimeError> = {
5396 let call_idx = pair.0;
5397 LINK_FAIL_AFTER.with(|cell| {
5398 let n = cell.get();
5399 if n > 0 && call_idx + 1 == n {
5400 cell.set(0); Some(RuntimeError::Internal("injected link failure".to_string()))
5402 } else {
5403 None
5404 }
5405 })
5406 };
5407 #[cfg(not(test))]
5408 let injected_err: Option<RuntimeError> = None;
5409
5410 let link_result = if let Some(e) = injected_err {
5411 Err(e)
5412 } else {
5413 self.link(
5414 token,
5415 note.id,
5416 target_id,
5417 EdgeRelation::Annotates,
5418 1.0,
5419 None,
5420 )
5421 .await
5422 };
5423
5424 match link_result {
5425 Ok(edge) => created_edges.push(edge.id.into()),
5426 Err(e) => {
5427 let edge_ids = created_edges
5432 .iter()
5433 .map(Uuid::to_string)
5434 .collect::<Vec<_>>()
5435 .join(", ");
5436 match self.compensate_note_creation_with_edges(¬e).await {
5437 Ok(true) => return Err(e),
5438 Ok(false) => {
5439 return Err(RuntimeError::Internal(format!(
5440 "create_note: annotates link failed: {e}; note {} changed before \
5441 compensation, retaining its incident edges [{edge_ids}]",
5442 note.id
5443 )));
5444 }
5445 Err(cleanup_error) => {
5446 return Err(RuntimeError::Internal(format!(
5447 "create_note: annotates link failed: {e}; compensation failed \
5448 for note {} and retained edges [{edge_ids}]: {cleanup_error}; \
5449 note and edges remain for reconciliation",
5450 note.id
5451 )));
5452 }
5453 }
5454 }
5455 }
5456 }
5457
5458 let created_event = khive_storage::event::Event::new(
5463 note.namespace.clone(),
5464 "create",
5465 EventKind::NoteCreated,
5466 SubstrateKind::Note,
5467 "",
5468 )
5469 .with_target(note.id)
5470 .with_payload(serde_json::json!({
5471 "id": note.id,
5472 "namespace": note.namespace,
5473 "kind": note.kind,
5474 "salience": note.salience,
5475 }));
5476 let event_result = match self.events(token) {
5477 Ok(store) => store
5478 .append_event(created_event)
5479 .await
5480 .map_err(RuntimeError::from),
5481 Err(error) => Err(error),
5482 };
5483 let mut degradations = Vec::new();
5484 if let Err(error) = event_result {
5485 record_post_commit_degradation(
5486 &mut degradations,
5487 "create_note",
5488 note.id,
5489 "event_append",
5490 error,
5491 );
5492 }
5493
5494 vector_fences.sort_by(|a, b| a.0.cmp(&b.0));
5495 Ok((note, embedding_report, degradations, vector_fences))
5496 }
5497
5498 pub async fn list_notes(
5504 &self,
5505 token: &NamespaceToken,
5506 kind: Option<&str>,
5507 limit: u32,
5508 offset: u32,
5509 ) -> RuntimeResult<Vec<Note>> {
5510 let visible = token.visible_namespaces();
5511 if visible.len() == 1 {
5512 let page = self
5514 .notes(token)?
5515 .query_notes_count_free(
5516 token.namespace().as_str(),
5517 kind,
5518 PageRequest {
5519 offset: offset.into(),
5520 limit,
5521 },
5522 )
5523 .await?;
5524 return Ok(page.items);
5525 }
5526 use khive_storage::note::NoteFilter;
5528 let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
5529 let filter = NoteFilter {
5530 kind: kind.map(|k| k.to_string()),
5531 namespaces: ns_strs,
5532 ..Default::default()
5533 };
5534 let page = self
5535 .notes(token)?
5536 .query_notes_filtered_count_free(
5537 token.namespace().as_str(),
5538 &filter,
5539 PageRequest {
5540 offset: offset.into(),
5541 limit,
5542 },
5543 )
5544 .await?;
5545 Ok(page.items)
5546 }
5547
5548 pub async fn list_notes_after(
5554 &self,
5555 token: &NamespaceToken,
5556 kind: Option<&str>,
5557 after: Option<Uuid>,
5558 limit: u32,
5559 ) -> RuntimeResult<(Vec<Note>, Option<Uuid>)> {
5560 let store = self.notes(token)?;
5561 let after = match after {
5562 Some(id) => {
5563 let note = self
5564 .get_note_including_deleted(token, id)
5565 .await?
5566 .ok_or_else(|| RuntimeError::NotFound(format!("note cursor {id}")))?;
5567 Self::ensure_namespace_visible(¬e.namespace, token)?;
5568 let sequence = store.note_sequence(id).await?.ok_or_else(|| {
5569 RuntimeError::Internal(format!(
5570 "note cursor {id} has no insertion-sequence ledger row"
5571 ))
5572 })?;
5573 Some(SeekCursor { sequence, id })
5574 }
5575 None => None,
5576 };
5577 let filter = khive_storage::note::NoteFilter {
5578 kind: kind.map(str::to_string),
5579 namespaces: token
5580 .visible_namespaces()
5581 .iter()
5582 .map(|namespace| namespace.as_str().to_owned())
5583 .collect(),
5584 ..Default::default()
5585 };
5586 let page = store
5587 .query_notes_filtered_after(token.namespace().as_str(), &filter, after, limit)
5588 .await?;
5589 Ok((page.items, page.next_after.map(|cursor| cursor.id)))
5590 }
5591
5592 pub async fn count_notes(
5594 &self,
5595 token: &NamespaceToken,
5596 kind: Option<&str>,
5597 ) -> RuntimeResult<u64> {
5598 let namespaces: Vec<String> = token
5599 .visible_namespaces()
5600 .iter()
5601 .map(|namespace| namespace.as_str().to_owned())
5602 .collect();
5603 Ok(self
5604 .notes(token)?
5605 .count_notes_in_namespaces(&namespaces, kind)
5606 .await?)
5607 }
5608
5609 #[allow(clippy::too_many_arguments)]
5629 pub async fn search_notes(
5630 &self,
5631 token: &NamespaceToken,
5632 query_text: &str,
5633 query_vector: Option<Vec<f32>>,
5634 limit: u32,
5635 note_kind: Option<&str>,
5636 include_superseded: bool,
5637 tags_any: &[String],
5638 properties_filter: Option<&serde_json::Value>,
5639 ) -> RuntimeResult<Vec<NoteSearchHit>> {
5640 self.search_notes_with_text_mode(
5641 token,
5642 query_text,
5643 query_vector,
5644 limit,
5645 note_kind,
5646 include_superseded,
5647 tags_any,
5648 properties_filter,
5649 TextQueryMode::Plain,
5650 )
5651 .await
5652 }
5653
5654 #[allow(clippy::too_many_arguments)]
5656 pub async fn search_notes_with_text_mode(
5657 &self,
5658 token: &NamespaceToken,
5659 query_text: &str,
5660 query_vector: Option<Vec<f32>>,
5661 limit: u32,
5662 note_kind: Option<&str>,
5663 include_superseded: bool,
5664 tags_any: &[String],
5665 properties_filter: Option<&serde_json::Value>,
5666 text_mode: TextQueryMode,
5667 ) -> RuntimeResult<Vec<NoteSearchHit>> {
5668 let (hits, _vector_error) = self
5669 .search_notes_inner(
5670 token,
5671 query_text,
5672 query_vector,
5673 limit,
5674 note_kind,
5675 include_superseded,
5676 tags_any,
5677 properties_filter,
5678 text_mode,
5679 false,
5680 )
5681 .await?;
5682 Ok(hits)
5683 }
5684
5685 #[allow(clippy::too_many_arguments)]
5692 pub async fn search_notes_outcome(
5693 &self,
5694 token: &NamespaceToken,
5695 query_text: &str,
5696 limit: u32,
5697 note_kind: Option<&str>,
5698 include_superseded: bool,
5699 tags_any: &[String],
5700 properties_filter: Option<&serde_json::Value>,
5701 ) -> RuntimeResult<NoteSearchOutcome> {
5702 self.search_notes_outcome_with_text_mode(
5703 token,
5704 query_text,
5705 limit,
5706 note_kind,
5707 include_superseded,
5708 tags_any,
5709 properties_filter,
5710 TextQueryMode::Plain,
5711 )
5712 .await
5713 }
5714
5715 #[allow(clippy::too_many_arguments)]
5717 pub async fn search_notes_outcome_with_text_mode(
5718 &self,
5719 token: &NamespaceToken,
5720 query_text: &str,
5721 limit: u32,
5722 note_kind: Option<&str>,
5723 include_superseded: bool,
5724 tags_any: &[String],
5725 properties_filter: Option<&serde_json::Value>,
5726 text_mode: TextQueryMode,
5727 ) -> RuntimeResult<NoteSearchOutcome> {
5728 let (hits, vector_error) = self
5729 .search_notes_inner(
5730 token,
5731 query_text,
5732 None,
5733 limit,
5734 note_kind,
5735 include_superseded,
5736 tags_any,
5737 properties_filter,
5738 text_mode,
5739 true,
5740 )
5741 .await?;
5742 Ok(NoteSearchOutcome { hits, vector_error })
5743 }
5744
5745 #[allow(clippy::too_many_arguments)]
5746 async fn search_notes_inner(
5747 &self,
5748 token: &NamespaceToken,
5749 query_text: &str,
5750 query_vector: Option<Vec<f32>>,
5751 limit: u32,
5752 note_kind: Option<&str>,
5753 include_superseded: bool,
5754 tags_any: &[String],
5755 properties_filter: Option<&serde_json::Value>,
5756 text_mode: TextQueryMode,
5757 tolerate_vector_error: bool,
5758 ) -> RuntimeResult<(Vec<NoteSearchHit>, Option<String>)> {
5759 const RRF_K: usize = 60;
5760 let candidates = limit.saturating_mul(4).max(limit);
5761 let visible_ns: Vec<String> = token
5762 .visible_namespaces()
5763 .iter()
5764 .map(|ns| ns.as_str().to_owned())
5765 .collect();
5766
5767 #[cfg(any(test, feature = "fault-injection"))]
5780 let fts_search_inject = {
5781 let mut g = FTS_SEARCH_FAIL_NS.lock().unwrap();
5782 match g.as_deref() {
5783 Some(armed) if visible_ns.iter().any(|ns| ns == armed) => {
5784 *g = None;
5785 true
5786 }
5787 _ => false,
5788 }
5789 };
5790 #[cfg(not(any(test, feature = "fault-injection")))]
5791 let fts_search_inject = false;
5792
5793 let text_store = self.text_for_notes(token)?;
5794 let text_fut = async {
5795 if fts_search_inject {
5796 return Err(khive_storage::StorageError::Timeout {
5797 operation: "fts_search".into(),
5798 });
5799 }
5800 text_store
5801 .search(TextSearchRequest {
5802 query: query_text.to_string(),
5803 mode: text_mode,
5804 filter: Some(TextFilter {
5805 namespaces: visible_ns.clone(),
5806 record_kinds: note_kind
5813 .map(|kind| vec![kind.to_string()])
5814 .unwrap_or_default(),
5815 ..TextFilter::default()
5816 }),
5817 top_k: candidates,
5818 snippet_chars: 200,
5819 })
5820 .await
5821 };
5822 let text_fut = crate::stage_seam::text_stage(text_fut);
5823
5824 let vector_fut = async {
5826 if query_vector.is_some() || self.config().embedding_model.is_some() {
5827 self.note_search_vector_search(token, query_vector, query_text, candidates)
5828 .await
5829 } else {
5830 Ok(vec![])
5831 }
5832 };
5833 let (text_search_result, vector_result) = tokio::join!(text_fut, vector_fut);
5834
5835 let text_hits = crate::error::fts_text_leg_or_err(
5841 text_search_result.map_err(RuntimeError::from),
5842 "search_notes",
5843 query_text,
5844 )?;
5845
5846 let mut vector_error: Option<String> = None;
5847 let vector_hits = match vector_result {
5848 Ok(hits) => hits,
5849 Err(e) if tolerate_vector_error => {
5850 vector_error = Some(e.to_string());
5851 Vec::new()
5852 }
5853 Err(e) => return Err(e),
5854 };
5855
5856 let fuse_k = text_hits.len() + vector_hits.len();
5862 let fused = crate::fusion::rrf_fuse_k(self, text_hits, vector_hits, RRF_K, fuse_k).await?;
5863
5864 let candidate_ids: Vec<Uuid> = fused.iter().map(|hit| hit.entity_id).collect();
5865 if candidate_ids.is_empty() {
5866 return Ok((vec![], vector_error));
5867 }
5868
5869 let note_store = self.notes(token)?;
5876 let search_pool = self.backend().pool_arc();
5877 let mailbox_view = crate::MailboxView {
5878 actor_id: token.actor().id.clone(),
5879 delegated: false,
5880 };
5881 let mut alive_notes: HashMap<Uuid, Note> = HashMap::new();
5882 for note in note_store.get_notes_batch(&candidate_ids).await? {
5883 search_pool.record_note_candidate_hydration_row();
5884 if note.deleted_at.is_some() {
5885 continue;
5886 }
5887 if !mailbox_view.permits_message_note(token, ¬e) {
5888 continue;
5889 }
5890 if let Some(want_kind) = note_kind {
5891 if note.kind != want_kind {
5892 continue;
5893 }
5894 }
5895 if !tags_any.is_empty() {
5900 let note_tags: Vec<String> = note
5901 .properties
5902 .as_ref()
5903 .and_then(|p| p.get("tags"))
5904 .and_then(serde_json::Value::as_array)
5905 .map(|arr| {
5906 arr.iter()
5907 .filter_map(serde_json::Value::as_str)
5908 .map(str::to_owned)
5909 .collect()
5910 })
5911 .unwrap_or_default();
5912 if !note_tags
5913 .iter()
5914 .any(|t| tags_any.iter().any(|w| t.eq_ignore_ascii_case(w)))
5915 {
5916 continue;
5917 }
5918 }
5919 if let Some(pf) = properties_filter {
5921 if !note_props_match(note.properties.as_ref(), pf) {
5922 continue;
5923 }
5924 }
5925 alive_notes.insert(note.id, note);
5926 }
5927
5928 if !include_superseded && !alive_notes.is_empty() {
5931 let graph = self.graph(token)?;
5932 let note_ids: Vec<Uuid> = alive_notes.keys().copied().collect();
5933 let superseded: std::collections::HashSet<Uuid> = graph
5934 .batch_neighbors(
5935 ¬e_ids,
5936 NeighborQuery {
5937 direction: Direction::In,
5938 relations: Some(vec![EdgeRelation::Supersedes]),
5939 limit: Some(1),
5940 min_weight: None,
5941 },
5942 )
5943 .await?
5944 .into_iter()
5945 .map(|(note_id, _)| note_id)
5946 .collect();
5947 alive_notes.retain(|id, _| !superseded.contains(id));
5948 }
5949
5950 let mut hits: Vec<NoteSearchHit> = fused
5952 .into_iter()
5953 .filter_map(|hit| {
5954 let note = alive_notes.get(&hit.entity_id)?;
5955 let weighted = salience_weighted_rank(hit.score, note.salience);
5956 Some(NoteSearchHit {
5957 note_id: hit.entity_id,
5958 score: weighted,
5959 rank_score_kind: hit.rank_score_kind,
5960 signals: hit.signals,
5961 source: hit.source,
5962 title: hit.title.or_else(|| note_title(note)),
5963 snippet: hit.snippet.or_else(|| note_snippet(note)),
5964 })
5965 })
5966 .collect();
5967
5968 hits.sort_by(|a, b| b.score.cmp(&a.score).then(a.note_id.cmp(&b.note_id)));
5969 hits.truncate(limit as usize);
5970 Ok((hits, vector_error))
5971 }
5972
5973 pub async fn resolve_prefix(
5980 &self,
5981 token: &NamespaceToken,
5982 prefix: &str,
5983 ) -> RuntimeResult<Option<Uuid>> {
5984 let namespaces = [token.namespace().as_str().to_owned()];
5985 self.resolve_prefix_inner(Some(&namespaces), prefix, false, false)
5986 .await
5987 }
5988
5989 pub async fn resolve_prefix_including_deleted(
5990 &self,
5991 token: &NamespaceToken,
5992 prefix: &str,
5993 ) -> RuntimeResult<Option<Uuid>> {
5994 let namespaces = [token.namespace().as_str().to_owned()];
5995 self.resolve_prefix_inner(Some(&namespaces), prefix, true, false)
5996 .await
5997 }
5998
5999 pub async fn resolve_prefix_unfiltered(&self, prefix: &str) -> RuntimeResult<Option<Uuid>> {
6007 self.resolve_prefix_inner(None, prefix, false, false).await
6008 }
6009
6010 pub async fn resolve_prefix_unfiltered_including_deleted(
6013 &self,
6014 prefix: &str,
6015 ) -> RuntimeResult<Option<Uuid>> {
6016 self.resolve_prefix_inner(None, prefix, true, false).await
6017 }
6018
6019 pub(crate) async fn resolve_prefix_for_kg_read(
6022 &self,
6023 prefix: &str,
6024 include_deleted: bool,
6025 ) -> RuntimeResult<Option<Uuid>> {
6026 self.resolve_prefix_inner(None, prefix, include_deleted, true)
6027 .await
6028 }
6029
6030 async fn resolve_prefix_inner(
6042 &self,
6043 namespaces: Option<&[String]>,
6044 prefix: &str,
6045 include_deleted: bool,
6046 require_tables: bool,
6047 ) -> RuntimeResult<Option<Uuid>> {
6048 if !prefix.chars().all(|c| c.is_ascii_hexdigit() || c == '-') {
6055 return Ok(None);
6056 }
6057
6058 #[cfg(any(test, feature = "fault-injection"))]
6062 if consume_fault(&PREFIX_RESOLVE_FAIL_NS, prefix) {
6063 return Err(RuntimeError::Storage(
6064 khive_storage::StorageError::Timeout {
6065 operation: "resolve_prefix".into(),
6066 },
6067 ));
6068 }
6069
6070 let Some((lower, upper)) = uuid_prefix_bounds(prefix) else {
6071 return Ok(None);
6072 };
6073
6074 let tables = [
6075 ("entities", true),
6076 ("notes", true),
6077 ("events", false),
6078 ("graph_edges", false),
6079 ];
6080
6081 let mut matches: Vec<String> = Vec::new();
6090 let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
6091 let mut reader = self.sql().reader().await.map_err(RuntimeError::Storage)?;
6092
6093 for (table, has_deleted_at) in tables {
6094 let sql = resolve_prefix_statement(
6095 table,
6096 has_deleted_at,
6097 include_deleted,
6098 namespaces,
6099 &lower,
6100 &upper,
6101 );
6102 match reader.query_all(sql).await {
6103 Ok(rows) => {
6104 for row in rows {
6105 if let Some(col) = row.columns.first() {
6106 if let SqlValue::Text(s) = &col.value {
6107 if seen.insert(s.clone()) {
6108 matches.push(s.clone());
6109 }
6110 }
6111 }
6112 }
6113 }
6114 Err(e) => {
6115 let msg = e.to_string();
6116 if !require_tables && msg.contains("no such table") {
6117 continue;
6118 }
6119 return Err(RuntimeError::Storage(e));
6120 }
6121 }
6122 if matches.len() > 1 {
6123 break;
6124 }
6125 }
6126
6127 if matches.len() <= 1 {
6132 if let Some(sidecar_sql) = self.events_sidecar_sql_read_only()? {
6133 let mut sql = resolve_prefix_statement(
6136 "events",
6137 false,
6138 include_deleted,
6139 namespaces,
6140 &lower,
6141 &upper,
6142 );
6143 sql.label = Some("resolve_prefix.events_sidecar".into());
6144 let mut sidecar_reader =
6145 sidecar_sql.reader().await.map_err(RuntimeError::Storage)?;
6146 match sidecar_reader.query_all(sql).await {
6147 Ok(rows) => {
6148 for row in rows {
6149 if let Some(col) = row.columns.first() {
6150 if let SqlValue::Text(s) = &col.value {
6151 if seen.insert(s.clone()) {
6152 matches.push(s.clone());
6153 }
6154 }
6155 }
6156 }
6157 }
6158 Err(e) => {
6159 let msg = e.to_string();
6160 if require_tables || !msg.contains("no such table") {
6161 return Err(RuntimeError::Storage(e));
6162 }
6163 }
6164 }
6165 }
6166 }
6167
6168 match matches.len() {
6169 0 => Ok(None),
6170 1 => {
6171 let uuid = Uuid::from_str(&matches[0])
6172 .map_err(|e| RuntimeError::Internal(format!("stored UUID is invalid: {e}")))?;
6173 Ok(Some(uuid))
6174 }
6175 _ => {
6176 let uuids: Vec<uuid::Uuid> = matches
6177 .iter()
6178 .filter_map(|s| Uuid::from_str(s).ok())
6179 .collect();
6180 Err(RuntimeError::AmbiguousPrefix {
6181 prefix: prefix.to_string(),
6182 matches: uuids,
6183 })
6184 }
6185 }
6186 }
6187
6188 pub async fn resolve_by_id(
6199 &self,
6200 token: &NamespaceToken,
6201 id: Uuid,
6202 ) -> RuntimeResult<Option<Resolved>> {
6203 if let Some(entity) = self.entities(token)?.get_entity(id).await? {
6205 return Ok(Some(Resolved::Entity(entity)));
6206 }
6207
6208 if let Some(note) = self.notes(token)?.get_note(id).await? {
6210 return Ok(Some(Resolved::Note(note)));
6211 }
6212
6213 Ok(None)
6216 }
6217
6218 pub async fn resolve_by_id_including_deleted(
6224 &self,
6225 token: &NamespaceToken,
6226 id: Uuid,
6227 ) -> RuntimeResult<Option<Resolved>> {
6228 if let Some(entity) = self
6230 .entities(token)?
6231 .get_entity_including_deleted(id)
6232 .await?
6233 {
6234 return Ok(Some(Resolved::Entity(entity)));
6235 }
6236
6237 if let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? {
6239 return Ok(Some(Resolved::Note(note)));
6240 }
6241
6242 Ok(None)
6245 }
6246
6247 pub async fn resolve(
6252 &self,
6253 token: &NamespaceToken,
6254 id: Uuid,
6255 ) -> RuntimeResult<Option<Resolved>> {
6256 match self.get_entity(token, id).await {
6258 Ok(entity) => return Ok(Some(Resolved::Entity(entity))),
6259 Err(RuntimeError::NotFound(_) | RuntimeError::NamespaceMismatch { .. }) => {}
6260 Err(e) => return Err(e),
6261 }
6262
6263 if let Some(note) = self.notes(token)?.get_note(id).await? {
6265 if Self::ensure_namespace_visible(¬e.namespace, token).is_ok() {
6266 return Ok(Some(Resolved::Note(note)));
6267 }
6268 }
6269
6270 if let Some(event) = self.events(token)?.get_event(id).await? {
6272 if Self::ensure_namespace_visible(&event.namespace, token).is_ok() {
6273 return Ok(Some(Resolved::Event(event)));
6274 }
6275 }
6276
6277 Ok(None)
6278 }
6279
6280 pub async fn resolve_edge_endpoint(
6291 &self,
6292 token: &NamespaceToken,
6293 id: Uuid,
6294 ) -> RuntimeResult<Option<Resolved>> {
6295 if let Some(resolved) = self.resolve_by_id(token, id).await? {
6296 return Ok(Some(resolved));
6297 }
6298 if let Some(event) = self.events(token)?.get_event(id).await? {
6299 return Ok(Some(Resolved::Event(event)));
6300 }
6301 Ok(None)
6302 }
6303
6304 pub async fn resolve_primary(
6309 &self,
6310 token: &NamespaceToken,
6311 id: Uuid,
6312 ) -> RuntimeResult<Option<Resolved>> {
6313 let ns = token.namespace().as_str();
6314
6315 if let Some(entity) = self.entities(token)?.get_entity(id).await? {
6317 if Self::ensure_namespace(&entity.namespace, ns).is_ok() {
6318 return Ok(Some(Resolved::Entity(entity)));
6319 }
6320 }
6321
6322 if let Some(note) = self.notes(token)?.get_note(id).await? {
6324 if Self::ensure_namespace(¬e.namespace, ns).is_ok() {
6325 return Ok(Some(Resolved::Note(note)));
6326 }
6327 }
6328
6329 if let Some(event) = self.events(token)?.get_event(id).await? {
6331 if Self::ensure_namespace(&event.namespace, ns).is_ok() {
6332 return Ok(Some(Resolved::Event(event)));
6333 }
6334 }
6335
6336 Ok(None)
6337 }
6338
6339 pub async fn resolve_including_deleted(
6344 &self,
6345 token: &NamespaceToken,
6346 id: Uuid,
6347 ) -> RuntimeResult<Option<Resolved>> {
6348 let ns = token.namespace().as_str();
6349
6350 if let Some(entity) = self
6351 .entities(token)?
6352 .get_entity_including_deleted(id)
6353 .await?
6354 {
6355 if Self::ensure_namespace(&entity.namespace, ns).is_ok() {
6356 return Ok(Some(Resolved::Entity(entity)));
6357 }
6358 }
6359
6360 if let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? {
6361 if Self::ensure_namespace(¬e.namespace, ns).is_ok() {
6362 return Ok(Some(Resolved::Note(note)));
6363 }
6364 }
6365
6366 if let Some(event) = self.events(token)?.get_event(id).await? {
6367 if Self::ensure_namespace(&event.namespace, ns).is_ok() {
6368 return Ok(Some(Resolved::Event(event)));
6369 }
6370 }
6371
6372 Ok(None)
6373 }
6374
6375 async fn atomic_hard_delete_with_edge_purge(
6387 &self,
6388 row_statement: SqlStatement,
6389 node_id: Uuid,
6390 namespace: &str,
6391 actor: &str,
6392 substrate: SubstrateKind,
6393 ) -> RuntimeResult<bool> {
6394 let mut statements = vec![PlanStatement {
6395 statement: row_statement,
6396 guard: Some(AffectedRowGuard::exactly(1)),
6397 }];
6398 if matches!(substrate, SubstrateKind::Entity | SubstrateKind::Note) {
6399 statements.push(PlanStatement {
6400 statement: khive_db::stores::attachment::delete_record_attachments_statement(
6401 node_id,
6402 if substrate == SubstrateKind::Entity {
6403 AttachmentSubstrate::Entity
6404 } else {
6405 AttachmentSubstrate::Note
6406 },
6407 ),
6408 guard: None,
6409 });
6410 }
6411 statements.extend(
6412 hard_delete_lineage_warning_statements(namespace, actor, node_id, substrate)
6413 .into_iter()
6414 .map(|statement| PlanStatement {
6415 statement,
6416 guard: None,
6417 }),
6418 );
6419 statements.push(PlanStatement {
6420 statement: purge_incident_edges_statement(node_id),
6421 guard: None,
6422 });
6423 let plan = AtomicOpPlan::Delete(DeletePlan {
6424 target_id: node_id,
6425 statements,
6426 post_commit: PostCommitEffect::None,
6427 });
6428 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6429 Ok(AtomicRunOutcome::Committed { .. }) => Ok(true),
6430 Ok(AtomicRunOutcome::RolledBack {
6431 failure: AtomicOpFailure::NoteConflict(conflict),
6432 ..
6433 }) => Err(conflict.into_error().into()),
6434 Ok(AtomicRunOutcome::RolledBack {
6435 failure: AtomicOpFailure::EntityConflict(conflict),
6436 ..
6437 }) => Err(conflict.into_error().into()),
6438 Ok(AtomicRunOutcome::RolledBack {
6439 failure: AtomicOpFailure::GuardFailed { .. },
6440 ..
6441 }) => Ok(false),
6442 Ok(AtomicRunOutcome::RolledBack {
6443 failure: AtomicOpFailure::SqlError { message, .. },
6444 ..
6445 }) => Err(RuntimeError::Internal(format!(
6446 "hard delete + edge purge for {node_id} failed: {message}"
6447 ))),
6448 Err(e) => Err(RuntimeError::Internal(format!(
6449 "hard delete + edge purge for {node_id}: atomic unit seam failure: {}",
6450 e.0
6451 ))),
6452 }
6453 }
6454
6455 pub async fn restore_entity(
6461 &self,
6462 token: &NamespaceToken,
6463 id: Uuid,
6464 ) -> RuntimeResult<Option<(Entity, bool)>> {
6465 let Some(entity) = self
6466 .entities(token)?
6467 .get_entity_including_deleted(id)
6468 .await?
6469 else {
6470 return Ok(None);
6471 };
6472 if entity.namespace != token.namespace().as_str() {
6473 return Ok(None);
6474 }
6475 if let Some(kept_id) = entity.merged_into {
6487 if entity.deleted_at.is_none() {
6488 return Err(live_merged_entity_refused(id, kept_id));
6489 }
6490 return Err(merge_tombstone_restore_refused(id, kept_id));
6491 }
6492 if entity.deleted_at.is_none() {
6493 return Ok(Some((entity, false)));
6494 }
6495 let updated_at =
6496 Utc::now()
6497 .timestamp_micros()
6498 .max(entity.updated_at.checked_add(1).ok_or_else(|| {
6499 RuntimeError::Internal(format!(
6500 "entity {id} updated_at is already at i64::MAX and cannot advance"
6501 ))
6502 })?);
6503 let mut restored = entity;
6504 restored.deleted_at = None;
6505 restored.updated_at = updated_at;
6506 restored.version = restored
6507 .version
6508 .checked_add(1)
6509 .ok_or_else(|| RuntimeError::InvalidInput("entity version overflow".into()))?;
6510 let mut statements = vec![PlanStatement {
6511 statement: SqlStatement {
6512 sql: "UPDATE entities SET deleted_at=NULL, updated_at=?1, version=version+1 \
6513 WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL AND version=?4"
6514 .into(),
6515 params: vec![
6516 SqlValue::Integer(updated_at),
6517 SqlValue::Text(id.to_string()),
6518 SqlValue::Text(token.namespace().as_str().to_owned()),
6519 SqlValue::Integer(restored.version - 1),
6520 ],
6521 label: Some("entity-restore".into()),
6522 },
6523 guard: Some(AffectedRowGuard::exactly(1)),
6524 }];
6525 for statement in khive_db::stores::text::delete_document_statements(
6530 "fts_entities",
6531 &restored.namespace,
6532 id,
6533 )
6534 .into_iter()
6535 .chain(insert_document_statements(
6536 "fts_entities",
6537 &entity_fts_document(&restored),
6538 )) {
6539 statements.push(PlanStatement {
6540 statement,
6541 guard: None,
6542 });
6543 }
6544 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6545 graph_effects: Vec::new(),
6546 target_id: id,
6547 statements,
6548 post_commit: PostCommitEffect::None,
6549 edge_natural_key: None,
6550 idempotent_noop: false,
6551 entity_guard: None,
6552 note_guard: None,
6553 note_vector_purge: None,
6554 note_embedding_inheritance: None,
6555 }));
6556 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6557 Ok(AtomicRunOutcome::Committed { .. }) => {
6558 #[cfg(any(test, feature = "fault-injection"))]
6561 if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
6562 return Err(restore_reindex_failed(
6563 "entity",
6564 id,
6565 RuntimeError::Internal("injected FTS failure".to_string()),
6566 ));
6567 }
6568 self.reindex_entity(token, &restored)
6569 .await
6570 .map_err(|e| restore_reindex_failed("entity", id, e))?;
6571 Ok(Some((restored, true)))
6572 }
6573 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6574 format!("entity restore rolled back: {failure:?}"),
6575 )),
6576 Err(error) => Err(RuntimeError::Storage(error.0)),
6577 }
6578 }
6579
6580 pub async fn restore_note(
6586 &self,
6587 token: &NamespaceToken,
6588 id: Uuid,
6589 ) -> RuntimeResult<Option<(Note, bool)>> {
6590 let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? else {
6591 return Ok(None);
6592 };
6593 if note.namespace != token.namespace().as_str() {
6594 return Ok(None);
6595 }
6596 if note.deleted_at.is_none() {
6597 return Ok(Some((note, false)));
6598 }
6599 if let Some(key) = note.key.as_deref() {
6600 if let Some(holder) = self
6601 .notes(token)?
6602 .get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
6603 .await?
6604 .into_iter()
6605 .find(|holder| holder.id != note.id)
6606 {
6607 return Err(restore_key_conflict(key, &holder));
6608 }
6609 }
6610 let updated_at =
6611 Utc::now()
6612 .timestamp_micros()
6613 .max(note.updated_at.checked_add(1).ok_or_else(|| {
6614 RuntimeError::Internal(format!(
6615 "note {id} updated_at is already at i64::MAX and cannot advance"
6616 ))
6617 })?);
6618 let mut params = vec![
6619 SqlValue::Text("active".into()),
6620 SqlValue::Integer(updated_at),
6621 SqlValue::Text(id.to_string()),
6622 SqlValue::Text(note.namespace.clone()),
6623 SqlValue::Text(note.kind.clone()),
6624 ];
6625 let key_clause = if let Some(key) = note.key.as_deref() {
6626 params.push(SqlValue::Text(key.to_owned()));
6627 format!(
6628 " AND (key IS NULL OR NOT EXISTS (SELECT 1 FROM notes live \
6629 WHERE live.namespace=?4 AND live.kind=?5 AND live.key=?{} \
6630 AND live.deleted_at IS NULL AND live.id != notes.id))",
6631 params.len()
6632 )
6633 } else {
6634 String::new()
6635 };
6636 let mut restored = note.clone();
6637 restored.status = "active".into();
6638 restored.deleted_at = None;
6639 restored.updated_at = updated_at;
6640 restored.version = restored
6641 .version
6642 .checked_add(1)
6643 .ok_or_else(|| RuntimeError::Internal(format!("note {id} version is exhausted")))?;
6644 let mut statements = vec![PlanStatement {
6645 statement: SqlStatement {
6646 sql: format!(
6647 "UPDATE notes SET status=?1, deleted_at=NULL, updated_at=?2 \
6648 WHERE id=?3 AND namespace=?4 AND kind=?5 AND deleted_at IS NOT NULL{key_clause}"
6649 ),
6650 params,
6651 label: Some("note-restore".into()),
6652 },
6653 guard: Some(AffectedRowGuard::exactly(1)),
6654 }];
6655 for statement in
6657 khive_db::stores::text::delete_document_statements("fts_notes", &restored.namespace, id)
6658 .into_iter()
6659 .chain(insert_document_statements(
6660 "fts_notes",
6661 ¬e_fts_document(&restored),
6662 ))
6663 {
6664 statements.push(PlanStatement {
6665 statement,
6666 guard: None,
6667 });
6668 }
6669 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6670 graph_effects: Vec::new(),
6671 target_id: id,
6672 statements,
6673 post_commit: PostCommitEffect::None,
6674 edge_natural_key: None,
6675 idempotent_noop: false,
6676 entity_guard: None,
6677 note_guard: None,
6678 note_vector_purge: None,
6679 note_embedding_inheritance: None,
6680 }));
6681 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6682 Ok(AtomicRunOutcome::Committed { .. }) => {
6683 #[cfg(any(test, feature = "fault-injection"))]
6684 if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
6685 return Err(restore_reindex_failed(
6686 "note",
6687 id,
6688 RuntimeError::Internal("injected FTS failure".to_string()),
6689 ));
6690 }
6691 let reindexed = self.reindex_note_with_report(token, &restored).await;
6692 let report = reindexed.map_err(|e| restore_reindex_failed("note", id, e))?;
6693 let degradations = report.post_commit_degradations();
6694 legacy_post_commit_result("restore_note", id, Some((restored, true)), degradations)
6695 }
6696 Ok(AtomicRunOutcome::RolledBack {
6697 failure: AtomicOpFailure::GuardFailed { .. },
6698 ..
6699 }) => {
6700 if let Some(key) = note.key.as_deref() {
6701 if let Some(holder) = self
6702 .notes(token)?
6703 .get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
6704 .await?
6705 .into_iter()
6706 .find(|holder| holder.id != note.id)
6707 {
6708 return Err(restore_key_conflict(key, &holder));
6709 }
6710 }
6711 Err(RuntimeError::NotFound(format!(
6712 "note {id} is no longer a caller-owned tombstone"
6713 )))
6714 }
6715 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6716 format!("note restore rolled back: {failure:?}"),
6717 )),
6718 Err(error) => Err(RuntimeError::Storage(error.0)),
6719 }
6720 }
6721
6722 pub async fn restore_edge(
6724 &self,
6725 token: &NamespaceToken,
6726 id: Uuid,
6727 ) -> RuntimeResult<Option<(Edge, bool)>> {
6728 let Some(edge) = self.get_edge_including_deleted(token, id).await? else {
6729 return Ok(None);
6730 };
6731 if edge.namespace != token.namespace().as_str() {
6732 return Ok(None);
6733 }
6734 if edge.deleted_at.is_none() {
6735 return Ok(Some((edge, false)));
6736 }
6737 let updated_at = Utc::now();
6738 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6739 graph_effects: Vec::new(),
6740 target_id: id,
6741 statements: vec![PlanStatement {
6742 statement: SqlStatement {
6743 sql: "UPDATE graph_edges SET deleted_at=NULL, updated_at=?1 \
6744 WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL"
6745 .into(),
6746 params: vec![
6747 SqlValue::Integer(updated_at.timestamp_micros()),
6748 SqlValue::Text(id.to_string()),
6749 SqlValue::Text(edge.namespace.clone()),
6750 ],
6751 label: Some("edge-restore".into()),
6752 },
6753 guard: Some(AffectedRowGuard::exactly(1)),
6754 }],
6755 post_commit: PostCommitEffect::None,
6756 edge_natural_key: None,
6757 idempotent_noop: false,
6758 entity_guard: None,
6759 note_guard: None,
6760 note_vector_purge: None,
6761 note_embedding_inheritance: None,
6762 }));
6763 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6764 Ok(AtomicRunOutcome::Committed { .. }) => {
6765 let mut restored = edge;
6766 restored.deleted_at = None;
6767 restored.updated_at = updated_at;
6768 Ok(Some((restored, true)))
6769 }
6770 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6771 format!("edge restore rolled back: {failure:?}"),
6772 )),
6773 Err(error) => Err(RuntimeError::Storage(error.0)),
6774 }
6775 }
6776
6777 pub async fn delete_note(
6788 &self,
6789 token: &NamespaceToken,
6790 id: Uuid,
6791 hard: bool,
6792 ) -> RuntimeResult<bool> {
6793 let (deleted, degradations) = self
6794 .delete_note_with_post_commit_report(token, id, hard)
6795 .await?;
6796 legacy_post_commit_result("delete_note", id, deleted, degradations)
6797 }
6798
6799 pub async fn delete_note_with_post_commit_report(
6803 &self,
6804 token: &NamespaceToken,
6805 id: Uuid,
6806 hard: bool,
6807 ) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
6808 let note_store = self.notes(token)?;
6809 let note = if hard {
6810 match note_store.get_note_including_deleted(id).await? {
6811 Some(n) => n,
6812 None => return Ok((false, Vec::new())),
6813 }
6814 } else {
6815 match note_store.get_note(id).await? {
6816 Some(n) => n,
6817 None => return Ok((false, Vec::new())),
6818 }
6819 };
6820 if let Some(error) = self.stream_member_error(¬e).await? {
6821 return Err(error);
6822 }
6823 let mode = if hard {
6824 DeleteMode::Hard
6825 } else {
6826 DeleteMode::Soft
6827 };
6828
6829 let record_tok = token.with_namespace(
6831 khive_types::Namespace::parse(¬e.namespace)
6832 .map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
6833 );
6834 let record_ns = note.namespace.clone();
6835 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
6836
6837 let deleted = if hard {
6842 self.atomic_hard_delete_with_edge_purge(
6843 note_hard_delete_statement(id),
6844 id,
6845 &record_ns,
6846 &actor,
6847 SubstrateKind::Note,
6848 )
6849 .await?
6850 } else {
6851 note_store.delete_note(id, mode).await?
6852 };
6853 let mut degradations = Vec::new();
6854 if deleted {
6855 let fts_result = match self.text_for_notes(&record_tok) {
6856 Ok(store) => store
6857 .delete_document(&record_ns, id)
6858 .await
6859 .map_err(RuntimeError::from),
6860 Err(error) => Err(error),
6861 };
6862 if let Err(error) = fts_result {
6863 record_post_commit_degradation(
6864 &mut degradations,
6865 "delete_note",
6866 id,
6867 "fts_cleanup",
6868 error,
6869 );
6870 }
6871 for model_name in self.registered_embedding_model_names() {
6873 let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
6874 Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
6875 Err(error) => Err(error),
6876 };
6877 if let Err(error) = vector_result {
6878 record_post_commit_degradation(
6879 &mut degradations,
6880 "delete_note",
6881 id,
6882 "vector_cleanup",
6883 format!("{model_name}: {error}"),
6884 );
6885 }
6886 }
6887 let event = khive_storage::event::Event::new(
6888 record_ns.clone(),
6889 "delete",
6890 EventKind::NoteDeleted,
6891 SubstrateKind::Note,
6892 "",
6893 )
6894 .with_target(id)
6895 .with_payload(serde_json::json!({"id": id, "namespace": record_ns, "hard": hard}));
6896 let event_result = match self.events(&record_tok) {
6897 Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
6898 Err(error) => Err(error),
6899 };
6900 if let Err(error) = event_result {
6901 record_post_commit_degradation(
6902 &mut degradations,
6903 "delete_note",
6904 id,
6905 "event_append",
6906 error,
6907 );
6908 }
6909 self.fire_note_mutation_hook(¬e.kind, id).await;
6916 }
6917 Ok((deleted, degradations))
6918 }
6919
6920 pub async fn delete_note_row_first_for_compensation(
6939 &self,
6940 token: &NamespaceToken,
6941 id: Uuid,
6942 ) -> RuntimeResult<()> {
6943 let note_store = self.notes(token)?;
6944 let Some(note) = note_store.get_note_including_deleted(id).await? else {
6945 return Ok(());
6946 };
6947 let record_tok = NamespaceToken::for_namespace(
6948 khive_types::Namespace::parse(¬e.namespace)
6949 .map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
6950 );
6951 let record_ns = note.namespace.clone();
6952
6953 note_store.delete_note(id, DeleteMode::Hard).await?;
6955
6956 #[cfg(any(test, feature = "fault-injection"))]
6957 {
6958 let armed = ROLLBACK_CLEANUP_FAIL_NS.lock().unwrap().take();
6959 if armed.as_deref() == Some(record_ns.as_str()) {
6960 return Err(RuntimeError::Internal(
6961 "row removed but compensation cleanup failed: injected=true".to_string(),
6962 ));
6963 }
6964 }
6965
6966 let mut cleanup_errors = Vec::new();
6967 if let Err(e) = self.graph(&record_tok)?.purge_incident_edges(id).await {
6968 cleanup_errors.push(format!("graph={e}"));
6969 }
6970 if let Err(e) = self
6971 .text_for_notes(&record_tok)?
6972 .delete_document(&record_ns, id)
6973 .await
6974 {
6975 cleanup_errors.push(format!("fts={e}"));
6976 }
6977 for model_name in self.registered_embedding_model_names() {
6978 if let Err(e) = self
6979 .vectors_for_model(&record_tok, &model_name)?
6980 .delete(id)
6981 .await
6982 {
6983 cleanup_errors.push(format!("vector[{model_name}]={e}"));
6984 }
6985 }
6986 if cleanup_errors.is_empty() {
6987 Ok(())
6988 } else {
6989 Err(RuntimeError::Internal(format!(
6990 "row removed but compensation cleanup failed: {}",
6991 cleanup_errors.join("; ")
6992 )))
6993 }
6994 }
6995}
6996
6997#[derive(Clone, Debug, Serialize)]
6999pub struct QueryResult {
7000 pub rows: Vec<SqlRow>,
7001 #[serde(skip_serializing_if = "Vec::is_empty")]
7002 pub warnings: Vec<String>,
7003 pub offset: usize,
7005 pub page_size: usize,
7007 pub has_more: bool,
7009 #[serde(skip_serializing_if = "Option::is_none")]
7011 pub next_offset: Option<usize>,
7012 pub truncated: bool,
7014}
7015
7016#[derive(Debug)]
7018enum SymmetricEdgeUpdateOutcome {
7019 Absorbed(String),
7023 Updated,
7025 Stale,
7030}
7031
7032impl KhiveRuntime {
7033 pub async fn query(&self, token: &NamespaceToken, query: &str) -> RuntimeResult<Vec<SqlRow>> {
7041 Ok(self
7042 .query_with_metadata(token, query, khive_query::CompileOptions::default())
7043 .await?
7044 .rows)
7045 }
7046
7047 pub async fn query_with_metadata(
7049 &self,
7050 token: &NamespaceToken,
7051 query: &str,
7052 mut opts: khive_query::CompileOptions,
7053 ) -> RuntimeResult<QueryResult> {
7054 use khive_query::QueryValue;
7055 use khive_storage::types::SqlValue;
7056
7057 let (language, ast) = khive_query::language::parse_auto_with_language(query)?;
7058 if opts.max_limit == 0 {
7059 return Err(RuntimeError::InvalidInput(
7060 "query page size must be at least 1".into(),
7061 ));
7062 }
7063 let offset = ast.offset;
7064 let page_size = ast.limit.unwrap_or(opts.max_limit).min(opts.max_limit);
7065 opts.scopes = token
7066 .visible_namespaces()
7067 .iter()
7068 .map(|ns| ns.as_str().to_string())
7069 .collect();
7070 let compiled = khive_query::compile(&ast, &opts)?;
7071 let mut warnings = compiled.warnings;
7072 let truncation_check = compiled.truncation_check;
7073
7074 warnings.extend(self.with_pack_edge_rules(|pack_rules| {
7075 static_impossible_edge_pattern_warnings(language, &ast.pattern, pack_rules)
7076 }));
7077
7078 let params: Vec<SqlValue> = compiled
7081 .params
7082 .into_iter()
7083 .map(|qv| match qv {
7084 QueryValue::Null => SqlValue::Null,
7085 QueryValue::Integer(n) => SqlValue::Integer(n),
7086 QueryValue::Float(f) => SqlValue::Float(f),
7087 QueryValue::Text(s) => SqlValue::Text(s),
7088 QueryValue::Blob(b) => SqlValue::Blob(b),
7089 })
7090 .collect();
7091
7092 let mut reader = self.sql().reader().await?;
7093 let stmt = SqlStatement {
7094 sql: compiled.sql,
7095 params,
7096 label: None,
7097 };
7098 let mut rows = reader.query_all(stmt).await?;
7099
7100 let mut truncated = false;
7104 if let Some(check) = truncation_check {
7105 if rows.len() > check.max_limit {
7106 rows.truncate(check.max_limit);
7107 truncated = true;
7108 }
7109 }
7110
7111 let next_offset = if truncated && language == khive_query::QueryLanguage::Gql {
7112 let next = offset.checked_add(rows.len()).ok_or_else(|| {
7113 RuntimeError::InvalidInput("GQL next_offset exceeds usize::MAX".into())
7114 })?;
7115 if next == offset {
7116 return Err(RuntimeError::InvalidInput(
7117 "query page did not advance; page size must be at least 1".into(),
7118 ));
7119 }
7120 i64::try_from(next).map_err(|_| {
7121 RuntimeError::InvalidInput("GQL next_offset exceeds i64::MAX".into())
7122 })?;
7123 Some(next)
7124 } else {
7125 None
7126 };
7127
7128 if truncated {
7129 let Some(check) = truncation_check else {
7130 return Err(RuntimeError::Internal(
7131 "truncated query result is missing sentinel metadata".into(),
7132 ));
7133 };
7134 let bound = match check.requested_limit {
7135 Some(requested) => {
7136 format!("requested query LIMIT {requested} exceeds the effective page size")
7137 }
7138 None => "the query has no explicit LIMIT".to_string(),
7139 };
7140 let warning = match language {
7141 khive_query::QueryLanguage::Gql => {
7142 let Some(next) = next_offset else {
7143 return Err(RuntimeError::Internal(
7144 "truncated GQL result is missing its continuation offset".into(),
7145 ));
7146 };
7147 format!(
7148 "result page capped at {} rows because {bound}; more matches exist. \
7149 Continue the same GQL query with `SKIP {next}` (the machine-readable \
7150 `next_offset`) and keep the same page size.",
7151 check.max_limit
7152 )
7153 }
7154 khive_query::QueryLanguage::Sparql => format!(
7155 "result page capped at {} rows because {bound}; more matches exist. \
7156 SPARQL OFFSET paging is not part of the supported dialect.",
7157 check.max_limit
7158 ),
7159 };
7160 warnings.push(warning);
7161 }
7162
7163 Ok(QueryResult {
7164 rows,
7165 warnings,
7166 offset,
7167 page_size,
7168 has_more: truncated,
7169 next_offset,
7170 truncated,
7171 })
7172 }
7173
7174 pub async fn delete_entity(
7183 &self,
7184 token: &NamespaceToken,
7185 id: Uuid,
7186 hard: bool,
7187 ) -> RuntimeResult<bool> {
7188 let (deleted, degradations) = self
7189 .delete_entity_with_post_commit_report(token, id, hard)
7190 .await?;
7191 legacy_post_commit_result("delete_entity", id, deleted, degradations)
7192 }
7193
7194 pub async fn delete_entity_with_post_commit_report(
7196 &self,
7197 token: &NamespaceToken,
7198 id: Uuid,
7199 hard: bool,
7200 ) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
7201 let entity = if hard {
7202 match self
7203 .entities(token)?
7204 .get_entity_including_deleted(id)
7205 .await?
7206 {
7207 Some(e) => e,
7208 None => return Ok((false, Vec::new())),
7209 }
7210 } else {
7211 match self.entities(token)?.get_entity(id).await? {
7212 Some(e) => e,
7213 None => return Ok((false, Vec::new())),
7214 }
7215 };
7216 let mode = if hard {
7217 DeleteMode::Hard
7218 } else {
7219 DeleteMode::Soft
7220 };
7221
7222 let record_tok = token.with_namespace(
7224 khive_types::Namespace::parse(&entity.namespace)
7225 .map_err(|e| RuntimeError::Internal(format!("entity namespace invalid: {e}")))?,
7226 );
7227 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
7228
7229 let deleted = if hard {
7234 self.atomic_hard_delete_with_edge_purge(
7236 entity_hard_delete_statement(id),
7237 id,
7238 &entity.namespace,
7239 &actor,
7240 SubstrateKind::Entity,
7241 )
7242 .await?
7243 } else {
7244 self.entities(token)?.delete_entity(id, mode).await?
7245 };
7246 let mut degradations = Vec::new();
7247 if deleted {
7248 let ns = entity.namespace.clone();
7249 let fts_result = match self.text(&record_tok) {
7250 Ok(store) => store
7251 .delete_document(&ns, id)
7252 .await
7253 .map_err(RuntimeError::from),
7254 Err(error) => Err(error),
7255 };
7256 if let Err(error) = fts_result {
7257 record_post_commit_degradation(
7258 &mut degradations,
7259 "delete_entity",
7260 id,
7261 "fts_cleanup",
7262 error,
7263 );
7264 }
7265 for model_name in self.registered_embedding_model_names() {
7266 let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
7267 Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
7268 Err(error) => Err(error),
7269 };
7270 if let Err(error) = vector_result {
7271 record_post_commit_degradation(
7272 &mut degradations,
7273 "delete_entity",
7274 id,
7275 "vector_cleanup",
7276 format!("{model_name}: {error}"),
7277 );
7278 }
7279 }
7280 let event = khive_storage::event::Event::new(
7281 ns.clone(),
7282 "delete",
7283 EventKind::EntityDeleted,
7284 SubstrateKind::Entity,
7285 "",
7286 )
7287 .with_target(id)
7288 .with_payload(serde_json::json!({"id": id, "namespace": ns, "hard": hard}));
7289 let event_result = match self.events(&record_tok) {
7290 Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
7291 Err(error) => Err(error),
7292 };
7293 if let Err(error) = event_result {
7294 record_post_commit_degradation(
7295 &mut degradations,
7296 "delete_entity",
7297 id,
7298 "event_append",
7299 error,
7300 );
7301 }
7302 }
7303 Ok((deleted, degradations))
7304 }
7305
7306 pub(crate) async fn delete_entity_attachments_on_core(&self, id: Uuid) -> RuntimeResult<bool> {
7307 let core = self.core();
7308 drop(core.attachments()?);
7309 let statement = khive_db::stores::attachment::delete_record_attachments_statement(
7310 id,
7311 AttachmentSubstrate::Entity,
7312 );
7313 Ok(core.sql().writer().await?.execute(statement).await? > 0)
7314 }
7315
7316 pub async fn count_entities(
7318 &self,
7319 token: &NamespaceToken,
7320 kind: Option<&str>,
7321 ) -> RuntimeResult<u64> {
7322 let ns_strs: Vec<String> = token
7323 .visible_namespaces()
7324 .iter()
7325 .map(|ns| ns.as_str().to_owned())
7326 .collect();
7327 let filter = EntityFilter {
7328 kinds: match kind {
7329 Some(k) => vec![k.to_string()],
7330 None => vec![],
7331 },
7332 namespaces: ns_strs,
7333 ..Default::default()
7334 };
7335 Ok(self
7336 .entities(token)?
7337 .count_entities(token.namespace().as_str(), filter)
7338 .await?)
7339 }
7340
7341 pub async fn entity_stats_counts(
7343 &self,
7344 token: &NamespaceToken,
7345 ) -> RuntimeResult<EntityStatsCounts> {
7346 entity_stats_counts(self.entities(token)?.as_ref(), token).await
7347 }
7348
7349 pub async fn get_edge(
7356 &self,
7357 _token: &NamespaceToken,
7358 edge_id: Uuid,
7359 ) -> RuntimeResult<Option<Edge>> {
7360 let mut reader = self.sql().reader().await?;
7361 let record_ns = reader
7362 .query_scalar(SqlStatement {
7363 sql: "SELECT namespace FROM graph_edges \
7364 WHERE id = ?1 AND deleted_at IS NULL LIMIT 1"
7365 .into(),
7366 params: vec![SqlValue::Text(edge_id.to_string())],
7367 label: Some("get_edge_namespace".into()),
7368 })
7369 .await?;
7370
7371 let Some(SqlValue::Text(record_ns)) = record_ns else {
7372 return Ok(None);
7373 };
7374 let record_tok = NamespaceToken::for_namespace(
7377 khive_types::Namespace::parse(&record_ns)
7378 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
7379 );
7380 Ok(self
7381 .graph(&record_tok)?
7382 .get_edge(LinkId::from(edge_id))
7383 .await?)
7384 }
7385
7386 pub async fn get_edges_by_id(
7395 &self,
7396 _token: &NamespaceToken,
7397 ids: &[Uuid],
7398 ) -> RuntimeResult<Vec<Option<Edge>>> {
7399 let mut edges = Vec::with_capacity(ids.len());
7400 for chunk in ids.chunks(900) {
7401 let window = self.prepare_edge_read_window(chunk).await?;
7402 edges.extend(
7403 Self::hydrate_edge_read_window(chunk, window, |record_token| {
7404 self.graph(record_token)
7405 })
7406 .await?,
7407 );
7408 }
7409 Ok(edges)
7410 }
7411
7412 async fn prepare_edge_read_window(&self, ids: &[Uuid]) -> RuntimeResult<EdgeReadWindow> {
7413 let placeholders = (1..=ids.len())
7414 .map(|index| format!("?{index}"))
7415 .collect::<Vec<_>>()
7416 .join(",");
7417 let mut reader = self.sql().reader().await?;
7418 let rows = reader
7419 .query_all(SqlStatement {
7420 sql: format!(
7421 "SELECT id, namespace FROM graph_edges WHERE id IN ({placeholders}) AND deleted_at IS NULL"
7422 ),
7423 params: ids.iter().map(|id| SqlValue::Text(id.to_string())).collect(),
7424 label: Some("get_edge_namespace".into()),
7425 })
7426 .await?;
7427 let mut namespaces = HashMap::with_capacity(rows.len());
7428 for row in rows {
7429 let Some(SqlValue::Text(id)) = row.columns.first().map(|column| &column.value) else {
7430 return Err(RuntimeError::Internal(
7431 "edge namespace lookup returned an invalid id".into(),
7432 ));
7433 };
7434 let id = Uuid::parse_str(id).map_err(|e| {
7435 RuntimeError::Internal(format!("edge namespace lookup returned an invalid id: {e}"))
7436 })?;
7437 let value = row
7438 .columns
7439 .get(1)
7440 .map(|column| column.value.clone())
7441 .unwrap_or(SqlValue::Null);
7442 namespaces.insert(id, value);
7443 }
7444 let mut window = EdgeReadWindow {
7445 outcomes: (0..ids.len()).map(|_| Some(Ok(None))).collect(),
7446 groups: Vec::new(),
7447 };
7448 let mut group_indices = HashMap::new();
7449 for (index, id) in ids.iter().enumerate() {
7450 let Some(SqlValue::Text(record_ns)) = namespaces.get(id) else {
7451 continue;
7452 };
7453 match khive_types::Namespace::parse(record_ns) {
7454 Ok(namespace) => {
7455 let next_group = window.groups.len();
7456 let group = *group_indices.entry(record_ns.clone()).or_insert(next_group);
7457 if group == next_group {
7458 window.groups.push((namespace, Vec::new()));
7459 }
7460 window.groups[group].1.push(index);
7461 window.outcomes[index] = None;
7462 }
7463 Err(error) => {
7464 window.outcomes[index] = Some(Err(RuntimeError::Internal(format!(
7465 "edge namespace invalid: {error}"
7466 ))));
7467 }
7468 }
7469 }
7470 Ok(window)
7471 }
7472
7473 async fn hydrate_edge_read_window<F>(
7474 ids: &[Uuid],
7475 mut window: EdgeReadWindow,
7476 mut graph: F,
7477 ) -> RuntimeResult<Vec<Option<Edge>>>
7478 where
7479 F: FnMut(&NamespaceToken) -> RuntimeResult<std::sync::Arc<dyn khive_storage::GraphStore>>,
7480 {
7481 for (namespace, indices) in window.groups {
7482 let record_token = NamespaceToken::for_namespace(namespace);
7483 let group_ids: Vec<LinkId> = indices
7484 .iter()
7485 .map(|&index| LinkId::from(ids[index]))
7486 .collect();
7487 let outcomes = match graph(&record_token) {
7488 Ok(store) => store
7489 .get_edge_read_outcomes(&group_ids)
7490 .await
7491 .map_err(RuntimeError::from),
7492 Err(error) => Err(error),
7493 };
7494 match outcomes {
7495 Ok(outcomes) if outcomes.len() == indices.len() => {
7496 for (index, outcome) in indices.into_iter().zip(outcomes) {
7497 window.outcomes[index] = Some(outcome.map_err(RuntimeError::from));
7498 }
7499 }
7500 Ok(_) => {
7501 window.outcomes[indices[0]] = Some(Err(RuntimeError::Internal(
7502 "edge batch returned an invalid outcome count".into(),
7503 )));
7504 }
7505 Err(error) => {
7506 window.outcomes[indices[0]] = Some(Err(error));
7507 }
7508 }
7509 }
7510 window
7511 .outcomes
7512 .into_iter()
7513 .map(|outcome| {
7514 outcome.unwrap_or_else(|| {
7515 Err(RuntimeError::Internal(
7516 "edge batch omitted an input outcome".into(),
7517 ))
7518 })
7519 })
7520 .collect()
7521 }
7522
7523 pub async fn get_edge_visible(
7528 &self,
7529 token: &NamespaceToken,
7530 edge_id: Uuid,
7531 ) -> RuntimeResult<Option<Edge>> {
7532 self.get_edge(token, edge_id).await
7533 }
7534
7535 pub async fn get_edge_including_deleted(
7541 &self,
7542 _token: &NamespaceToken,
7543 edge_id: Uuid,
7544 ) -> RuntimeResult<Option<Edge>> {
7545 let mut reader = self.sql().reader().await?;
7546 let record_ns = reader
7547 .query_scalar(SqlStatement {
7548 sql: "SELECT namespace FROM graph_edges WHERE id = ?1 LIMIT 1".into(),
7549 params: vec![SqlValue::Text(edge_id.to_string())],
7550 label: Some("get_edge_including_deleted_namespace".into()),
7551 })
7552 .await?;
7553
7554 let Some(SqlValue::Text(record_ns)) = record_ns else {
7555 return Ok(None);
7556 };
7557 let record_tok = NamespaceToken::for_namespace(
7559 khive_types::Namespace::parse(&record_ns)
7560 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
7561 );
7562 Ok(self
7563 .graph(&record_tok)?
7564 .get_edge_including_deleted(LinkId::from(edge_id))
7565 .await?)
7566 }
7567
7568 pub async fn get_edge_by_natural_key_including_deleted(
7582 &self,
7583 token: &NamespaceToken,
7584 namespace: &str,
7585 source_id: Uuid,
7586 target_id: Uuid,
7587 relation: EdgeRelation,
7588 ) -> RuntimeResult<Option<Edge>> {
7589 Ok(self
7590 .graph(token)?
7591 .get_edge_by_natural_key_including_deleted(namespace, source_id, target_id, relation)
7592 .await?)
7593 }
7594
7595 pub const EDGE_LIST_MAX_LIMIT: u32 = 1000;
7600
7601 pub async fn list_edges(
7606 &self,
7607 token: &NamespaceToken,
7608 filter: crate::curation::EdgeListFilter,
7609 limit: u32,
7610 offset: u32,
7611 ) -> RuntimeResult<Vec<Edge>> {
7612 let limit = limit.min(Self::EDGE_LIST_MAX_LIMIT);
7613 let visible = token.visible_namespaces();
7614
7615 if let [ns] = visible {
7618 let temp = NamespaceToken::for_namespace(ns.clone());
7619 let page = self
7620 .graph(&temp)?
7621 .query_edges(
7622 filter.into(),
7623 vec![SortOrder {
7624 field: EdgeSortField::CreatedAt,
7625 direction: khive_storage::types::SortDirection::Asc,
7626 }],
7627 PageRequest {
7628 offset: offset.into(),
7629 limit,
7630 },
7631 )
7632 .await?;
7633 return Ok(page.items);
7634 }
7635
7636 let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
7643 let sort = vec![SortOrder {
7644 field: EdgeSortField::CreatedAt,
7645 direction: khive_storage::types::SortDirection::Asc,
7646 }];
7647 let graph = self.graph(token)?;
7648 match graph
7649 .query_edges_in_namespaces(
7650 &ns_strs,
7651 filter.clone().into(),
7652 sort.clone(),
7653 PageRequest {
7654 offset: offset.into(),
7655 limit,
7656 },
7657 )
7658 .await
7659 {
7660 Ok(page) => Ok(page.items),
7661 Err(khive_storage::StorageError::Unsupported { operation, .. })
7662 if operation == "query_edges_in_namespaces" =>
7663 {
7664 let fetch_limit = offset.saturating_add(limit);
7677 let mut namespace_prefixes = Vec::new();
7678 for ns in visible {
7679 let temp = NamespaceToken::for_namespace(ns.clone());
7680 let page = self
7681 .graph(&temp)?
7682 .query_edges(
7683 filter.clone().into(),
7684 sort.clone(),
7685 PageRequest {
7686 offset: 0,
7687 limit: fetch_limit,
7688 },
7689 )
7690 .await?;
7691 namespace_prefixes.push(page.items);
7692 }
7693 Ok(Self::merge_paged_namespace_edges(
7694 namespace_prefixes,
7695 offset,
7696 limit,
7697 ))
7698 }
7699 Err(error) => Err(error.into()),
7700 }
7701 }
7702
7703 fn merge_paged_namespace_edges(
7718 namespace_prefixes: Vec<Vec<Edge>>,
7719 offset: u32,
7720 limit: u32,
7721 ) -> Vec<Edge> {
7722 let mut results: Vec<Edge> = namespace_prefixes.into_iter().flatten().collect();
7723 results.sort_by_key(|e| (e.created_at, Uuid::from(e.id)));
7724 let start = (offset as usize).min(results.len());
7725 let end = (start + limit as usize).min(results.len());
7726 results[start..end].to_vec()
7727 }
7728
7729 pub async fn list_edges_after(
7741 &self,
7742 token: &NamespaceToken,
7743 filter: crate::curation::EdgeListFilter,
7744 after: Option<Uuid>,
7745 limit: u32,
7746 ) -> RuntimeResult<(Vec<Edge>, Option<Uuid>)> {
7747 let limit = limit.clamp(1, Self::EDGE_LIST_MAX_LIMIT);
7748 let visible = token.visible_namespaces();
7749 let limit_usize = limit as usize;
7750 let cursor_store = self.graph(token)?;
7751 let after = match after {
7752 Some(id) => {
7753 let edge = self
7754 .get_edge_including_deleted(token, id)
7755 .await?
7756 .ok_or_else(|| RuntimeError::NotFound(format!("edge cursor {id}")))?;
7757 Self::ensure_namespace_visible(&edge.namespace, token)?;
7758 let sequence = cursor_store.edge_sequence(id).await?.ok_or_else(|| {
7759 RuntimeError::Internal(format!(
7760 "edge cursor {id} has no insertion-sequence ledger row"
7761 ))
7762 })?;
7763 Some(SeekCursor { sequence, id })
7764 }
7765 None => None,
7766 };
7767
7768 if let [ns] = visible {
7769 let temp = NamespaceToken::for_namespace(ns.clone());
7770 let page = self
7771 .graph(&temp)?
7772 .query_edges_sequence_after(filter.into(), after, limit)
7773 .await?;
7774 return Ok((page.items, page.next_after.map(|cursor| cursor.id)));
7775 }
7776
7777 let probe_limit = limit.saturating_add(1);
7781 let mut results = Vec::new();
7782 for ns in visible {
7783 let temp = NamespaceToken::for_namespace(ns.clone());
7784 let page = self
7785 .graph(&temp)?
7786 .query_edges_sequence_after(filter.clone().into(), after, probe_limit)
7787 .await?;
7788 results.extend(page.items);
7789 }
7790 let ids = results
7791 .iter()
7792 .map(|edge| Uuid::from(edge.id))
7793 .collect::<Vec<_>>();
7794 let sequences = cursor_store
7795 .edge_sequences(&ids)
7796 .await?
7797 .into_iter()
7798 .collect::<HashMap<_, _>>();
7799 if let Some(missing) = ids.iter().find(|id| !sequences.contains_key(id)) {
7800 return Err(RuntimeError::Internal(format!(
7801 "edge {missing} has no insertion-sequence ledger row"
7802 )));
7803 }
7804 results.sort_by_key(|edge| {
7805 let id = Uuid::from(edge.id);
7806 (sequences[&id], id)
7807 });
7808 results.dedup_by_key(|e| Uuid::from(e.id));
7809 let has_more = results.len() > limit_usize;
7810 if has_more {
7811 results.truncate(limit_usize);
7812 }
7813 let next_after = if has_more {
7814 results.last().map(|e| Uuid::from(e.id))
7815 } else {
7816 None
7817 };
7818 Ok((results, next_after))
7819 }
7820
7821 pub async fn count_edges_by_relation(
7825 &self,
7826 token: &NamespaceToken,
7827 ) -> RuntimeResult<std::collections::HashMap<String, u64>> {
7828 let namespaces: Vec<String> = token
7829 .visible_namespaces()
7830 .iter()
7831 .map(|namespace| namespace.as_str().to_owned())
7832 .collect();
7833 let graph = self.graph(token)?;
7834 let counts = match graph
7835 .count_edges_by_relation_in_namespaces(&namespaces)
7836 .await
7837 {
7838 Ok(counts) => counts,
7839 Err(khive_storage::StorageError::Unsupported { operation, .. })
7840 if operation == "count_edges_by_relation_in_namespaces" =>
7841 {
7842 let mut totals = HashMap::new();
7843 for namespace in token.visible_namespaces() {
7844 let scoped = NamespaceToken::for_namespace(namespace.clone());
7845 for (relation, count) in self.graph(&scoped)?.count_edges_by_relation().await? {
7846 *totals.entry(relation).or_insert(0) += count;
7847 }
7848 }
7849 return Ok(totals
7850 .into_iter()
7851 .map(|(relation, count)| (relation.to_string(), count))
7852 .collect());
7853 }
7854 Err(error) => return Err(error.into()),
7855 };
7856 Ok(counts
7857 .into_iter()
7858 .map(|(relation, count)| (relation.to_string(), count))
7859 .collect())
7860 }
7861
7862 pub async fn count_edges_by_endpoint_base(
7871 &self,
7872 token: &NamespaceToken,
7873 ) -> RuntimeResult<khive_storage::types::EdgeEndpointBaseCounts> {
7874 use khive_storage::types::EdgeEndpointBaseCounts;
7875
7876 let namespaces: Vec<String> = token
7877 .visible_namespaces()
7878 .iter()
7879 .map(|namespace| namespace.as_str().to_owned())
7880 .collect();
7881 let graph = self.graph(token)?;
7882 match graph
7883 .count_edges_by_endpoint_base_in_namespaces(&namespaces)
7884 .await
7885 {
7886 Ok(counts) => Ok(counts),
7887 Err(khive_storage::StorageError::Unsupported { operation, .. })
7888 if operation == "count_edges_by_endpoint_base_in_namespaces"
7889 || operation == "count_edges_by_endpoint_base" =>
7890 {
7891 let mut totals = EdgeEndpointBaseCounts::default();
7892 for namespace in token.visible_namespaces() {
7893 let scoped = NamespaceToken::for_namespace(namespace.clone());
7894 let counts = self.graph(&scoped)?.count_edges_by_endpoint_base().await?;
7895 totals.entity_entity =
7896 totals.entity_entity.saturating_add(counts.entity_entity);
7897 totals.entity_note = totals.entity_note.saturating_add(counts.entity_note);
7898 totals.note_entity = totals.note_entity.saturating_add(counts.note_entity);
7899 totals.note_note = totals.note_note.saturating_add(counts.note_note);
7900 totals.unresolved = totals.unresolved.saturating_add(counts.unresolved);
7901 }
7902 Ok(totals)
7903 }
7904 Err(error) => Err(error.into()),
7905 }
7906 }
7907
7908 #[allow(clippy::too_many_arguments)]
7920 fn update_edge_symmetric_dml(
7921 conn: &rusqlite::Connection,
7922 ns: &str,
7923 edge_id_str: &str,
7924 canon_src_str: &str,
7925 canon_tgt_str: &str,
7926 relation_str: &str,
7927 weight: f64,
7928 metadata: Option<String>,
7929 expected_updated_at_micros: i64,
7930 expected_deleted_at_micros: Option<i64>,
7931 ) -> Result<SymmetricEdgeUpdateOutcome, SqliteError> {
7932 let minimum_updated_at_micros =
7946 expected_updated_at_micros.checked_add(1).ok_or_else(|| {
7947 SqliteError::InvalidData(format!(
7948 "update_edge: edge {edge_id_str} updated_at is already at i64::MAX \
7949 and cannot advance"
7950 ))
7951 })?;
7952 let now_ts = chrono::Utc::now()
7953 .timestamp_micros()
7954 .max(minimum_updated_at_micros);
7955
7956 let conflict_id: Option<String> = conn
7959 .query_row(
7960 khive_db::stores::graph::EDGE_SYMMETRIC_CONFLICT_PROBE_SQL,
7961 rusqlite::params![
7962 &ns,
7963 &canon_src_str,
7964 &canon_tgt_str,
7965 &relation_str,
7966 &edge_id_str
7967 ],
7968 |row| row.get(0),
7969 )
7970 .optional()
7971 .map_err(SqliteError::Rusqlite)?;
7972
7973 if let Some(existing_id) = conflict_id {
7974 let affected = conn
7993 .execute(
7994 khive_db::stores::graph::EDGE_SYMMETRIC_DELETE_NONCANONICAL_GUARDED_SQL,
7995 rusqlite::params![
7996 &ns,
7997 &edge_id_str,
7998 expected_updated_at_micros,
7999 expected_deleted_at_micros,
8000 ],
8001 )
8002 .map_err(SqliteError::Rusqlite)?;
8003 if affected == 0 {
8004 return Ok(SymmetricEdgeUpdateOutcome::Stale);
8005 }
8006 Ok(SymmetricEdgeUpdateOutcome::Absorbed(existing_id))
8007 } else {
8008 let affected = conn
8014 .execute(
8015 khive_db::stores::graph::EDGE_SYMMETRIC_UPDATE_INPLACE_SQL,
8016 rusqlite::params![
8017 &canon_src_str,
8018 &canon_tgt_str,
8019 &relation_str,
8020 weight,
8021 now_ts,
8022 metadata,
8023 &ns,
8024 &edge_id_str,
8025 expected_updated_at_micros,
8026 expected_deleted_at_micros,
8027 ],
8028 )
8029 .map_err(SqliteError::Rusqlite)?;
8030 if affected == 0 {
8031 return Ok(SymmetricEdgeUpdateOutcome::Stale);
8032 }
8033 Ok(SymmetricEdgeUpdateOutcome::Updated)
8034 }
8035 }
8036
8037 pub async fn update_edge(
8051 &self,
8052 token: &NamespaceToken,
8053 edge_id: Uuid,
8054 patch: crate::curation::EdgePatch,
8055 ) -> RuntimeResult<Edge> {
8056 let graph_for_fetch = self.graph(token)?;
8059 let mut edge = graph_for_fetch
8060 .get_edge(LinkId::from(edge_id))
8061 .await?
8062 .ok_or_else(|| crate::RuntimeError::NotFound(format!("edge {edge_id}")))?;
8063 let expected_updated_at = edge.updated_at;
8064 let expected_deleted_at = edge.deleted_at;
8065 #[cfg(test)]
8066 crate::curation::race_seam::pause_after_read().await;
8067
8068 let record_ns: String = edge.namespace.clone();
8073 let record_tok = token.with_namespace(
8074 khive_types::Namespace::parse(&record_ns)
8075 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
8076 );
8077 let graph = self.graph(&record_tok)?;
8078
8079 let mut changed_fields: Vec<&'static str> = Vec::new();
8080 if let Some(r) = patch.relation {
8081 self.validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, r)
8084 .await?;
8085 edge.relation = r;
8086 changed_fields.push("relation");
8087 }
8088 if let Some(w) = patch.weight {
8089 if !w.is_finite() || !(0.0..=1.0).contains(&w) {
8092 return Err(RuntimeError::InvalidInput(format!(
8093 "edge weight must be a finite value in [0.0, 1.0]; got {w}"
8094 )));
8095 }
8096 edge.weight = w;
8097 changed_fields.push("weight");
8098 }
8099 if let Some(props) = patch.properties {
8100 crate::secret_gate::reject_reserved_secret_gate_property(Some(&props))?;
8101 edge.metadata = Some(props);
8102 }
8103
8104 let (canon_src, canon_tgt) =
8114 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
8115
8116 if edge.relation.is_symmetric() {
8117 let ns = record_ns.clone();
8121 let edge_id_str = edge_id.to_string();
8122 let relation_str = edge.relation.to_string();
8123 let canon_src_str = canon_src.to_string();
8124 let canon_tgt_str = canon_tgt.to_string();
8125 let weight = edge.weight;
8126 let metadata = edge
8127 .metadata
8128 .as_ref()
8129 .map(|v| serde_json::to_string(v).unwrap_or_default());
8130
8131 let expected_updated_at_micros = expected_updated_at.timestamp_micros();
8132 let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
8133
8134 let pool = self.backend().pool_arc();
8135 let writer_task = pool
8136 .writer_task_for_runtime_write(RuntimeWriteOperation::UpdateSymmetricEdge)
8137 .map_err(RuntimeError::Storage)?;
8138
8139 let outcome: SymmetricEdgeUpdateOutcome = if let Some(writer_task) = writer_task {
8140 writer_task
8141 .send(move |conn| {
8142 Self::update_edge_symmetric_dml(
8143 conn,
8144 &ns,
8145 &edge_id_str,
8146 &canon_src_str,
8147 &canon_tgt_str,
8148 &relation_str,
8149 weight,
8150 metadata,
8151 expected_updated_at_micros,
8152 expected_deleted_at_micros,
8153 )
8154 .map_err(|e| {
8155 khive_storage::StorageError::driver(
8156 khive_storage::StorageCapability::Graph,
8157 "update_edge",
8158 e,
8159 )
8160 })
8161 })
8162 .await
8163 .map_err(RuntimeError::Storage)?
8164 } else {
8165 tokio::task::spawn_blocking(move || {
8166 let guard = pool.writer()?;
8167 guard.transaction(|conn| {
8168 Self::update_edge_symmetric_dml(
8169 conn,
8170 &ns,
8171 &edge_id_str,
8172 &canon_src_str,
8173 &canon_tgt_str,
8174 &relation_str,
8175 weight,
8176 metadata,
8177 expected_updated_at_micros,
8178 expected_deleted_at_micros,
8179 )
8180 })
8181 })
8182 .await
8183 .map_err(|e| {
8184 RuntimeError::Internal(format!("update_edge: spawn_blocking join: {e}"))
8185 })?
8186 .map_err(RuntimeError::Sqlite)?
8187 };
8188
8189 match outcome {
8190 SymmetricEdgeUpdateOutcome::Absorbed(sid) => {
8191 let surviving_uuid = Uuid::parse_str(&sid).map_err(|e| {
8198 RuntimeError::Internal(format!(
8199 "update_edge: surviving id parse failed: {e}"
8200 ))
8201 })?;
8202 edge = self
8203 .get_edge_including_deleted(&record_tok, surviving_uuid)
8204 .await?
8205 .ok_or_else(|| {
8206 RuntimeError::Internal(format!(
8207 "update_edge: surviving canonical row {surviving_uuid} vanished after update"
8208 ))
8209 })?;
8210 }
8211 SymmetricEdgeUpdateOutcome::Updated => {
8212 edge.source_id = canon_src;
8214 edge.target_id = canon_tgt;
8215 }
8216 SymmetricEdgeUpdateOutcome::Stale => {
8217 return Err(crate::curation::stale_edge_snapshot_error(edge_id));
8218 }
8219 }
8220 } else {
8221 let minimum_updated_at_micros = expected_updated_at
8234 .timestamp_micros()
8235 .checked_add(1)
8236 .ok_or_else(|| {
8237 RuntimeError::Internal(format!(
8238 "edge {edge_id} updated_at is already at i64::MAX and cannot advance"
8239 ))
8240 })?;
8241 let now_micros = chrono::Utc::now()
8242 .timestamp_micros()
8243 .max(minimum_updated_at_micros);
8244 edge.updated_at =
8245 chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
8246 RuntimeError::Internal(format!(
8247 "edge {edge_id}: computed updated_at {now_micros} is not a valid timestamp"
8248 ))
8249 })?;
8250 let persisted = graph
8251 .replace_edge_if_unchanged(edge.clone(), expected_updated_at, expected_deleted_at)
8252 .await?;
8253 if !persisted {
8254 return Err(crate::curation::stale_edge_snapshot_error(edge_id));
8255 }
8256 }
8257
8258 let event_store = self.events(&record_tok)?;
8260 let event = khive_storage::event::Event::new(
8261 record_ns.clone(),
8262 "update",
8263 EventKind::EdgeUpdated,
8264 SubstrateKind::Entity,
8265 "",
8266 )
8267 .with_target(edge_id)
8268 .with_payload(
8269 serde_json::json!({"id": edge_id, "namespace": record_ns, "changed_fields": changed_fields}),
8270 );
8271 event_store.append_event(event).await.map_err(|e| {
8272 RuntimeError::Internal(format!("update_edge: event store write failed: {e}"))
8273 })?;
8274
8275 Ok(edge)
8276 }
8277
8278 pub async fn delete_edge(
8289 &self,
8290 token: &NamespaceToken,
8291 edge_id: Uuid,
8292 hard: bool,
8293 ) -> RuntimeResult<bool> {
8294 let mode = if hard {
8295 DeleteMode::Hard
8296 } else {
8297 DeleteMode::Soft
8298 };
8299
8300 let edge = if hard {
8306 self.get_edge_including_deleted(token, edge_id).await?
8307 } else {
8308 self.get_edge(token, edge_id).await?
8309 };
8310 let Some(edge) = edge else {
8311 return Ok(false);
8312 };
8313
8314 let record_ns: String = edge.namespace.clone();
8316 let record_tok = token.with_namespace(
8317 khive_types::Namespace::parse(&record_ns)
8318 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
8319 );
8320 let graph = self.graph(&record_tok)?;
8321 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
8322
8323 let deleted = if hard {
8330 self.atomic_hard_delete_with_edge_purge(
8331 edge_hard_delete_statement(edge_id),
8332 edge_id,
8333 &record_ns,
8334 &actor,
8335 SubstrateKind::Entity,
8336 )
8337 .await?
8338 } else {
8339 graph.delete_edge(LinkId::from(edge_id), mode).await?
8340 };
8341 if deleted {
8342 let event_store = self.events(&record_tok)?;
8344 let event = khive_storage::event::Event::new(
8345 record_ns.clone(),
8346 "delete",
8347 EventKind::EdgeDeleted,
8348 SubstrateKind::Entity,
8349 "",
8350 )
8351 .with_target(edge_id)
8352 .with_payload(serde_json::json!({"id": edge_id, "namespace": record_ns, "hard": hard}));
8353 event_store.append_event(event).await.map_err(|e| {
8354 RuntimeError::Internal(format!("delete_edge: event store write failed: {e}"))
8355 })?;
8356 }
8357 Ok(deleted)
8358 }
8359
8360 pub async fn count_edges(
8362 &self,
8363 token: &NamespaceToken,
8364 filter: crate::curation::EdgeListFilter,
8365 ) -> RuntimeResult<u64> {
8366 let namespaces: Vec<String> = token
8367 .visible_namespaces()
8368 .iter()
8369 .map(|namespace| namespace.as_str().to_owned())
8370 .collect();
8371 let graph = self.graph(token)?;
8372 match graph
8373 .count_edges_in_namespaces(&namespaces, filter.clone().into())
8374 .await
8375 {
8376 Ok(count) => Ok(count),
8377 Err(khive_storage::StorageError::Unsupported { operation, .. })
8378 if operation == "count_edges_in_namespaces" =>
8379 {
8380 let mut total = 0;
8381 for namespace in token.visible_namespaces() {
8382 let scoped = NamespaceToken::for_namespace(namespace.clone());
8383 total += self
8384 .graph(&scoped)?
8385 .count_edges(filter.clone().into())
8386 .await?;
8387 }
8388 Ok(total)
8389 }
8390 Err(error) => Err(error.into()),
8391 }
8392 }
8393
8394 pub async fn build_edge(&self, token: &NamespaceToken, spec: &LinkSpec) -> RuntimeResult<Edge> {
8405 self.build_edge_with_endpoint_kinds(token, spec)
8406 .await
8407 .map(|(edge, _)| edge)
8408 }
8409
8410 async fn build_edge_with_endpoint_kinds(
8411 &self,
8412 token: &NamespaceToken,
8413 spec: &LinkSpec,
8414 ) -> RuntimeResult<(Edge, (EdgeEndpointKind, EdgeEndpointKind))> {
8415 validate_edge_metadata(spec.relation, spec.metadata.as_ref())?;
8416 let ns_str = match &spec.namespace {
8417 Some(s) => {
8418 let spec_ns = crate::Namespace::parse(s)
8419 .map_err(|e| RuntimeError::InvalidInput(format!("invalid namespace: {e}")))?;
8420 if &spec_ns != token.namespace() {
8421 return Err(RuntimeError::InvalidInput(
8422 "LinkSpec namespace does not match token namespace".into(),
8423 ));
8424 }
8425 s.as_str()
8426 }
8427 None => token.namespace().as_str(),
8428 };
8429 let endpoint_kinds = self
8430 .validate_edge_relation_endpoints(token, spec.source_id, spec.target_id, spec.relation)
8431 .await?;
8432 let (source_id, target_id) =
8433 canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
8434 let endpoint_kinds = canonical_edge_endpoint_kinds(
8435 spec.source_id,
8436 source_id,
8437 endpoint_kinds.0,
8438 endpoint_kinds.1,
8439 );
8440 let metadata = if spec.relation == EdgeRelation::DependsOn {
8441 match (
8446 self.resolve_edge_endpoint(token, source_id).await?,
8447 self.resolve_edge_endpoint(token, target_id).await?,
8448 ) {
8449 (Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
8450 merge_dependency_kind(&src_e.kind, &tgt_e.kind, spec.metadata.clone())
8451 }
8452 _ => spec.metadata.clone(),
8453 }
8454 } else {
8455 spec.metadata.clone()
8456 };
8457 validate_edge_metadata(spec.relation, metadata.as_ref())?;
8458 let now = chrono::Utc::now();
8459 Ok((
8460 Edge {
8461 id: LinkId::from(Uuid::new_v4()),
8462 namespace: ns_str.to_string(),
8463 source_id,
8464 target_id,
8465 relation: spec.relation,
8466 weight: spec.weight,
8467 created_at: now,
8468 updated_at: now,
8469 deleted_at: None,
8470 metadata,
8471 target_backend: None,
8472 },
8473 endpoint_kinds,
8474 ))
8475 }
8476
8477 pub async fn link_many(
8494 &self,
8495 token: &NamespaceToken,
8496 specs: Vec<LinkSpec>,
8497 ) -> RuntimeResult<Vec<Edge>> {
8498 self.link_many_observed(token, specs)
8499 .await
8500 .map(|rows| rows.into_iter().map(|row| row.edge).collect())
8501 }
8502
8503 pub async fn link_many_observed(
8507 &self,
8508 token: &NamespaceToken,
8509 specs: Vec<LinkSpec>,
8510 ) -> RuntimeResult<Vec<EdgeUpsertResult>> {
8511 self.link_many_guarded_observed(
8512 token,
8513 specs,
8514 GraphMutationPreconditions::default(),
8515 Vec::new(),
8516 )
8517 .await
8518 .map(|(rows, _)| rows)
8519 }
8520
8521 #[doc(hidden)]
8525 pub async fn link_many_guarded_observed(
8526 &self,
8527 token: &NamespaceToken,
8528 specs: Vec<LinkSpec>,
8529 preconditions: GraphMutationPreconditions,
8530 retirements: Vec<Edge>,
8531 ) -> RuntimeResult<(Vec<EdgeUpsertResult>, Vec<LinkId>)> {
8532 let namespace = token.namespace().as_str();
8533 let foreign_document = preconditions
8534 .document
8535 .as_ref()
8536 .is_some_and(|guard| guard.namespace != namespace);
8537 let foreign_edge = preconditions.edges.iter().any(|guard| {
8538 guard.namespace != namespace
8539 || guard
8540 .expected
8541 .as_ref()
8542 .is_some_and(|edge| edge.namespace != namespace)
8543 });
8544 if foreign_document
8545 || foreign_edge
8546 || retirements.iter().any(|edge| edge.namespace != namespace)
8547 {
8548 return Err(RuntimeError::InvalidInput(
8549 "guarded link namespace does not match token namespace".into(),
8550 ));
8551 }
8552 if specs.is_empty()
8553 && preconditions.document.is_none()
8554 && preconditions.edges.is_empty()
8555 && retirements.is_empty()
8556 {
8557 return Ok((Vec::new(), Vec::new()));
8558 }
8559 let mut edges = Vec::with_capacity(specs.len());
8560 let mut endpoint_kinds = Vec::with_capacity(specs.len());
8561 for spec in &specs {
8562 let (edge, kinds) = self.build_edge_with_endpoint_kinds(token, spec).await?;
8563 edges.push(edge);
8564 endpoint_kinds.push(kinds);
8565 }
8566 let requests = edges
8575 .into_iter()
8576 .zip(specs.iter())
8577 .map(|(edge, spec)| EdgeUpsertRequest {
8578 edge,
8579 resurrect: spec.resurrect,
8580 })
8581 .collect();
8582 let attribution = crate::EventAttribution::from_token(token);
8583 let outcome = compose_graph_mutation_events(
8584 self.backend(),
8585 GraphMutationRequest::Batch {
8586 requests,
8587 guard_endpoints: true,
8588 },
8589 preconditions,
8590 retirements,
8591 move |outcome| {
8592 let GraphMutationOutcome::Batch(batch) = &outcome.mutation else {
8593 return Err(Self::link_composition_shape_error(
8594 "expected a written batch",
8595 ));
8596 };
8597 if batch.rows.len() != endpoint_kinds.len() {
8598 return Err(Self::link_composition_shape_error(
8599 "edge result count differs from validated endpoint count",
8600 ));
8601 }
8602 let mut events = Vec::with_capacity(batch.rows.len() + outcome.retired.len());
8603 for (row, (source_kind, target_kind)) in batch.rows.iter().zip(endpoint_kinds) {
8604 events.push(Self::link_mutation_event(
8605 &attribution,
8606 row,
8607 source_kind,
8608 target_kind,
8609 ));
8610 }
8611 for edge in &outcome.retired {
8612 let edge_id = Uuid::from(edge.id);
8613 events.push(
8614 attribution.stamp(
8615 Event::new(
8616 edge.namespace.clone(),
8617 "delete",
8618 EventKind::EdgeDeleted,
8619 SubstrateKind::Entity,
8620 "",
8621 )
8622 .with_target(edge_id)
8623 .with_payload(serde_json::json!({
8624 "id": edge_id, "namespace": edge.namespace, "hard": false,
8625 })),
8626 ),
8627 );
8628 }
8629 Ok(events)
8630 },
8631 )
8632 .await?;
8633 let retired = outcome.retired.into_iter().map(|edge| edge.id).collect();
8634 let GraphMutationOutcome::Batch(outcome) = outcome.mutation else {
8635 return Err(RuntimeError::Internal(
8636 "link_many: unexpected composition outcome".into(),
8637 ));
8638 };
8639 if let Some(refusal) = outcome.refusal {
8640 return match refusal.reason {
8641 EdgeUpsertRefusal::MissingEndpoints(missing) => {
8642 Err(RuntimeError::GuardedWriteFailed(guarded_link_batch_failure(
8643 &specs[refusal.entry_index],
8644 refusal.entry_index,
8645 missing,
8646 )))
8647 }
8648 EdgeUpsertRefusal::ResurrectionRequired { edge } => {
8649 Err(RuntimeError::InvalidInput(format!(
8650 "batch entry {} targets soft-deleted edge {}; pass resurrect=true for that link",
8651 refusal.entry_index, edge.id
8652 )))
8653 }
8654 };
8655 }
8656 Ok((outcome.rows, retired))
8657 }
8658
8659 pub async fn link_commit_annotation_if_absent(
8664 &self,
8665 token: &NamespaceToken,
8666 commit_id: Uuid,
8667 project_id: Uuid,
8668 guard: CommitAnnotationGuard,
8669 ) -> RuntimeResult<CommitAnnotationInsertOutcome> {
8670 if !matches!(guard.expected_sha.len(), 40 | 64)
8671 || !guard
8672 .expected_sha
8673 .bytes()
8674 .all(|byte| byte.is_ascii_hexdigit())
8675 {
8676 return Err(RuntimeError::InvalidInput(
8677 "expected full commit SHA".into(),
8678 ));
8679 }
8680 let edge = self
8681 .build_edge(
8682 token,
8683 &LinkSpec {
8684 namespace: None,
8685 source_id: commit_id,
8686 target_id: project_id,
8687 relation: EdgeRelation::Annotates,
8688 weight: 1.0,
8689 metadata: None,
8690 resurrect: false,
8691 },
8692 )
8693 .await?;
8694 let attribution = crate::EventAttribution::from_token(token);
8695 let outcome = compose_graph_mutation_events(
8696 self.backend(),
8697 GraphMutationRequest::CommitAnnotation { edge, guard },
8698 GraphMutationPreconditions::default(),
8699 Vec::new(),
8700 move |outcome| match &outcome.mutation {
8701 GraphMutationOutcome::CommitAnnotation(CommitAnnotationInsertOutcome::Created(
8702 edge,
8703 )) => Ok(vec![Self::link_mutation_event(
8704 &attribution,
8705 &EdgeUpsertResult {
8706 edge: edge.clone(),
8707 disposition: EdgeUpsertDisposition::Created,
8708 previous: None,
8709 },
8710 EdgeEndpointKind::Note,
8711 EdgeEndpointKind::Entity,
8712 )]),
8713 _ => Err(Self::link_composition_shape_error(
8714 "expected a created annotation",
8715 )),
8716 },
8717 )
8718 .await?;
8719 let GraphMutationOutcome::CommitAnnotation(result) = outcome.mutation else {
8720 return Err(RuntimeError::Internal(
8721 "link annotation: unexpected composition outcome".into(),
8722 ));
8723 };
8724 Ok(result)
8725 }
8726
8727 pub async fn create_many(
8738 &self,
8739 token: &NamespaceToken,
8740 specs: Vec<EntityCreateSpec>,
8741 ) -> RuntimeResult<Vec<Entity>> {
8742 if specs.is_empty() {
8743 return Ok(vec![]);
8744 }
8745 let ns = token.namespace().as_str();
8746
8747 let mut entities = Vec::with_capacity(specs.len());
8751 for (index, spec) in specs.iter().enumerate() {
8752 entities.push(self.validate_bulk_entity(ns, spec, &format!("entity[{index}]"))?);
8753 }
8754
8755 #[cfg(any(test, feature = "fault-injection"))]
8756 let fts_many_inject = consume_fault(&FTS_FAIL_MANY_NS, ns);
8757 #[cfg(not(any(test, feature = "fault-injection")))]
8758 let fts_many_inject = false;
8759
8760 #[cfg(any(test, feature = "fault-injection"))]
8761 let fts_many_inject_partial = consume_fault(&FTS_FAIL_MANY_PARTIAL_NS, ns);
8762 #[cfg(not(any(test, feature = "fault-injection")))]
8763 let fts_many_inject_partial = false;
8764
8765 let injected_failure_index = if fts_many_inject {
8766 Some(0)
8767 } else if fts_many_inject_partial {
8768 Some(usize::from(entities.len() > 1))
8769 } else {
8770 None
8771 };
8772
8773 let _ = self.entities(token)?;
8774 let _ = self.text(token)?;
8775
8776 let plans = entities
8777 .iter()
8778 .enumerate()
8779 .map(|(index, entity)| {
8780 let mut plan = bulk_entity_plan(entity)?;
8781 if injected_failure_index == Some(index) {
8782 plan.statements.truncate(1);
8784 plan.statements.push(PlanStatement {
8785 statement: SqlStatement {
8786 sql:
8787 "INSERT INTO __khive_create_many_injected_failure__ DEFAULT VALUES"
8788 .to_string(),
8789 params: vec![],
8790 label: Some("fts-insert-injected-failure".to_string()),
8791 },
8792 guard: None,
8793 });
8794 }
8795 Ok(AtomicOpPlan::AddEntity(plan))
8796 })
8797 .collect::<RuntimeResult<Vec<_>>>()?;
8798
8799 match run_atomic_unit(self.sql().as_ref(), plans).await {
8800 Ok(AtomicRunOutcome::Committed { .. }) => Ok(entities),
8801 Ok(AtomicRunOutcome::RolledBack {
8802 failed_op_index,
8803 failure,
8804 }) => Err(RuntimeError::Internal(format!(
8805 "create_many: atomic batch rolled back at entity index {failed_op_index}: \
8806 {failure:?}"
8807 ))),
8808 Err(e) => Err(RuntimeError::Internal(format!(
8809 "create_many: atomic batch failed: {}",
8810 e.0
8811 ))),
8812 }
8813 }
8814
8815 fn validate_bulk_entity(
8820 &self,
8821 ns: &str,
8822 spec: &EntityCreateSpec,
8823 record: &str,
8824 ) -> RuntimeResult<Entity> {
8825 self.validate_entity_kind(&spec.kind)?;
8826 let validated_type =
8831 self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
8832 if spec.name.trim().is_empty() {
8833 return Err(RuntimeError::InvalidInput("name must not be empty".into()));
8834 }
8835 crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
8836 crate::secret_gate::check_at(&spec.name, record, "name")?;
8837 if let Some(d) = &spec.description {
8838 crate::secret_gate::check_at(d, record, "description")?;
8839 }
8840 if let Some(ref p) = spec.properties {
8841 crate::secret_gate::check_json_at(p, record, "properties")?;
8842 }
8843 crate::secret_gate::check_tags_at(&spec.tags, record, "tags")?;
8844
8845 let mut entity =
8846 Entity::new(ns, &spec.kind, &spec.name).with_entity_type(validated_type.as_deref());
8847 if let Some(d) = &spec.description {
8848 entity = entity.with_description(d);
8849 }
8850 if let Some(p) = spec.properties.clone() {
8851 entity = entity.with_properties(p);
8852 }
8853 if !spec.tags.is_empty() {
8854 entity = entity.with_tags(spec.tags.clone());
8855 }
8856 Ok(entity)
8857 }
8858
8859 pub async fn prepare_bulk_entity_plan(
8866 &self,
8867 token: &NamespaceToken,
8868 spec: EntityCreateSpec,
8869 ) -> RuntimeResult<(Entity, AtomicOpPlan)> {
8870 let entity = self.validate_bulk_entity(token.namespace().as_str(), &spec, "entity")?;
8871 let _ = self.entities(token)?;
8872 let _ = self.text(token)?;
8873
8874 let plan = AtomicOpPlan::AddEntity(bulk_entity_plan(&entity)?);
8875 Ok((entity, plan))
8876 }
8877
8878 pub async fn prepare_bulk_note_plan(
8887 &self,
8888 token: &NamespaceToken,
8889 spec: NoteCreateSpec,
8890 ) -> RuntimeResult<(Note, AtomicOpPlan)> {
8891 let mut candidate = Note::new(token.namespace().as_str(), &spec.kind, &spec.content);
8892 candidate.name = spec.name.clone();
8893 candidate.properties = spec.properties.clone();
8894 crate::note_write::validate_head(&candidate)?;
8895 let mut prepared = crate::atomic_message::prepare_atomic_notes(
8896 self,
8897 vec![crate::atomic_message::AtomicNoteSpec {
8898 token,
8899 id: None,
8900 kind: &spec.kind,
8901 name: spec.name.as_deref(),
8902 content: &spec.content,
8903 properties: spec.properties,
8904 }],
8905 crate::atomic_message::AtomicNoteOptions {
8906 salience: spec.salience,
8907 embed: Some(false),
8908 ..Default::default()
8909 },
8910 )
8911 .await?;
8912 match (prepared.notes.pop(), prepared.plans.pop()) {
8913 (Some(note), Some(plan)) if prepared.notes.is_empty() && prepared.plans.is_empty() => {
8914 Ok((note, plan))
8915 }
8916 _ => Err(RuntimeError::Internal(
8917 "bulk note preparation must yield exactly one note and one plan".into(),
8918 )),
8919 }
8920 }
8921}
8922
8923#[derive(Clone, Debug)]
8928pub struct NoteCreateSpec {
8929 pub kind: String,
8930 pub name: Option<String>,
8931 pub content: String,
8932 pub salience: Option<f64>,
8933 pub properties: Option<serde_json::Value>,
8934}
8935
8936fn bulk_entity_plan(entity: &Entity) -> RuntimeResult<AddEntityPlan> {
8937 crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
8938 let mut statements = vec![PlanStatement {
8939 statement: entity_upsert_statement(entity),
8940 guard: Some(AffectedRowGuard::exactly(1)),
8941 }];
8942 statements.extend(
8944 insert_document_statements("fts_entities", &entity_fts_document(entity))
8945 .into_iter()
8946 .map(|statement| PlanStatement {
8947 statement,
8948 guard: None,
8949 }),
8950 );
8951 Ok(AddEntityPlan {
8952 entity_id: entity.id,
8953 statements,
8954 post_commit: PostCommitEffect::None,
8955 })
8956}
8957
8958fn guarded_link_batch_failure(
8959 spec: &LinkSpec,
8960 entry_index: usize,
8961 missing: khive_storage::MissingEndpoints,
8962) -> GuardedWriteFailure {
8963 let (source_id, target_id) =
8966 canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
8967 GuardedWriteFailure {
8968 entry_index: Some(entry_index),
8969 missing_source: missing.source.then_some(source_id),
8970 missing_target: missing.target.then_some(target_id),
8971 }
8972}
8973
8974#[derive(Clone, Debug)]
8977pub struct LinkSpec {
8978 pub namespace: Option<String>,
8979 pub source_id: Uuid,
8980 pub target_id: Uuid,
8981 pub relation: EdgeRelation,
8982 pub weight: f64,
8983 pub metadata: Option<serde_json::Value>,
8984 pub resurrect: bool,
8985}
8986
8987#[derive(Clone, Debug)]
8995pub struct EntityCreateSpec {
8996 pub kind: String,
8997 pub entity_type: Option<String>,
8998 pub name: String,
8999 pub description: Option<String>,
9000 pub properties: Option<serde_json::Value>,
9001 pub tags: Vec<String>,
9002}
9003
9004#[cfg(test)]
9010#[path = "operations_tests.rs"]
9011mod tests;