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(any(test, feature = "fault-injection"))]
319mod fault_injection;
320mod resolve_uuid_or_prefix;
321
322#[cfg(any(test, feature = "fault-injection"))]
323pub use fault_injection::{
324 arm_entity_compensation_fail_scoped, arm_fts_fail_many_partial_scoped,
325 arm_fts_fail_many_scoped, arm_fts_fail_scoped, arm_fts_search_fail,
326 arm_prefix_resolve_fail_scoped, arm_rollback_cleanup_fail, arm_vector_fail_after,
327 arm_vector_fail_scoped, FaultInjectionArm,
328};
329#[cfg(test)]
330use fault_injection::{arm_fault, FaultArmSet, LINK_FAIL_AFTER};
331#[cfg(any(test, feature = "fault-injection"))]
332use fault_injection::{
333 consume_fault, ENTITY_COMPENSATION_FAIL_NS, FTS_FAIL_MANY_NS, FTS_FAIL_MANY_PARTIAL_NS,
334 FTS_FAIL_NS, FTS_SEARCH_FAIL_NS, PREFIX_RESOLVE_FAIL_NS, ROLLBACK_CLEANUP_FAIL_NS,
335 VECTOR_FAIL_AFTER, VECTOR_FAIL_NS,
336};
337#[cfg(any(test, feature = "fault-injection"))]
338pub(crate) use fault_injection::{consume_fts_fail_fault, consume_vector_fail_fault};
339
340#[derive(Clone, Debug)]
342pub struct NoteSearchHit {
343 pub note_id: Uuid,
344 pub score: DeterministicScore,
345 pub rank_score_kind: crate::RankScoreKind,
346 pub signals: crate::SearchSignals,
347 pub source: crate::SearchSource,
348 pub title: Option<String>,
349 pub snippet: Option<String>,
350}
351
352fn salience_weighted_rank(score: DeterministicScore, salience: Option<f64>) -> DeterministicScore {
353 const SCALE_RAW: i128 = 1_i128 << 32;
354 let salience = DeterministicScore::from_f64(salience.unwrap_or(0.5));
355 let weight_raw = SCALE_RAW / 2 + i128::from(salience.to_raw()) / 2;
356 let weighted_raw = i128::from(score.to_raw()) * weight_raw / SCALE_RAW;
359 DeterministicScore::from_raw(weighted_raw.clamp(
360 i128::from(DeterministicScore::NEG_INF.to_raw()),
361 i128::from(DeterministicScore::MAX.to_raw()),
362 ) as i64)
363}
364
365#[derive(Clone, Debug)]
369pub struct NoteSearchOutcome {
370 pub hits: Vec<NoteSearchHit>,
371 pub vector_error: Option<String>,
372}
373
374pub fn hex_prefix_to_uuid_pattern(prefix: &str) -> String {
385 if prefix.contains('-') {
386 return prefix.to_string();
387 }
388 const BOUNDARIES: [usize; 4] = [8, 13, 18, 23]; let mut out = String::with_capacity(36);
390 for c in prefix.chars() {
391 if BOUNDARIES.contains(&out.len()) {
392 out.push('-');
393 }
394 out.push(c);
395 }
396 out
397}
398
399pub fn uuid_prefix_bounds(prefix: &str) -> Option<(String, String)> {
409 const HYPHEN_POSITIONS: [usize; 4] = [8, 13, 18, 23];
410
411 let compact = if prefix.contains('-') {
412 if prefix.len() > 36 {
413 return None;
414 }
415 let mut compact = String::with_capacity(32);
416 for (index, byte) in prefix.bytes().enumerate() {
417 if HYPHEN_POSITIONS.contains(&index) {
418 if byte != b'-' {
419 return None;
420 }
421 } else if byte.is_ascii_hexdigit() {
422 compact.push(char::from(byte.to_ascii_lowercase()));
423 } else {
424 return None;
425 }
426 }
427 compact
428 } else {
429 if prefix.is_empty()
430 || prefix.len() > 32
431 || !prefix.bytes().all(|byte| byte.is_ascii_hexdigit())
432 {
433 return None;
434 }
435 prefix.to_ascii_lowercase()
436 };
437
438 if compact.is_empty() || compact.len() > 32 {
439 return None;
440 }
441
442 let lower = hex_prefix_to_uuid_pattern(&compact);
443 let mut successor = compact.into_bytes();
444 let mut carried_past_start = true;
445 for index in (0..successor.len()).rev() {
446 let next = match successor[index] {
447 b'0'..=b'8' | b'a'..=b'e' => Some(successor[index] + 1),
448 b'9' => Some(b'a'),
449 b'f' => None,
450 _ => return None,
451 };
452 if let Some(next) = next {
453 successor[index] = next;
454 successor.truncate(index + 1);
455 carried_past_start = false;
456 break;
457 }
458 }
459
460 let upper = if carried_past_start {
461 "g".to_string()
462 } else {
463 let compact_upper = String::from_utf8(successor).ok()?;
464 hex_prefix_to_uuid_pattern(&compact_upper)
465 };
466 Some((lower, upper))
467}
468
469fn resolve_prefix_statement(
470 table: &str,
471 has_deleted_at: bool,
472 include_deleted: bool,
473 namespaces: Option<&[String]>,
474 lower: &str,
475 upper: &str,
476) -> SqlStatement {
477 let namespace_clause = namespaces.map(|namespaces| {
478 let placeholders: Vec<String> = (0..namespaces.len())
479 .map(|index| format!("?{}", index + 3))
480 .collect();
481 format!(" AND namespace IN ({})", placeholders.join(", "))
482 });
483 let deleted_filter = if has_deleted_at && !include_deleted {
484 " AND deleted_at IS NULL"
485 } else {
486 ""
487 };
488 let mut params = vec![
489 SqlValue::Text(lower.to_owned()),
490 SqlValue::Text(upper.to_owned()),
491 ];
492 if let Some(namespaces) = namespaces {
493 params.extend(
494 namespaces
495 .iter()
496 .map(|namespace| SqlValue::Text(namespace.clone())),
497 );
498 }
499
500 SqlStatement {
501 sql: format!(
502 "SELECT id FROM {table} \
503 WHERE id >= ?1 AND id < ?2{namespace_clause}{deleted_filter} ORDER BY id LIMIT 2",
504 namespace_clause = namespace_clause.as_deref().unwrap_or("")
505 ),
506 params,
507 label: Some("resolve_prefix".into()),
508 }
509}
510
511fn text_preview(text: &str, max_chars: usize) -> Option<String> {
512 let trimmed = text.trim();
513 if trimmed.is_empty() {
514 None
515 } else {
516 Some(trimmed.chars().take(max_chars).collect())
517 }
518}
519
520fn normalize_symmetric_direction(
526 direction: Direction,
527 relations: Option<&[EdgeRelation]>,
528) -> Direction {
529 let Some(rels) = relations else {
530 return direction;
531 };
532 if rels.is_empty() {
533 return direction;
534 }
535 let all_symmetric = rels
536 .iter()
537 .all(|r| matches!(r, EdgeRelation::CompetesWith | EdgeRelation::ComposedWith));
538 if all_symmetric {
539 Direction::Both
540 } else {
541 direction
542 }
543}
544
545fn direction_sort_rank(direction: &Direction) -> u8 {
552 match direction {
553 Direction::Out => 0,
554 Direction::In => 1,
555 Direction::Both => 2,
556 }
557}
558
559fn note_title(note: &Note) -> Option<String> {
560 note.name
561 .clone()
562 .filter(|s| !s.trim().is_empty())
563 .or_else(|| Some(format!("[{}]", note.kind.as_str())))
564}
565
566fn note_snippet(note: &Note) -> Option<String> {
567 text_preview(¬e.content, 200)
568}
569
570#[derive(Clone, Debug)]
572pub enum Resolved {
573 Entity(Entity),
574 Note(Note),
575 Event(Event),
576 PackRecord {
583 pack: String,
584 kind: String,
585 data: serde_json::Value,
586 },
587}
588
589#[derive(Clone, Copy, Debug, Eq, PartialEq)]
596pub enum EdgeEndpointKind {
597 Entity,
598 Note,
599 Event,
600 Edge,
601}
602
603impl EdgeEndpointKind {
604 pub const fn name(self) -> &'static str {
606 match self {
607 Self::Entity => "entity",
608 Self::Note => "note",
609 Self::Event => "event",
610 Self::Edge => "edge",
611 }
612 }
613}
614
615fn resolved_pair(r: Option<&Resolved>) -> Option<(&'static str, &str, Option<&str>)> {
621 match r? {
622 Resolved::Entity(e) => Some(("entity", e.kind.as_str(), e.entity_type.as_deref())),
623 Resolved::Note(n) => Some(("note", n.kind.as_str(), None)),
624 Resolved::Event(_) => None,
625 Resolved::PackRecord { .. } => None,
626 }
627}
628
629pub fn endpoint_matches(
637 spec: &EndpointKind,
638 substrate: &str,
639 kind: &str,
640 entity_type: Option<&str>,
641) -> bool {
642 match spec {
643 EndpointKind::EntityOfKind(k) => substrate == "entity" && *k == kind,
644 EndpointKind::NoteOfKind(k) => substrate == "note" && *k == kind,
645 EndpointKind::EntityOfType {
646 kind: k,
647 entity_type: t,
648 } => substrate == "entity" && *k == kind && entity_type == Some(*t),
649 }
650}
651
652fn pattern_endpoint_matches(
666 spec: &EndpointKind,
667 substrate: &str,
668 kind: &str,
669 entity_type: Option<&str>,
670) -> bool {
671 match spec {
672 EndpointKind::EntityOfType {
673 kind: k,
674 entity_type: t,
675 } => substrate == "entity" && *k == kind && entity_type.is_none_or(|et| et == *t),
676 _ => endpoint_matches(spec, substrate, kind, entity_type),
677 }
678}
679
680pub fn accepted_pack_relations_for_entities(
694 rules: &[EdgeEndpointRule],
695 src_kind: &str,
696 src_entity_type: Option<&str>,
697 tgt_kind: &str,
698 tgt_entity_type: Option<&str>,
699) -> Vec<EdgeRelation> {
700 let mut relations: Vec<EdgeRelation> = rules
701 .iter()
702 .filter(|r| {
703 endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
704 && endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
705 })
706 .map(|r| r.relation)
707 .collect();
708 relations.sort_by_key(|r| r.as_str());
709 relations.dedup();
710 relations
711}
712
713pub fn accepted_entity_relations_for_entities(
724 rules: &[EdgeEndpointRule],
725 src_kind: &str,
726 src_entity_type: Option<&str>,
727 tgt_kind: &str,
728 tgt_entity_type: Option<&str>,
729) -> Vec<EdgeRelation> {
730 let mut relations: Vec<EdgeRelation> = BASE_ENTITY_ENDPOINT_RULES
731 .iter()
732 .filter(|(src, _relation, tgt)| (*src == "*" || *src == src_kind) && *tgt == tgt_kind)
733 .map(|(_src, relation, _tgt)| *relation)
734 .collect();
735 relations.extend(
736 accepted_pack_relations_for_entities(
737 rules,
738 src_kind,
739 src_entity_type,
740 tgt_kind,
741 tgt_entity_type,
742 )
743 .into_iter()
744 .filter(|relation| {
745 *relation != EdgeRelation::Annotates && !crate::pack::is_special_relation(*relation)
746 }),
747 );
748 relations.sort_by_key(|relation| relation.as_str());
749 relations.dedup();
750 relations
751}
752
753fn accepted_entity_relations_description(
754 rules: &[EdgeEndpointRule],
755 src_kind: &str,
756 src_entity_type: Option<&str>,
757 tgt_kind: &str,
758 tgt_entity_type: Option<&str>,
759) -> String {
760 let relations = accepted_entity_relations_for_entities(
761 rules,
762 src_kind,
763 src_entity_type,
764 tgt_kind,
765 tgt_entity_type,
766 );
767 if relations.is_empty() {
768 "none".to_string()
769 } else {
770 relations
771 .iter()
772 .map(EdgeRelation::as_str)
773 .collect::<Vec<_>>()
774 .join(", ")
775 }
776}
777
778fn accepted_pack_relations_for_pattern_entities(
784 rules: &[EdgeEndpointRule],
785 src_kind: &str,
786 src_entity_type: Option<&str>,
787 tgt_kind: &str,
788 tgt_entity_type: Option<&str>,
789) -> Vec<EdgeRelation> {
790 let mut relations: Vec<EdgeRelation> = rules
791 .iter()
792 .filter(|r| {
793 pattern_endpoint_matches(&r.source, "entity", src_kind, src_entity_type)
794 && pattern_endpoint_matches(&r.target, "entity", tgt_kind, tgt_entity_type)
795 })
796 .map(|r| r.relation)
797 .collect();
798 relations.sort_by_key(|r| r.as_str());
799 relations.dedup();
800 relations
801}
802
803fn accepted_entity_kind_pairs_for_relation(
814 pack_rules: &[EdgeEndpointRule],
815 relation: EdgeRelation,
816) -> Vec<(&'static str, &'static str)> {
817 let mut pairs = Vec::new();
818 for src in khive_types::EntityKind::ALL {
819 for tgt in khive_types::EntityKind::ALL {
820 let allowed = base_entity_rule_allows(src.name(), relation, tgt.name())
821 || (!crate::pack::is_special_relation(relation)
822 && accepted_pack_relations_for_pattern_entities(
823 pack_rules,
824 src.name(),
825 None,
826 tgt.name(),
827 None,
828 )
829 .contains(&relation));
830 if allowed {
831 pairs.push((src.name(), tgt.name()));
832 }
833 }
834 }
835 pairs
836}
837
838fn static_impossible_edge_pattern_warnings(
853 language: khive_query::QueryLanguage,
854 pattern: &khive_query::ast::MatchPattern,
855 pack_rules: &[EdgeEndpointRule],
856) -> Vec<String> {
857 use khive_query::ast::{EdgeDirection, PatternElement};
858
859 if language != khive_query::QueryLanguage::Gql {
860 return Vec::new();
861 }
862
863 let elements = &pattern.elements;
864 let mut warnings = Vec::new();
865
866 for (i, el) in elements.iter().enumerate() {
867 let PatternElement::Edge(edge) = el else {
868 continue;
869 };
870 if edge.relations.len() != 1 || edge.min_hops != 1 || edge.max_hops != 1 {
871 continue;
872 }
873 let (left, right) = match (elements.get(i.wrapping_sub(1)), elements.get(i + 1)) {
874 (Some(PatternElement::Node(l)), Some(PatternElement::Node(r))) => (l, r),
875 _ => continue,
876 };
877 let (src_node, tgt_node) = match edge.direction {
878 EdgeDirection::Out => (left, right),
879 EdgeDirection::In => (right, left),
880 EdgeDirection::Both => continue,
881 };
882 let (Some(src_raw), Some(tgt_raw)) = (src_node.kind.as_deref(), tgt_node.kind.as_deref())
883 else {
884 continue;
885 };
886 let (Ok(src_kind), Ok(tgt_kind)) = (
887 src_raw.parse::<khive_types::EntityKind>(),
888 tgt_raw.parse::<khive_types::EntityKind>(),
889 ) else {
890 continue;
891 };
892 let Ok(relation) = edge.relations[0].parse::<EdgeRelation>() else {
893 continue;
894 };
895
896 let possible = base_entity_rule_allows(src_kind.name(), relation, tgt_kind.name())
897 || (!crate::pack::is_special_relation(relation)
898 && accepted_pack_relations_for_pattern_entities(
899 pack_rules,
900 src_kind.name(),
901 src_node.entity_type.as_deref(),
902 tgt_kind.name(),
903 tgt_node.entity_type.as_deref(),
904 )
905 .contains(&relation));
906 if possible {
907 continue;
908 }
909
910 let accepted = accepted_entity_kind_pairs_for_relation(pack_rules, relation);
911 let accepted_str = if accepted.is_empty() {
912 "none".to_string()
913 } else {
914 accepted
915 .iter()
916 .map(|(s, t)| format!("{s}->{t}"))
917 .collect::<Vec<_>>()
918 .join(", ")
919 };
920 warnings.push(format!(
921 "pattern ({src})-[:{relation}]->({tgt}) can never match: '{relation}' does not accept \
922 {src}->{tgt} endpoints; accepted source->target kinds for '{relation}': {accepted_str}",
923 src = src_kind.name(),
924 tgt = tgt_kind.name(),
925 ));
926 }
927
928 warnings
929}
930
931fn pack_rule_allows(
934 rules: &[EdgeEndpointRule],
935 relation: EdgeRelation,
936 src: Option<&Resolved>,
937 tgt: Option<&Resolved>,
938) -> bool {
939 let Some((src_sub, src_kind, src_type)) = resolved_pair(src) else {
940 return false;
941 };
942 let Some((tgt_sub, tgt_kind, tgt_type)) = resolved_pair(tgt) else {
943 return false;
944 };
945 rules.iter().any(|r| {
946 r.relation == relation
947 && endpoint_matches(&r.source, src_sub, src_kind, src_type)
948 && endpoint_matches(&r.target, tgt_sub, tgt_kind, tgt_type)
949 })
950}
951
952pub const BASE_ENTITY_ENDPOINT_RULES: &[(&str, EdgeRelation, &str)] = &[
962 ("concept", EdgeRelation::Contains, "concept"),
964 ("project", EdgeRelation::Contains, "project"),
965 ("project", EdgeRelation::Contains, "artifact"),
966 ("org", EdgeRelation::Contains, "project"),
967 ("org", EdgeRelation::Contains, "service"),
968 ("concept", EdgeRelation::PartOf, "concept"),
969 ("project", EdgeRelation::PartOf, "project"),
970 ("project", EdgeRelation::PartOf, "org"),
971 ("*", EdgeRelation::InstanceOf, "concept"),
972 ("service", EdgeRelation::InstanceOf, "project"),
973 ("document", EdgeRelation::LinksTo, "document"),
978 ("concept", EdgeRelation::LocatedIn, "concept"),
981 ("org", EdgeRelation::LocatedIn, "concept"),
982 ("person", EdgeRelation::Owns, "org"),
984 ("org", EdgeRelation::Owns, "org"),
985 ("concept", EdgeRelation::Extends, "concept"),
987 ("concept", EdgeRelation::VariantOf, "concept"),
988 ("artifact", EdgeRelation::VariantOf, "artifact"),
989 ("concept", EdgeRelation::IntroducedBy, "document"),
990 ("concept", EdgeRelation::IntroducedBy, "person"),
991 ("artifact", EdgeRelation::IntroducedBy, "document"),
992 ("project", EdgeRelation::IntroducedBy, "document"),
993 ("service", EdgeRelation::IntroducedBy, "document"),
996 ("document", EdgeRelation::IntroducedBy, "person"),
997 ("document", EdgeRelation::IntroducedBy, "org"),
998 ("concept", EdgeRelation::IntroducedBy, "org"),
999 ("artifact", EdgeRelation::DerivedFrom, "dataset"),
1001 ("artifact", EdgeRelation::DerivedFrom, "document"),
1002 ("artifact", EdgeRelation::DerivedFrom, "project"),
1003 ("artifact", EdgeRelation::DerivedFrom, "artifact"),
1004 ("document", EdgeRelation::DerivedFrom, "document"),
1007 ("document", EdgeRelation::Precedes, "document"),
1009 ("dataset", EdgeRelation::Precedes, "dataset"),
1010 ("artifact", EdgeRelation::Precedes, "artifact"),
1011 ("service", EdgeRelation::Precedes, "service"),
1012 ("project", EdgeRelation::Precedes, "project"),
1013 ("project", EdgeRelation::DependsOn, "project"),
1015 ("service", EdgeRelation::DependsOn, "project"),
1016 ("service", EdgeRelation::DependsOn, "service"),
1017 ("service", EdgeRelation::DependsOn, "artifact"),
1018 ("service", EdgeRelation::DependsOn, "dataset"),
1019 ("artifact", EdgeRelation::DependsOn, "project"),
1020 ("artifact", EdgeRelation::DependsOn, "service"),
1021 ("document", EdgeRelation::DependsOn, "document"),
1022 ("concept", EdgeRelation::Enables, "concept"),
1023 ("service", EdgeRelation::Enables, "concept"),
1024 ("dataset", EdgeRelation::Enables, "concept"),
1025 ("project", EdgeRelation::Implements, "concept"),
1027 ("service", EdgeRelation::Implements, "concept"),
1028 ("concept", EdgeRelation::CompetesWith, "concept"),
1030 ("project", EdgeRelation::CompetesWith, "project"),
1031 ("service", EdgeRelation::CompetesWith, "service"),
1032 ("org", EdgeRelation::CompetesWith, "org"),
1033 ("concept", EdgeRelation::ComposedWith, "concept"),
1034 ("project", EdgeRelation::ComposedWith, "project"),
1035 ("concept", EdgeRelation::Supersedes, "concept"),
1037 ("document", EdgeRelation::Supersedes, "document"),
1038 ("artifact", EdgeRelation::Supersedes, "artifact"),
1039 ("service", EdgeRelation::Supersedes, "service"),
1040 ("dataset", EdgeRelation::Supersedes, "dataset"),
1041 ("concept", EdgeRelation::Supports, "concept"),
1043 ("document", EdgeRelation::Supports, "concept"),
1044 ("dataset", EdgeRelation::Supports, "concept"),
1045 ("artifact", EdgeRelation::Supports, "concept"),
1046 ("concept", EdgeRelation::Refutes, "concept"),
1047 ("document", EdgeRelation::Refutes, "concept"),
1048 ("dataset", EdgeRelation::Refutes, "concept"),
1049 ("artifact", EdgeRelation::Refutes, "concept"),
1050];
1051
1052pub fn base_entity_endpoint_rules() -> &'static [(&'static str, EdgeRelation, &'static str)] {
1058 BASE_ENTITY_ENDPOINT_RULES
1059}
1060
1061pub fn base_entity_rule_allows(src_kind: &str, relation: EdgeRelation, tgt_kind: &str) -> bool {
1067 BASE_ENTITY_ENDPOINT_RULES.iter().any(|(src, rel, tgt)| {
1068 *rel == relation && (*src == "*" || *src == src_kind) && *tgt == tgt_kind
1069 })
1070}
1071
1072pub(crate) fn canonical_edge_endpoints(
1078 relation: EdgeRelation,
1079 source_id: Uuid,
1080 target_id: Uuid,
1081) -> (Uuid, Uuid) {
1082 relation.canonical_endpoints(source_id, target_id)
1083}
1084
1085pub(crate) fn canonical_edge_endpoint_kinds(
1087 requested_source_id: Uuid,
1088 canonical_source_id: Uuid,
1089 source_kind: EdgeEndpointKind,
1090 target_kind: EdgeEndpointKind,
1091) -> (EdgeEndpointKind, EdgeEndpointKind) {
1092 if requested_source_id == canonical_source_id {
1093 (source_kind, target_kind)
1094 } else {
1095 (target_kind, source_kind)
1096 }
1097}
1098
1099pub(crate) fn infer_dependency_kind(src_kind: &str, tgt_kind: &str) -> Option<&'static str> {
1105 match (src_kind, tgt_kind) {
1106 ("project", "project") => Some("build"),
1107 ("service", "service") => Some("runtime"),
1108 ("service", "dataset") => Some("data"),
1109 ("service", "artifact") => Some("artifact"),
1110 ("artifact", "project") | ("artifact", "service") => Some("tooling"),
1111 ("document", "document") => Some("normative"),
1112 _ => None,
1113 }
1114}
1115
1116pub(crate) fn merge_dependency_kind(
1126 src_kind: &str,
1127 tgt_kind: &str,
1128 metadata: Option<serde_json::Value>,
1129) -> Option<serde_json::Value> {
1130 let metadata = metadata.filter(|value| !value.is_null());
1133 if let Some(ref m) = metadata {
1134 if m.get("dependency_kind").is_some() {
1135 return metadata;
1136 }
1137 }
1138 let Some(inferred) = infer_dependency_kind(src_kind, tgt_kind) else {
1139 return metadata;
1140 };
1141 let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
1142 if let Some(o) = obj.as_object_mut() {
1143 o.insert("dependency_kind".to_string(), serde_json::json!(inferred));
1144 }
1145 Some(obj)
1146}
1147
1148pub fn merge_entry_metadata(
1162 metadata: Option<serde_json::Value>,
1163 dependency_kind: Option<String>,
1164) -> RuntimeResult<Option<serde_json::Value>> {
1165 validate_metadata_shape(metadata.as_ref())?;
1166 let metadata = metadata.filter(|value| !value.is_null());
1167 let Some(dk) = dependency_kind else {
1168 return Ok(metadata);
1169 };
1170 let mut obj = metadata.unwrap_or_else(|| serde_json::json!({}));
1171 let map = obj
1172 .as_object_mut()
1173 .ok_or_else(|| RuntimeError::InvalidInput("metadata must be a JSON object".into()))?;
1174 map.entry("dependency_kind".to_string())
1175 .or_insert_with(|| serde_json::json!(dk));
1176 Ok(Some(obj))
1177}
1178
1179const VALID_DEPENDENCY_KINDS: &[&str] = &[
1181 "build",
1182 "runtime",
1183 "data",
1184 "artifact",
1185 "tooling",
1186 "normative",
1187];
1188
1189pub(crate) fn validate_edge_weight(weight: f64) -> RuntimeResult<()> {
1195 if !khive_types::validate_edge_weight(weight) {
1196 return Err(RuntimeError::InvalidInput(format!(
1197 "edge weight must be finite and in [0.0, 1.0], got {weight}"
1198 )));
1199 }
1200 Ok(())
1201}
1202
1203fn validate_metadata_shape(metadata: Option<&serde_json::Value>) -> RuntimeResult<()> {
1204 if metadata.is_some_and(|value| !value.is_null() && !value.is_object()) {
1205 return Err(RuntimeError::InvalidInput(
1206 "metadata must be a JSON object".into(),
1207 ));
1208 }
1209 Ok(())
1210}
1211
1212pub(crate) fn validate_edge_metadata(
1217 relation: EdgeRelation,
1218 metadata: Option<&serde_json::Value>,
1219) -> RuntimeResult<()> {
1220 validate_metadata_shape(metadata)?;
1221 let Some(meta) = metadata.filter(|value| !value.is_null()) else {
1222 return Ok(());
1223 };
1224 let object = meta.as_object().expect("validated metadata object");
1225 if object
1226 .get("optional")
1227 .is_some_and(|value| !value.is_boolean())
1228 {
1229 return Err(RuntimeError::InvalidInput(
1230 "metadata.optional must be a boolean".into(),
1231 ));
1232 }
1233 if let Some(dk) = meta.get("dependency_kind") {
1234 if relation != EdgeRelation::DependsOn {
1235 return Err(RuntimeError::InvalidInput(format!(
1236 "dependency_kind is only valid on depends_on edges (got {})",
1237 relation.as_str()
1238 )));
1239 }
1240 let dk_str = dk
1241 .as_str()
1242 .ok_or_else(|| RuntimeError::InvalidInput("dependency_kind must be a string".into()))?;
1243 if !VALID_DEPENDENCY_KINDS.contains(&dk_str) {
1244 return Err(RuntimeError::InvalidInput(format!(
1245 "unknown dependency_kind {dk_str:?}; valid: {}",
1246 VALID_DEPENDENCY_KINDS.join(" | ")
1247 )));
1248 }
1249 }
1250 Ok(())
1251}
1252
1253fn note_graph_name(note: &Note) -> String {
1254 note.name
1255 .as_deref()
1256 .filter(|name| !name.trim().is_empty())
1257 .map(str::to_owned)
1258 .unwrap_or_else(|| format!("[{}]", note.kind))
1259}
1260
1261fn merge_traversal_paths_by_root(paths: Vec<GraphPath>, limit: Option<u32>) -> Vec<GraphPath> {
1266 let mut order: Vec<Uuid> = Vec::new();
1267 let mut merged: HashMap<Uuid, GraphPath> = HashMap::new();
1268 let mut node_index: HashMap<Uuid, HashMap<Uuid, usize>> = HashMap::new();
1272
1273 for path in paths {
1274 let existing = merged.entry(path.root_id).or_insert_with(|| {
1275 order.push(path.root_id);
1276 GraphPath {
1277 root_id: path.root_id,
1278 nodes: Vec::new(),
1279 total_weight: 0.0,
1280 }
1281 });
1282 let index = node_index.entry(path.root_id).or_default();
1283 for node in path.nodes {
1284 match index.get(&node.node_id) {
1285 Some(&i) => {
1286 if node.depth < existing.nodes[i].depth {
1287 existing.nodes[i] = node;
1288 }
1289 }
1290 None => {
1291 index.insert(node.node_id, existing.nodes.len());
1292 existing.nodes.push(node);
1293 }
1294 }
1295 }
1296 }
1297
1298 order
1299 .into_iter()
1300 .filter_map(|root_id| merged.remove(&root_id))
1301 .map(|mut path| {
1302 path.nodes.sort_by_key(|n| n.depth);
1304 if let Some(lim) = limit {
1305 let lim = lim as usize;
1306 let mut non_root_kept = 0usize;
1307 path.nodes.retain(|n| {
1308 if n.depth == 0 {
1309 return true;
1310 }
1311 if non_root_kept < lim {
1312 non_root_kept += 1;
1313 true
1314 } else {
1315 false
1316 }
1317 });
1318 }
1319 recompute_total_weight(&mut path);
1320 path
1321 })
1322 .collect()
1323}
1324
1325fn recompute_total_weight(path: &mut GraphPath) {
1334 path.total_weight = path.nodes.iter().map(|n| n.weight).fold(0.0_f64, f64::max);
1335}
1336
1337async fn drain_embed_join_set<T: Send + 'static>(
1351 mut join_set: tokio::task::JoinSet<(usize, RuntimeResult<T>)>,
1352 model_count: usize,
1353) -> RuntimeResult<Vec<T>> {
1354 let mut vectors: Vec<Option<T>> = (0..model_count).map(|_| None).collect();
1355
1356 while let Some(joined) = join_set.join_next().await {
1357 match joined {
1358 Ok((idx, Ok(vector))) => vectors[idx] = Some(vector),
1359 Ok((_idx, Err(e))) => {
1360 join_set.abort_all();
1361 return Err(e);
1362 }
1363 Err(join_err) => {
1364 join_set.abort_all();
1365 return Err(RuntimeError::Internal(format!(
1366 "embed task panicked: {join_err}"
1367 )));
1368 }
1369 }
1370 }
1371
1372 Ok(vectors
1373 .into_iter()
1374 .map(|v| v.expect("every model index observed exactly once by join_set drain"))
1375 .collect())
1376}
1377
1378impl KhiveRuntime {
1379 async fn compensate_entity_create(
1382 &self,
1383 token: &NamespaceToken,
1384 entity_id: Uuid,
1385 namespace: &str,
1386 vector_models: &[String],
1387 ) -> Vec<String> {
1388 let mut cleanup_errors = Vec::new();
1389
1390 #[cfg(any(test, feature = "fault-injection"))]
1391 let entity_delete_injected = consume_fault(&ENTITY_COMPENSATION_FAIL_NS, namespace);
1392 #[cfg(not(any(test, feature = "fault-injection")))]
1393 let entity_delete_injected = false;
1394
1395 if entity_delete_injected {
1396 cleanup_errors.push("entity row delete: injected compensation failure".to_string());
1397 } else {
1398 match self.entities(token) {
1399 Ok(store) => {
1400 if let Err(error) = store.delete_entity(entity_id, DeleteMode::Hard).await {
1401 cleanup_errors.push(format!("entity row delete: {error}"));
1402 }
1403 }
1404 Err(error) => cleanup_errors.push(format!("entity store access: {error}")),
1405 }
1406 }
1407
1408 match self.text(token) {
1409 Ok(fts) => {
1410 if let Err(error) = fts.delete_document(namespace, entity_id).await {
1411 cleanup_errors.push(format!("FTS document delete: {error}"));
1412 }
1413 }
1414 Err(error) => cleanup_errors.push(format!("FTS store access: {error}")),
1415 }
1416
1417 for model_name in vector_models {
1418 match self.vectors_for_model(token, model_name) {
1419 Ok(vectors) => {
1420 if let Err(error) = vectors.delete(entity_id).await {
1421 cleanup_errors
1422 .push(format!("vector delete for model {model_name}: {error}"));
1423 }
1424 }
1425 Err(error) => cleanup_errors.push(format!(
1426 "vector store access for model {model_name}: {error}"
1427 )),
1428 }
1429 }
1430
1431 cleanup_errors
1432 }
1433
1434 fn entity_create_failure(
1435 entity_id: Uuid,
1436 primary: RuntimeError,
1437 cleanup_errors: Vec<String>,
1438 ) -> RuntimeError {
1439 if cleanup_errors.is_empty() {
1440 primary
1441 } else {
1442 RuntimeError::Khive(KhiveError::internal(format!(
1443 "create_entity indexing failed for record {entity_id}; primary failure: \
1444 {primary}; compensation failure(s): {}; partial persistence is possible; \
1445 inspect and reconcile this record before retrying",
1446 cleanup_errors.join("; ")
1447 )))
1448 }
1449 }
1450
1451 pub async fn claim_entity_if_absent(
1454 &self,
1455 token: &NamespaceToken,
1456 spec: EntityClaimSpec,
1457 ) -> RuntimeResult<(Entity, bool)> {
1458 self.validate_entity_kind(&spec.kind)?;
1459 let entity_type =
1460 self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
1461 crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
1462 crate::secret_gate::check_at(&spec.name, "entity", "name")?;
1463 if let Some(description) = &spec.description {
1464 crate::secret_gate::check_at(description, "entity", "description")?;
1465 }
1466 if let Some(properties) = &spec.properties {
1467 crate::secret_gate::check_json_at(properties, "entity", "properties")?;
1468 }
1469 crate::secret_gate::check_tags_at(&spec.tags, "entity", "tags")?;
1470
1471 let mut proposed = Entity::new(token.namespace().as_str(), &spec.kind, &spec.name);
1472 proposed.id = spec.id;
1473 proposed.entity_type = entity_type.clone();
1474 proposed.description = spec.description;
1475 proposed.properties = spec.properties;
1476 proposed.tags = spec.tags;
1477
1478 let store = self.entities(token)?;
1479 let inserted = store.insert_entity_if_absent(proposed.clone()).await?;
1480 let entity = if inserted {
1481 proposed
1482 } else {
1483 store
1484 .get_entity_including_deleted(spec.id)
1485 .await?
1486 .ok_or_else(|| {
1487 RuntimeError::Internal(format!(
1488 "entity claim {} lost but the winning row is missing",
1489 spec.id
1490 ))
1491 })?
1492 };
1493 if entity.deleted_at.is_some() {
1494 return Err(RuntimeError::InvalidInput(format!(
1495 "entity claim {} is soft-deleted; restore it explicitly",
1496 entity.id
1497 )));
1498 }
1499 if entity.namespace != token.namespace().as_str()
1500 || entity.kind != spec.kind
1501 || entity.entity_type.as_deref() != entity_type.as_deref()
1502 || !entity.name.eq_ignore_ascii_case(&spec.name)
1503 || !entity
1504 .tags
1505 .iter()
1506 .any(|tag| tag.eq_ignore_ascii_case(&spec.identity_tag))
1507 {
1508 return Err(RuntimeError::InvalidInput(format!(
1509 "entity claim {} belongs to a different record",
1510 entity.id
1511 )));
1512 }
1513
1514 self.ensure_claimed_entity_create_event(token, &entity)
1515 .await?;
1516 self.reindex_claimed_entity(token, &entity).await?;
1517 Ok((entity, inserted))
1518 }
1519
1520 pub async fn ensure_claimed_entity_create_event(
1523 &self,
1524 token: &NamespaceToken,
1525 entity: &Entity,
1526 ) -> RuntimeResult<()> {
1527 if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
1528 return Err(RuntimeError::InvalidInput(format!(
1529 "entity {} is not a live row in the write namespace",
1530 entity.id
1531 )));
1532 }
1533 let events = self.events(token).map_err(|error| {
1534 RuntimeError::Internal(format!(
1535 "entity {} persists but its create event store is unavailable: {error}",
1536 entity.id
1537 ))
1538 })?;
1539 let filter = EventFilter {
1540 target_id: Some(entity.id),
1541 kinds: vec![EventKind::EntityCreated],
1542 verbs: vec!["create".into()],
1543 substrates: vec![SubstrateKind::Entity],
1544 after: Some(entity.created_at.saturating_sub(1)),
1545 ..EventFilter::default()
1546 };
1547 let page = PageRequest {
1548 offset: 0,
1549 limit: 1,
1550 };
1551 if !events
1552 .query_events(filter.clone(), page.clone())
1553 .await?
1554 .items
1555 .is_empty()
1556 {
1557 return Ok(());
1558 }
1559
1560 let mut event = Event::new(
1561 entity.namespace.clone(),
1562 "create",
1563 EventKind::EntityCreated,
1564 SubstrateKind::Entity,
1565 "",
1566 )
1567 .with_target(entity.id)
1568 .with_payload(serde_json::json!({
1569 "id": entity.id,
1570 "namespace": &entity.namespace,
1571 "kind": &entity.kind,
1572 }));
1573 let event_seed = Uuid::new_v5(&Uuid::NAMESPACE_URL, b"khive:claimed-entity-create:v1");
1574 let mut event_key = Vec::with_capacity(24);
1575 event_key.extend_from_slice(entity.id.as_bytes());
1576 event_key.extend_from_slice(&entity.created_at.to_be_bytes());
1577 event.id = Uuid::new_v5(&event_seed, &event_key);
1578 if let Err(error) = events.append_event(event).await {
1579 if events.query_events(filter, page).await?.items.is_empty() {
1580 return Err(RuntimeError::Internal(format!(
1581 "entity {} persists but its create event failed: {error}",
1582 entity.id
1583 )));
1584 }
1585 }
1586 Ok(())
1587 }
1588
1589 pub async fn reindex_claimed_entity(
1592 &self,
1593 token: &NamespaceToken,
1594 entity: &Entity,
1595 ) -> RuntimeResult<()> {
1596 if entity.namespace != token.namespace().as_str() || entity.deleted_at.is_some() {
1597 return Err(RuntimeError::InvalidInput(format!(
1598 "entity {} is not a live row in the write namespace",
1599 entity.id
1600 )));
1601 }
1602 let doc = entity_fts_document(entity);
1603 let embed_body = doc.body.clone();
1604 #[cfg(any(test, feature = "fault-injection"))]
1605 let fts_inject = consume_fault(&FTS_FAIL_NS, &entity.namespace);
1606 #[cfg(not(any(test, feature = "fault-injection")))]
1607 let fts_inject = false;
1608 let fts_result = if fts_inject {
1609 Err(RuntimeError::Internal("injected FTS failure".into()))
1610 } else {
1611 match self.text(token) {
1612 Ok(text) => text.upsert_document(doc).await.map_err(Into::into),
1613 Err(error) => Err(error),
1614 }
1615 };
1616 fts_result.map_err(|error| {
1617 RuntimeError::Internal(format!(
1618 "entity {} persists but its text index failed: {error}",
1619 entity.id
1620 ))
1621 })?;
1622
1623 for model_name in self.registered_embedding_model_names() {
1624 let outcome = self
1625 .embed_document_with_model_outcome_for_token(token, &model_name, &embed_body)
1626 .await
1627 .map_err(|error| {
1628 RuntimeError::Internal(format!(
1629 "entity {} persists but model {model_name} embedding failed: {error}",
1630 entity.id
1631 ))
1632 })?;
1633 #[cfg(any(test, feature = "fault-injection"))]
1634 let vector_inject = consume_fault(&VECTOR_FAIL_NS, &entity.namespace);
1635 #[cfg(not(any(test, feature = "fault-injection")))]
1636 let vector_inject = false;
1637 if vector_inject {
1638 return Err(RuntimeError::Internal(format!(
1639 "entity {} persists but model {model_name} vector indexing failed: injected vector failure",
1640 entity.id
1641 )));
1642 }
1643 self.vectors_for_model(token, &model_name)
1644 .map_err(|error| {
1645 RuntimeError::Internal(format!(
1646 "entity {} persists but model {model_name} vector store is unavailable: {error}",
1647 entity.id
1648 ))
1649 })?
1650 .insert(
1651 entity.id,
1652 SubstrateKind::Entity,
1653 &entity.namespace,
1654 "entity.body",
1655 vec![outcome.vector],
1656 )
1657 .await
1658 .map_err(|error| {
1659 RuntimeError::Internal(format!(
1660 "entity {} persists but model {model_name} vector indexing failed: {error}",
1661 entity.id
1662 ))
1663 })?;
1664 }
1665 Ok(())
1666 }
1667
1668 #[allow(clippy::too_many_arguments)]
1679 #[cfg(test)]
1680 pub(crate) async fn create_entity(
1681 &self,
1682 token: &NamespaceToken,
1683 kind: &str,
1684 entity_type: Option<&str>,
1685 name: &str,
1686 description: Option<&str>,
1687 properties: Option<serde_json::Value>,
1688 tags: Vec<String>,
1689 ) -> RuntimeResult<Entity> {
1690 let (entity, _, degradations) = self
1691 .create_entity_with_embedding_report_inner(
1692 token,
1693 kind,
1694 entity_type,
1695 name,
1696 description,
1697 properties,
1698 tags,
1699 Vec::new(),
1700 )
1701 .await?;
1702 legacy_post_commit_result("create_entity", entity.id, entity, degradations)
1703 }
1704
1705 #[allow(clippy::too_many_arguments)]
1718 pub async fn create_entity_with_attachments(
1719 &self,
1720 token: &NamespaceToken,
1721 kind: &str,
1722 entity_type: Option<&str>,
1723 name: &str,
1724 description: Option<&str>,
1725 properties: Option<serde_json::Value>,
1726 tags: Vec<String>,
1727 attachments: Vec<NewAttachment>,
1728 ) -> RuntimeResult<Entity> {
1729 let (entity, embedding, degradations) = self
1730 .create_entity_with_attachments_inner(
1731 token,
1732 kind,
1733 entity_type,
1734 name,
1735 description,
1736 properties,
1737 tags,
1738 attachments,
1739 )
1740 .await?;
1741 legacy_post_commit_result_with_embedding(
1742 "create_entity_with_attachments",
1743 entity.id,
1744 entity,
1745 embedding,
1746 degradations,
1747 )
1748 }
1749
1750 #[allow(clippy::too_many_arguments)]
1752 pub async fn create_entity_with_attachments_and_report(
1753 &self,
1754 token: &NamespaceToken,
1755 kind: &str,
1756 entity_type: Option<&str>,
1757 name: &str,
1758 description: Option<&str>,
1759 properties: Option<serde_json::Value>,
1760 tags: Vec<String>,
1761 attachments: Vec<NewAttachment>,
1762 ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
1763 let (entity, embedding, degradations) = self
1764 .create_entity_with_attachments_inner(
1765 token,
1766 kind,
1767 entity_type,
1768 name,
1769 description,
1770 properties,
1771 tags,
1772 attachments,
1773 )
1774 .await?;
1775 legacy_post_commit_result(
1776 "create_entity_with_attachments_and_report",
1777 entity.id,
1778 (entity, embedding),
1779 degradations,
1780 )
1781 }
1782
1783 #[allow(clippy::too_many_arguments)]
1784 async fn create_entity_with_attachments_inner(
1785 &self,
1786 token: &NamespaceToken,
1787 kind: &str,
1788 entity_type: Option<&str>,
1789 name: &str,
1790 description: Option<&str>,
1791 properties: Option<serde_json::Value>,
1792 tags: Vec<String>,
1793 attachments: Vec<NewAttachment>,
1794 ) -> RuntimeResult<(
1795 Entity,
1796 crate::retrieval::EmbeddingTruncationReport,
1797 Vec<PostCommitDegradation>,
1798 )> {
1799 drop(self.attachments()?);
1803 let blob_store = self.blob_store().ok_or_else(|| {
1804 RuntimeError::Unconfigured(
1805 "create_entity_with_attachments requires an installed BlobStore".to_string(),
1806 )
1807 })?;
1808 let mut roles = std::collections::HashSet::with_capacity(attachments.len());
1809 for attachment in &attachments {
1810 attachment.validate()?;
1811 if !roles.insert(attachment.role.as_str()) {
1812 return Err(RuntimeError::InvalidInput(format!(
1813 "duplicate attachment role {:?}",
1814 attachment.role
1815 )));
1816 }
1817 }
1818 for attachment in &attachments {
1819 if !blob_store.exists(&attachment.content_ref).await? {
1820 return Err(RuntimeError::InvalidInput(format!(
1821 "create_entity_with_attachments requires a published blob; no object exists for {}",
1822 attachment.content_ref
1823 )));
1824 }
1825 }
1826 let validated_type = self.validate_entity_type_for_kind(kind, entity_type)?;
1827 let (entity, embedding, degradations) = self
1828 .create_entity_with_embedding_report_inner(
1829 token,
1830 kind,
1831 validated_type.as_deref(),
1832 name,
1833 description,
1834 properties,
1835 tags,
1836 attachments,
1837 )
1838 .await?;
1839 Ok((entity, embedding, degradations))
1840 }
1841
1842 #[allow(clippy::too_many_arguments)]
1843 pub async fn create_entity_with_embedding_report(
1844 &self,
1845 token: &NamespaceToken,
1846 kind: &str,
1847 entity_type: Option<&str>,
1848 name: &str,
1849 description: Option<&str>,
1850 properties: Option<serde_json::Value>,
1851 tags: Vec<String>,
1852 ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
1853 let (entity, embedding, degradations) = self
1854 .create_entity_with_embedding_report_inner(
1855 token,
1856 kind,
1857 entity_type,
1858 name,
1859 description,
1860 properties,
1861 tags,
1862 Vec::new(),
1863 )
1864 .await?;
1865 legacy_post_commit_result(
1866 "create_entity_with_embedding_report",
1867 entity.id,
1868 (entity, embedding),
1869 degradations,
1870 )
1871 }
1872
1873 #[allow(clippy::too_many_arguments)]
1875 pub async fn create_entity_with_post_commit_report(
1876 &self,
1877 token: &NamespaceToken,
1878 kind: &str,
1879 entity_type: Option<&str>,
1880 name: &str,
1881 description: Option<&str>,
1882 properties: Option<serde_json::Value>,
1883 tags: Vec<String>,
1884 ) -> RuntimeResult<(
1885 Entity,
1886 crate::retrieval::EmbeddingTruncationReport,
1887 Vec<PostCommitDegradation>,
1888 )> {
1889 self.create_entity_with_embedding_report_inner(
1890 token,
1891 kind,
1892 entity_type,
1893 name,
1894 description,
1895 properties,
1896 tags,
1897 Vec::new(),
1898 )
1899 .await
1900 }
1901
1902 #[allow(clippy::too_many_arguments)]
1903 async fn create_entity_with_embedding_report_inner(
1904 &self,
1905 token: &NamespaceToken,
1906 kind: &str,
1907 entity_type: Option<&str>,
1908 name: &str,
1909 description: Option<&str>,
1910 properties: Option<serde_json::Value>,
1911 tags: Vec<String>,
1912 attachments: Vec<NewAttachment>,
1913 ) -> RuntimeResult<(
1914 Entity,
1915 crate::retrieval::EmbeddingTruncationReport,
1916 Vec<PostCommitDegradation>,
1917 )> {
1918 self.validate_entity_kind(kind)?;
1919 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
1920 crate::secret_gate::check_at(name, "entity", "name")?;
1922 if let Some(d) = description {
1923 crate::secret_gate::check_at(d, "entity", "description")?;
1924 }
1925 if let Some(ref p) = properties {
1926 crate::secret_gate::check_json_at(p, "entity", "properties")?;
1927 }
1928 crate::secret_gate::check_tags_at(&tags, "entity", "tags")?;
1929 let ns = token.namespace().as_str();
1930 let mut entity = Entity::new(ns, kind, name).with_entity_type(entity_type);
1931 if let Some(d) = description {
1932 entity = entity.with_description(d);
1933 }
1934 if let Some(p) = properties {
1935 entity = entity.with_properties(p);
1936 }
1937 if !tags.is_empty() {
1938 entity = entity.with_tags(tags);
1939 }
1940 let projected_content_ref = attachments
1941 .iter()
1942 .find(|attachment| attachment.role == "content")
1943 .map(|attachment| attachment.content_ref.to_string());
1944 let attachment_rows = attachments
1945 .into_iter()
1946 .map(|attachment| {
1947 Attachment::from_new(
1948 entity.id,
1949 AttachmentSubstrate::Entity,
1950 attachment,
1951 entity.created_at,
1952 )
1953 })
1954 .collect();
1955 self.entities(token)?
1956 .upsert_entity_with_attachments(entity.clone(), attachment_rows)
1957 .await?;
1958 entity.content_ref = projected_content_ref;
1959
1960 let doc = entity_fts_document(&entity);
1961 let embed_body = doc.body.clone();
1962
1963 {
1965 #[cfg(any(test, feature = "fault-injection"))]
1966 let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
1967 #[cfg(not(any(test, feature = "fault-injection")))]
1968 let fts_inject = false;
1969 let fts_result: RuntimeResult<()> = if fts_inject {
1970 Err(RuntimeError::Internal("injected FTS failure".to_string()))
1971 } else {
1972 match self.text(token) {
1973 Ok(fts) => fts.upsert_document(doc).await.map_err(RuntimeError::from),
1974 Err(e) => Err(e),
1975 }
1976 };
1977 if let Err(e) = fts_result {
1978 let cleanup_errors = self
1979 .compensate_entity_create(token, entity.id, ns, &[])
1980 .await;
1981 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
1982 }
1983 }
1984
1985 let embed_model_names = {
1988 let names = self.registered_embedding_model_names();
1989 if names.is_empty() {
1990 vec![]
1991 } else {
1992 names
1993 }
1994 };
1995
1996 let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
1997 if embed_model_names.len() == 1 {
1998 let model_name = &embed_model_names[0];
1999 let vec_result = self
2000 .embed_document_with_model_outcome_for_token(token, model_name, &embed_body)
2001 .await;
2002
2003 #[cfg(any(test, feature = "fault-injection"))]
2004 let vec_inject = consume_fault(&VECTOR_FAIL_NS, ns);
2005 #[cfg(not(any(test, feature = "fault-injection")))]
2006 let vec_inject = false;
2007 let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
2008 if vec_inject {
2009 Err(RuntimeError::Internal(
2010 "injected vector failure".to_string(),
2011 ))
2012 } else {
2013 vec_result
2014 };
2015
2016 let single_result: RuntimeResult<()> = match vec_result {
2017 Ok(outcome) => {
2018 embedding_report.observe(&outcome);
2019 match self.vectors_for_model(token, model_name) {
2020 Ok(vs) => vs
2021 .insert(
2022 entity.id,
2023 SubstrateKind::Entity,
2024 ns,
2025 "entity.body",
2026 vec![outcome.vector],
2027 )
2028 .await
2029 .map_err(RuntimeError::from),
2030 Err(e) => Err(e),
2031 }
2032 }
2033 Err(e) => Err(e),
2034 };
2035 if let Err(e) = single_result {
2036 let cleanup_errors = self
2037 .compensate_entity_create(
2038 token,
2039 entity.id,
2040 ns,
2041 std::slice::from_ref(model_name),
2042 )
2043 .await;
2044 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2045 }
2046 } else if !embed_model_names.is_empty() {
2047 let rt_clone = self.clone();
2050 let body_owned = embed_body.clone();
2051 let usage_ctx = crate::usage::current();
2052 let mut join_set = tokio::task::JoinSet::new();
2053 for (idx, model_name) in embed_model_names.iter().enumerate() {
2054 let rt = rt_clone.clone();
2055 let text = body_owned.clone();
2056 let name = model_name.clone();
2057 let ctx = usage_ctx.clone();
2058 let token = (*token).clone();
2059 join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
2060 let fut = rt.embed_document_with_model_outcome_for_token(&token, &name, &text);
2061 let result = match ctx {
2062 Some(ctx) => crate::usage::scope(ctx, fut).await,
2063 None => fut.await,
2064 };
2065 (idx, result)
2066 }));
2067 }
2068 let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
2072 Ok(outcomes) => outcomes,
2073 Err(e) => {
2074 let cleanup_errors = self
2075 .compensate_entity_create(token, entity.id, ns, &[])
2076 .await;
2077 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2078 }
2079 };
2080 let mut inserted_models: Vec<String> = Vec::with_capacity(embed_model_names.len());
2082 for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
2083 embedding_report.observe(&outcome);
2084 #[cfg(any(test, feature = "fault-injection"))]
2086 let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
2087 Some(0) => {
2088 cell.set(None);
2089 true
2090 }
2091 Some(n) => {
2092 cell.set(Some(n - 1));
2093 false
2094 }
2095 None => false,
2096 });
2097 #[cfg(not(any(test, feature = "fault-injection")))]
2098 let count_inject = false;
2099
2100 let insert_result = if count_inject {
2101 Err(RuntimeError::Internal(
2102 "injected vector insert failure".to_string(),
2103 ))
2104 } else {
2105 match self.vectors_for_model(token, model_name) {
2106 Ok(vs) => vs
2107 .insert(
2108 entity.id,
2109 SubstrateKind::Entity,
2110 ns,
2111 "entity.body",
2112 vec![outcome.vector],
2113 )
2114 .await
2115 .map_err(RuntimeError::from),
2116 Err(e) => Err(e),
2117 }
2118 };
2119 if let Err(e) = insert_result {
2120 let mut cleanup_models = inserted_models.clone();
2123 cleanup_models.push(model_name.clone());
2124 let cleanup_errors = self
2125 .compensate_entity_create(token, entity.id, ns, &cleanup_models)
2126 .await;
2127 return Err(Self::entity_create_failure(entity.id, e, cleanup_errors));
2128 }
2129 inserted_models.push(model_name.clone());
2130 }
2131 }
2132
2133 let created_event = khive_storage::event::Event::new(
2140 entity.namespace.clone(),
2141 "create",
2142 EventKind::EntityCreated,
2143 SubstrateKind::Entity,
2144 "",
2145 )
2146 .with_target(entity.id)
2147 .with_payload(serde_json::json!({
2148 "id": entity.id,
2149 "namespace": entity.namespace,
2150 "kind": entity.kind,
2151 }));
2152 let event_result = match self.events(token) {
2153 Ok(store) => store
2154 .append_event(created_event)
2155 .await
2156 .map_err(RuntimeError::from),
2157 Err(error) => Err(error),
2158 };
2159 let mut degradations = Vec::new();
2160 if let Err(error) = event_result {
2161 record_post_commit_degradation(
2162 &mut degradations,
2163 "create_entity",
2164 entity.id,
2165 "event_append",
2166 error,
2167 );
2168 }
2169
2170 Ok((entity, embedding_report, degradations))
2171 }
2172
2173 pub async fn get_entity(&self, token: &NamespaceToken, id: Uuid) -> RuntimeResult<Entity> {
2186 let store = self.entities(token)?;
2187 if let Some(entity) = store.get_entity(id).await? {
2188 return Ok(entity);
2189 }
2190 if let Some(tombstone) = store.get_entity_including_deleted(id).await? {
2191 if let Some(kept_id) = tombstone.merged_into {
2192 return Err(RuntimeError::NotFound(format!(
2193 "{id} was merged into {kept_id}; query the kept id"
2194 )));
2195 }
2196 }
2197 Err(RuntimeError::NotFound(format!("entity {id}")))
2198 }
2199
2200 pub async fn get_entity_including_deleted(
2204 &self,
2205 token: &NamespaceToken,
2206 id: Uuid,
2207 ) -> RuntimeResult<Option<Entity>> {
2208 self.entities(token)?
2209 .get_entity_including_deleted(id)
2210 .await
2211 .map_err(Into::into)
2212 }
2213
2214 pub async fn get_note_including_deleted(
2218 &self,
2219 token: &NamespaceToken,
2220 id: Uuid,
2221 ) -> RuntimeResult<Option<khive_storage::note::Note>> {
2222 self.notes(token)?
2223 .get_note_including_deleted(id)
2224 .await
2225 .map_err(Into::into)
2226 }
2227
2228 pub async fn get_entities_by_ids(
2232 &self,
2233 token: &NamespaceToken,
2234 ids: &[Uuid],
2235 ) -> RuntimeResult<Vec<Entity>> {
2236 if ids.is_empty() {
2237 return Ok(vec![]);
2238 }
2239 let filter = EntityFilter {
2240 ids: ids.to_vec(),
2241 ..Default::default()
2242 };
2243 let page = self
2244 .entities(token)?
2245 .query_entities(
2246 token.namespace().as_str(),
2247 filter,
2248 PageRequest {
2249 offset: 0,
2250 limit: ids.len() as u32,
2251 },
2252 )
2253 .await?;
2254 Ok(page.items)
2255 }
2256
2257 async fn get_entities_by_ids_visible(
2266 &self,
2267 token: &NamespaceToken,
2268 ids: &[Uuid],
2269 ) -> RuntimeResult<Vec<Entity>> {
2270 if ids.is_empty() {
2271 return Ok(vec![]);
2272 }
2273 let namespaces: Vec<String> = token
2274 .visible_namespaces()
2275 .iter()
2276 .map(|ns| ns.as_str().to_owned())
2277 .collect();
2278 let filter = EntityFilter {
2279 ids: ids.to_vec(),
2280 namespaces,
2281 ..Default::default()
2282 };
2283 let page = self
2284 .entities(token)?
2285 .query_entities(
2286 token.namespace().as_str(),
2287 filter,
2288 PageRequest {
2289 offset: 0,
2290 limit: ids.len() as u32,
2291 },
2292 )
2293 .await?;
2294 Ok(page.items)
2295 }
2296
2297 pub(crate) fn ensure_namespace(record_ns: &str, caller_primary_ns: &str) -> RuntimeResult<()> {
2306 if record_ns == caller_primary_ns {
2307 return Ok(());
2308 }
2309 Err(RuntimeError::NotFound("not found in this namespace".into()))
2310 }
2311
2312 pub(crate) fn ensure_namespace_visible(
2318 record_ns: &str,
2319 token: &NamespaceToken,
2320 ) -> RuntimeResult<()> {
2321 for ns in token.visible_namespaces() {
2322 if record_ns == ns.as_str() {
2323 return Ok(());
2324 }
2325 }
2326 Err(RuntimeError::NotFound("not found in this namespace".into()))
2327 }
2328
2329 pub async fn list_entities(
2336 &self,
2337 token: &NamespaceToken,
2338 kind: Option<&str>,
2339 entity_type: Option<&str>,
2340 limit: u32,
2341 offset: u32,
2342 ) -> RuntimeResult<Vec<Entity>> {
2343 let filter = EntityFilter {
2344 kinds: kind
2345 .map(|value| vec![value.to_string()])
2346 .unwrap_or_default(),
2347 entity_types: entity_type
2348 .map(|value| vec![value.to_string()])
2349 .unwrap_or_default(),
2350 legacy_entity_type_fallback: true,
2351 ..Default::default()
2352 };
2353 self.list_entities_filtered(token, filter, limit, offset)
2354 .await
2355 }
2356
2357 pub async fn list_entities_filtered(
2360 &self,
2361 token: &NamespaceToken,
2362 mut filter: EntityFilter,
2363 limit: u32,
2364 offset: u32,
2365 ) -> RuntimeResult<Vec<Entity>> {
2366 filter.namespaces = token
2367 .visible_namespaces()
2368 .iter()
2369 .map(|namespace| namespace.as_str().to_owned())
2370 .collect();
2371 let page = self
2372 .entities(token)?
2373 .query_entities_count_free(
2374 token.namespace().as_str(),
2375 filter,
2376 PageRequest {
2377 offset: offset.into(),
2378 limit,
2379 },
2380 )
2381 .await?;
2382 Ok(page.items)
2383 }
2384
2385 pub async fn list_entities_after(
2392 &self,
2393 token: &NamespaceToken,
2394 kind: Option<&str>,
2395 entity_type: Option<&str>,
2396 tags_any: &[String],
2397 after: Option<Uuid>,
2398 limit: u32,
2399 ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
2400 let filter = EntityFilter {
2401 kinds: kind
2402 .map(|value| vec![value.to_string()])
2403 .unwrap_or_default(),
2404 entity_types: entity_type
2405 .map(|value| vec![value.to_string()])
2406 .unwrap_or_default(),
2407 legacy_entity_type_fallback: true,
2408 tags_any: tags_any.to_vec(),
2409 ..Default::default()
2410 };
2411 self.list_entities_after_filtered(token, filter, after, limit)
2412 .await
2413 }
2414
2415 pub async fn list_entities_after_filtered(
2418 &self,
2419 token: &NamespaceToken,
2420 mut filter: EntityFilter,
2421 after: Option<Uuid>,
2422 limit: u32,
2423 ) -> RuntimeResult<(Vec<Entity>, Option<Uuid>)> {
2424 let store = self.entities(token)?;
2425 let after = match after {
2426 Some(id) => {
2427 let entity = self
2428 .get_entity_including_deleted(token, id)
2429 .await?
2430 .ok_or_else(|| RuntimeError::NotFound(format!("entity cursor {id}")))?;
2431 Self::ensure_namespace_visible(&entity.namespace, token)?;
2432 let sequence = store.entity_sequence(id).await?.ok_or_else(|| {
2433 RuntimeError::Internal(format!(
2434 "entity cursor {id} has no insertion-sequence ledger row"
2435 ))
2436 })?;
2437 Some(SeekCursor { sequence, id })
2438 }
2439 None => None,
2440 };
2441 filter.namespaces = token
2442 .visible_namespaces()
2443 .iter()
2444 .map(|namespace| namespace.as_str().to_owned())
2445 .collect();
2446 let page = store
2447 .query_entities_after(token.namespace().as_str(), filter, after, limit)
2448 .await?;
2449 Ok((page.items, page.next_after.map(|cursor| cursor.id)))
2450 }
2451
2452 pub async fn list_entities_tagged(
2459 &self,
2460 token: &NamespaceToken,
2461 kind: Option<&str>,
2462 domain_tag: Option<&str>,
2463 limit: u32,
2464 offset: u32,
2465 ) -> RuntimeResult<Vec<Entity>> {
2466 let ns_strs: Vec<String> = token
2467 .visible_namespaces()
2468 .iter()
2469 .map(|ns| ns.as_str().to_owned())
2470 .collect();
2471 let filter = EntityFilter {
2472 kinds: match kind {
2473 Some(k) => vec![k.to_string()],
2474 None => vec![],
2475 },
2476 tags_any: match domain_tag {
2477 Some(t) if !t.is_empty() => vec![t.to_string()],
2478 _ => vec![],
2479 },
2480 namespaces: ns_strs,
2481 ..Default::default()
2482 };
2483 let page = self
2484 .entities(token)?
2485 .query_entities_count_free(
2486 token.namespace().as_str(),
2487 filter,
2488 PageRequest {
2489 offset: offset.into(),
2490 limit,
2491 },
2492 )
2493 .await?;
2494 Ok(page.items)
2495 }
2496
2497 pub async fn count_entities_tagged(
2502 &self,
2503 token: &NamespaceToken,
2504 kind: Option<&str>,
2505 domain_tag: Option<&str>,
2506 ) -> RuntimeResult<u64> {
2507 let ns_strs: Vec<String> = token
2508 .visible_namespaces()
2509 .iter()
2510 .map(|ns| ns.as_str().to_owned())
2511 .collect();
2512 let filter = EntityFilter {
2513 kinds: match kind {
2514 Some(k) => vec![k.to_string()],
2515 None => vec![],
2516 },
2517 tags_any: match domain_tag {
2518 Some(t) if !t.is_empty() => vec![t.to_string()],
2519 _ => vec![],
2520 },
2521 namespaces: ns_strs,
2522 ..Default::default()
2523 };
2524 Ok(self
2525 .entities(token)?
2526 .count_entities(token.namespace().as_str(), filter)
2527 .await?)
2528 }
2529
2530 pub async fn list_events(
2532 &self,
2533 token: &NamespaceToken,
2534 filter: EventFilter,
2535 page: PageRequest,
2536 ) -> RuntimeResult<Page<Event>> {
2537 self.events(token)?
2538 .query_events(filter, page)
2539 .await
2540 .map_err(Into::into)
2541 }
2542
2543 pub(crate) async fn validate_edge_relation_endpoints(
2561 &self,
2562 token: &NamespaceToken,
2563 source_id: Uuid,
2564 target_id: Uuid,
2565 relation: EdgeRelation,
2566 ) -> RuntimeResult<(EdgeEndpointKind, EdgeEndpointKind)> {
2567 if source_id == target_id {
2568 return Err(RuntimeError::InvalidInput(
2569 "self-loop edges are not allowed: source_id and target_id must be different".into(),
2570 ));
2571 }
2572 if relation == EdgeRelation::Annotates {
2573 match self.resolve_edge_endpoint(token, source_id).await? {
2577 Some(Resolved::Note(_)) => {}
2578 Some(_) => {
2579 return Err(RuntimeError::InvalidInput(format!(
2580 "annotates source {source_id} must be a note"
2581 )));
2582 }
2583 None => {
2584 if self.get_edge(token, source_id).await?.is_some() {
2586 return Err(RuntimeError::InvalidInput(format!(
2587 "annotates source {source_id} must be a note"
2588 )));
2589 }
2590 return Err(RuntimeError::NotFound(format!(
2591 "link source {source_id} not found"
2592 )));
2593 }
2594 }
2595 let target_kind = match self.resolve_edge_endpoint(token, target_id).await? {
2597 Some(Resolved::Entity(_)) => EdgeEndpointKind::Entity,
2598 Some(Resolved::Note(_)) => EdgeEndpointKind::Note,
2599 Some(Resolved::Event(_)) => EdgeEndpointKind::Event,
2600 Some(Resolved::PackRecord { .. }) => {
2601 return Err(RuntimeError::InvalidInput(
2602 "pack-private record is not a valid edge endpoint for annotates".into(),
2603 ));
2604 }
2605 None => match self.get_edge(token, target_id).await {
2606 Ok(Some(_)) => EdgeEndpointKind::Edge,
2607 Ok(None) | Err(RuntimeError::NotFound(_)) => {
2608 return Err(RuntimeError::NotFound(format!(
2609 "link target {target_id} not found"
2610 )));
2611 }
2612 Err(error) => return Err(error),
2613 },
2614 };
2615 return Ok((EdgeEndpointKind::Note, target_kind));
2616 } else if crate::pack::is_special_relation(relation) {
2617 let rel_name = relation.as_str();
2621 let src = match self.resolve_edge_endpoint(token, source_id).await? {
2622 Some(r) => r,
2623 None => {
2624 if self.get_edge(token, source_id).await?.is_some() {
2625 return Err(RuntimeError::InvalidInput(format!(
2626 "{rel_name} source {source_id} must be a note or entity (got edge)"
2627 )));
2628 }
2629 return Err(RuntimeError::NotFound(format!(
2630 "link source {source_id} not found"
2631 )));
2632 }
2633 };
2634 let tgt = match self.resolve_edge_endpoint(token, target_id).await? {
2635 Some(r) => r,
2636 None => {
2637 if self.get_edge(token, target_id).await?.is_some() {
2638 return Err(RuntimeError::InvalidInput(format!(
2639 "{rel_name} target {target_id} must be a note or entity (got edge)"
2640 )));
2641 }
2642 return Err(RuntimeError::NotFound(format!(
2643 "link target {target_id} not found"
2644 )));
2645 }
2646 };
2647 return match (&src, &tgt) {
2648 (Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
2649 if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
2650 let legal_relations = accepted_entity_relations_description(
2651 &self.pack_edge_rules(),
2652 &src_e.kind,
2653 src_e.entity_type.as_deref(),
2654 &tgt_e.kind,
2655 tgt_e.entity_type.as_deref(),
2656 );
2657 let rule_hint = match relation {
2658 EdgeRelation::Supports | EdgeRelation::Refutes => {
2659 "requires concept|document|dataset|artifact -> concept \
2660 (or same-substrate note -> note)"
2661 }
2662 _ => "requires same-kind entity endpoints",
2663 };
2664 return Err(RuntimeError::InvalidInput(format!(
2665 "({}) -[{rel_name}]-> ({}) is not in the base endpoint \
2666 allowlist; {rel_name} {rule_hint}; currently legal relations for \
2667 {} -> {} under the loaded endpoint rules: {legal_relations}",
2668 src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
2669 )));
2670 }
2671 Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
2672 }
2673 (Resolved::Note(_), Resolved::Note(_)) => {
2674 Ok((EdgeEndpointKind::Note, EdgeEndpointKind::Note))
2675 }
2676 (Resolved::Event(_), _) => {
2677 return Err(RuntimeError::InvalidInput(format!(
2678 "{rel_name} does not apply to events; source {source_id} is an event"
2679 )));
2680 }
2681 (_, Resolved::Event(_)) => {
2682 return Err(RuntimeError::InvalidInput(format!(
2683 "{rel_name} does not apply to events; target {target_id} is an event"
2684 )));
2685 }
2686 (Resolved::Entity(_), Resolved::Note(_)) => {
2687 return Err(RuntimeError::InvalidInput(format!(
2688 "{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
2689 got source={source_id} (entity) target={target_id} (note)"
2690 )));
2691 }
2692 (Resolved::Note(_), Resolved::Entity(_)) => {
2693 return Err(RuntimeError::InvalidInput(format!(
2694 "{rel_name} endpoints must be the same substrate (note→note or entity→entity); \
2695 got source={source_id} (note) target={target_id} (entity)"
2696 )));
2697 }
2698 (Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
2699 return Err(RuntimeError::InvalidInput(format!(
2700 "pack-private record is not a valid edge endpoint for {rel_name}"
2701 )));
2702 }
2703 };
2704 } else {
2705 let src_res = self.resolve_edge_endpoint(token, source_id).await?;
2712 let tgt_res = self.resolve_edge_endpoint(token, target_id).await?;
2713 let pack_rules = self.pack_edge_rules();
2714
2715 if pack_rule_allows(&pack_rules, relation, src_res.as_ref(), tgt_res.as_ref()) {
2716 let kind = |resolved: Option<&Resolved>| match resolved {
2717 Some(Resolved::Entity(_)) => Some(EdgeEndpointKind::Entity),
2718 Some(Resolved::Note(_)) => Some(EdgeEndpointKind::Note),
2719 _ => None,
2720 };
2721 return match (kind(src_res.as_ref()), kind(tgt_res.as_ref())) {
2722 (Some(source_kind), Some(target_kind)) => Ok((source_kind, target_kind)),
2723 _ => Err(RuntimeError::Internal(
2724 "pack endpoint rule admitted an unsupported substrate".into(),
2725 )),
2726 };
2727 }
2728
2729 let (src_kind, src_entity_type) = match src_res.as_ref() {
2731 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2732 Some(_) => {
2733 return Err(RuntimeError::InvalidInput(format!(
2734 "link source {source_id} must be an entity for relation {relation:?} \
2735 (only `annotates` crosses substrates)"
2736 )));
2737 }
2738 None => {
2739 if self.get_edge(token, source_id).await?.is_some() {
2740 return Err(RuntimeError::InvalidInput(format!(
2741 "link source {source_id} must be an entity for relation {relation:?} \
2742 (only `annotates` crosses substrates)"
2743 )));
2744 }
2745 return Err(RuntimeError::NotFound(format!(
2746 "link source {source_id} not found"
2747 )));
2748 }
2749 };
2750 let (tgt_kind, tgt_entity_type) = match tgt_res.as_ref() {
2751 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2752 Some(_) => {
2753 return Err(RuntimeError::InvalidInput(format!(
2754 "link target {target_id} must be an entity for relation {relation:?} \
2755 (only `annotates` crosses substrates)"
2756 )));
2757 }
2758 None => {
2759 if self.get_edge(token, target_id).await?.is_some() {
2760 return Err(RuntimeError::InvalidInput(format!(
2761 "link target {target_id} must be an entity for relation {relation:?} \
2762 (only `annotates` crosses substrates)"
2763 )));
2764 }
2765 return Err(RuntimeError::NotFound(format!(
2766 "link target {target_id} not found"
2767 )));
2768 }
2769 };
2770 if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
2771 let legal_relations = accepted_entity_relations_description(
2772 &pack_rules,
2773 src_kind,
2774 src_entity_type,
2775 tgt_kind,
2776 tgt_entity_type,
2777 );
2778 return Err(RuntimeError::InvalidInput(format!(
2779 "({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
2780 allowlist; use pack EDGE_RULES to extend the allowlist; currently legal \
2781 relations for {src_kind} -> {tgt_kind} under the loaded endpoint rules: \
2782 {legal_relations}",
2783 relation.as_str()
2784 )));
2785 }
2786 }
2787 Ok((EdgeEndpointKind::Entity, EdgeEndpointKind::Entity))
2788 }
2789
2790 pub async fn validate_link_endpoints(
2795 &self,
2796 token: &NamespaceToken,
2797 source_id: Uuid,
2798 target_id: Uuid,
2799 relation: EdgeRelation,
2800 ) -> RuntimeResult<()> {
2801 self.validate_edge_relation_endpoints(token, source_id, target_id, relation)
2802 .await
2803 .map(|_| ())
2804 }
2805
2806 pub fn validate_link_endpoints_by_resolved(
2817 &self,
2818 source_id: Uuid,
2819 target_id: Uuid,
2820 relation: EdgeRelation,
2821 src: Option<&Resolved>,
2822 tgt: Option<&Resolved>,
2823 ) -> RuntimeResult<()> {
2824 if source_id == target_id {
2825 return Err(RuntimeError::InvalidInput(
2826 "self-loop edges are not allowed: source_id and target_id must be different".into(),
2827 ));
2828 }
2829
2830 if relation == EdgeRelation::Annotates {
2831 match src {
2832 Some(Resolved::Note(_)) => {}
2833 Some(_) => {
2834 return Err(RuntimeError::InvalidInput(format!(
2835 "annotates source {source_id} must be a note"
2836 )));
2837 }
2838 None => {
2839 return Err(RuntimeError::NotFound(format!(
2840 "link source {source_id} not found"
2841 )));
2842 }
2843 }
2844 if tgt.is_none() {
2845 return Err(RuntimeError::NotFound(format!(
2846 "link target {target_id} not found"
2847 )));
2848 }
2849 return Ok(());
2850 }
2851
2852 if crate::pack::is_special_relation(relation) {
2853 let rel_name = relation.as_str();
2854 let src = src.ok_or_else(|| {
2855 RuntimeError::NotFound(format!("link source {source_id} not found"))
2856 })?;
2857 let tgt = tgt.ok_or_else(|| {
2858 RuntimeError::NotFound(format!("link target {target_id} not found"))
2859 })?;
2860 match (src, tgt) {
2861 (Resolved::Entity(src_e), Resolved::Entity(tgt_e)) => {
2862 if !base_entity_rule_allows(&src_e.kind, relation, &tgt_e.kind) {
2863 let legal_relations = accepted_entity_relations_description(
2864 &self.pack_edge_rules(),
2865 &src_e.kind,
2866 src_e.entity_type.as_deref(),
2867 &tgt_e.kind,
2868 tgt_e.entity_type.as_deref(),
2869 );
2870 let rule_hint = match relation {
2871 EdgeRelation::Supports | EdgeRelation::Refutes => {
2872 "requires concept|document|dataset|artifact -> concept \
2873 (or same-substrate note -> note)"
2874 }
2875 _ => "requires same-kind entity endpoints",
2876 };
2877 return Err(RuntimeError::InvalidInput(format!(
2878 "({}) -[{rel_name}]-> ({}) is not in the base endpoint \
2879 allowlist; {rel_name} {rule_hint}; currently legal relations for \
2880 {} -> {} under the loaded endpoint rules: {legal_relations}",
2881 src_e.kind, tgt_e.kind, src_e.kind, tgt_e.kind
2882 )));
2883 }
2884 }
2885 (Resolved::Note(_), Resolved::Note(_)) => {}
2886 (Resolved::Entity(_), Resolved::Note(_)) => {
2887 return Err(RuntimeError::InvalidInput(format!(
2888 "{rel_name} endpoints must be the same substrate \
2889 (note→note or entity→entity); got source={source_id} (entity) \
2890 target={target_id} (note)"
2891 )));
2892 }
2893 (Resolved::Note(_), Resolved::Entity(_)) => {
2894 return Err(RuntimeError::InvalidInput(format!(
2895 "{rel_name} endpoints must be the same substrate \
2896 (note→note or entity→entity); got source={source_id} (note) \
2897 target={target_id} (entity)"
2898 )));
2899 }
2900 (Resolved::PackRecord { .. }, _) | (_, Resolved::PackRecord { .. }) => {
2901 return Err(RuntimeError::InvalidInput(format!(
2902 "pack-private record is not a valid edge endpoint for {rel_name}"
2903 )));
2904 }
2905 _ => {
2906 return Err(RuntimeError::InvalidInput(format!(
2907 "{rel_name} endpoints must be notes or entities (not events)"
2908 )));
2909 }
2910 }
2911 return Ok(());
2912 }
2913
2914 let pack_rules = self.pack_edge_rules();
2917 if pack_rule_allows(&pack_rules, relation, src, tgt) {
2918 return Ok(());
2919 }
2920
2921 let (src_kind, src_entity_type) = match src {
2922 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2923 Some(_) => {
2924 return Err(RuntimeError::InvalidInput(format!(
2925 "link source {source_id} must be an entity for relation {relation:?} \
2926 (only `annotates` crosses substrates)"
2927 )));
2928 }
2929 None => {
2930 return Err(RuntimeError::NotFound(format!(
2931 "link source {source_id} not found"
2932 )));
2933 }
2934 };
2935 let (tgt_kind, tgt_entity_type) = match tgt {
2936 Some(Resolved::Entity(e)) => (e.kind.as_str(), e.entity_type.as_deref()),
2937 Some(_) => {
2938 return Err(RuntimeError::InvalidInput(format!(
2939 "link target {target_id} must be an entity for relation {relation:?} \
2940 (only `annotates` crosses substrates)"
2941 )));
2942 }
2943 None => {
2944 return Err(RuntimeError::NotFound(format!(
2945 "link target {target_id} not found"
2946 )));
2947 }
2948 };
2949
2950 if !base_entity_rule_allows(src_kind, relation, tgt_kind) {
2951 let legal_relations = accepted_entity_relations_description(
2952 &pack_rules,
2953 src_kind,
2954 src_entity_type,
2955 tgt_kind,
2956 tgt_entity_type,
2957 );
2958 return Err(RuntimeError::InvalidInput(format!(
2959 "({src_kind}) -[{}]-> ({tgt_kind}) is not in the base endpoint \
2960 allowlist; use pack EDGE_RULES to extend the allowlist; currently legal relations \
2961 for {src_kind} -> {tgt_kind} under the loaded endpoint rules: {legal_relations}",
2962 relation.as_str()
2963 )));
2964 }
2965
2966 Ok(())
2967 }
2968
2969 pub fn validate_annotates_endpoint_kinds(
2981 &self,
2982 source_id: Uuid,
2983 target_id: Uuid,
2984 source: Option<EdgeEndpointKind>,
2985 target: Option<EdgeEndpointKind>,
2986 ) -> RuntimeResult<()> {
2987 if source_id == target_id {
2988 return Err(RuntimeError::InvalidInput(
2989 "self-loop edges are not allowed: source_id and target_id must be different".into(),
2990 ));
2991 }
2992 match source {
2993 Some(EdgeEndpointKind::Note) => {}
2994 Some(_) => {
2995 return Err(RuntimeError::InvalidInput(format!(
2996 "annotates source {source_id} must be a note"
2997 )));
2998 }
2999 None => {
3000 return Err(RuntimeError::NotFound(format!(
3001 "link source {source_id} not found"
3002 )));
3003 }
3004 }
3005 if target.is_none() {
3006 return Err(RuntimeError::NotFound(format!(
3007 "link target {target_id} not found"
3008 )));
3009 }
3010 Ok(())
3011 }
3012
3013 pub async fn link(
3033 &self,
3034 token: &NamespaceToken,
3035 source_id: Uuid,
3036 target_id: Uuid,
3037 relation: EdgeRelation,
3038 weight: f64,
3039 metadata: Option<serde_json::Value>,
3040 ) -> RuntimeResult<Edge> {
3041 self.link_observed(
3042 token, source_id, target_id, relation, weight, metadata, false,
3043 )
3044 .await
3045 .map(|result| result.edge)
3046 }
3047
3048 #[allow(clippy::too_many_arguments)]
3053 pub async fn link_observed(
3054 &self,
3055 token: &NamespaceToken,
3056 source_id: Uuid,
3057 target_id: Uuid,
3058 relation: EdgeRelation,
3059 weight: f64,
3060 metadata: Option<serde_json::Value>,
3061 resurrect: bool,
3062 ) -> RuntimeResult<EdgeUpsertResult> {
3063 validate_edge_weight(weight)?;
3064 validate_edge_metadata(relation, metadata.as_ref())?;
3065 let (source_kind, target_kind) = self
3066 .validate_edge_relation_endpoints(token, source_id, target_id, relation)
3067 .await?;
3068 let (canonical_source, canonical_target) =
3069 canonical_edge_endpoints(relation, source_id, target_id);
3070 let (source_kind, target_kind) =
3071 canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
3072 let (source_id, target_id) = (canonical_source, canonical_target);
3073 let metadata = if relation == EdgeRelation::DependsOn {
3074 match (
3079 self.resolve_edge_endpoint(token, source_id).await?,
3080 self.resolve_edge_endpoint(token, target_id).await?,
3081 ) {
3082 (Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
3083 merge_dependency_kind(&src_e.kind, &tgt_e.kind, metadata)
3084 }
3085 _ => metadata,
3086 }
3087 } else {
3088 metadata
3089 };
3090 validate_edge_metadata(relation, metadata.as_ref())?;
3091 let now = chrono::Utc::now();
3092 let ns = token.namespace().as_str();
3093 let edge = Edge {
3094 id: LinkId::from(Uuid::new_v4()),
3095 namespace: ns.to_string(),
3096 source_id,
3097 target_id,
3098 relation,
3099 weight,
3100 created_at: now,
3101 updated_at: now,
3102 deleted_at: None,
3103 metadata,
3104 target_backend: None,
3105 };
3106 let attribution = crate::EventAttribution::from_token(token);
3116 let outcome = compose_graph_mutation_events(
3117 self.backend(),
3118 GraphMutationRequest::Single {
3119 request: EdgeUpsertRequest { edge, resurrect },
3120 guard_endpoints: true,
3121 },
3122 GraphMutationPreconditions::default(),
3123 Vec::new(),
3124 move |outcome| match &outcome.mutation {
3125 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
3126 Ok(vec![Self::link_mutation_event(
3127 &attribution,
3128 result,
3129 source_kind,
3130 target_kind,
3131 )])
3132 }
3133 _ => Err(Self::link_composition_shape_error(
3134 "expected a written singleton",
3135 )),
3136 },
3137 )
3138 .await?;
3139 let GraphMutationOutcome::Single(outcome) = outcome.mutation else {
3140 return Err(RuntimeError::Internal(
3141 "link: unexpected composition outcome".into(),
3142 ));
3143 };
3144 let result = match outcome {
3145 GuardedEdgeUpsertOutcome::Written(result) => result,
3146 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::MissingEndpoints(missing)) => {
3147 return Err(RuntimeError::GuardedWriteFailed(GuardedWriteFailure {
3148 entry_index: None,
3149 missing_source: missing.source.then_some(source_id),
3150 missing_target: missing.target.then_some(target_id),
3151 }));
3152 }
3153 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::ResurrectionRequired { edge }) => {
3154 return Err(RuntimeError::InvalidInput(format!(
3155 "edge {} is soft-deleted; pass resurrect=true to link explicitly",
3156 edge.id
3157 )))
3158 }
3159 };
3160 Ok(result)
3161 }
3162
3163 #[allow(clippy::too_many_arguments)]
3171 pub async fn link_with_target_backend(
3172 &self,
3173 token: &NamespaceToken,
3174 source_id: Uuid,
3175 target_id: Uuid,
3176 source_kind: EdgeEndpointKind,
3177 target_kind: EdgeEndpointKind,
3178 relation: EdgeRelation,
3179 weight: f64,
3180 metadata: Option<serde_json::Value>,
3181 target_backend: Option<String>,
3182 ) -> RuntimeResult<Edge> {
3183 self.link_with_target_backend_observed(
3184 token,
3185 source_id,
3186 target_id,
3187 source_kind,
3188 target_kind,
3189 relation,
3190 weight,
3191 metadata,
3192 target_backend,
3193 false,
3194 )
3195 .await
3196 .map(|result| result.edge)
3197 }
3198
3199 #[allow(clippy::too_many_arguments)]
3203 pub async fn link_with_target_backend_observed(
3204 &self,
3205 token: &NamespaceToken,
3206 source_id: Uuid,
3207 target_id: Uuid,
3208 source_kind: EdgeEndpointKind,
3209 target_kind: EdgeEndpointKind,
3210 relation: EdgeRelation,
3211 weight: f64,
3212 metadata: Option<serde_json::Value>,
3213 target_backend: Option<String>,
3214 resurrect: bool,
3215 ) -> RuntimeResult<EdgeUpsertResult> {
3216 validate_edge_weight(weight)?;
3217 let (canonical_source, canonical_target) =
3218 canonical_edge_endpoints(relation, source_id, target_id);
3219 let (source_kind, target_kind) =
3220 canonical_edge_endpoint_kinds(source_id, canonical_source, source_kind, target_kind);
3221 let (source_id, target_id) = (canonical_source, canonical_target);
3222 validate_edge_metadata(relation, metadata.as_ref())?;
3223 let now = chrono::Utc::now();
3224 let ns = token.namespace().as_str();
3225 let edge = Edge {
3226 id: LinkId::from(Uuid::new_v4()),
3227 namespace: ns.to_string(),
3228 source_id,
3229 target_id,
3230 relation,
3231 weight,
3232 created_at: now,
3233 updated_at: now,
3234 deleted_at: None,
3235 metadata,
3236 target_backend,
3237 };
3238 let attribution = crate::EventAttribution::from_token(token);
3239 let outcome = compose_graph_mutation_events(
3240 self.backend(),
3241 GraphMutationRequest::Single {
3242 request: EdgeUpsertRequest { edge, resurrect },
3243 guard_endpoints: false,
3244 },
3245 GraphMutationPreconditions::default(),
3246 Vec::new(),
3247 move |outcome| match &outcome.mutation {
3248 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => {
3249 Ok(vec![Self::link_mutation_event(
3250 &attribution,
3251 result,
3252 source_kind,
3253 target_kind,
3254 )])
3255 }
3256 _ => Err(Self::link_composition_shape_error(
3257 "expected a written singleton",
3258 )),
3259 },
3260 )
3261 .await?;
3262 match outcome.mutation {
3263 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(result)) => Ok(result),
3264 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Refused(
3265 EdgeUpsertRefusal::ResurrectionRequired { edge },
3266 )) => {
3267 let error = khive_storage::StorageError::Conflict {
3268 capability: khive_storage::StorageCapability::Graph,
3269 operation: "upsert_edge_observed".into(),
3270 message: format!(
3271 "edge {} is soft-deleted; explicit resurrection is required",
3272 edge.id,
3273 ),
3274 };
3275 Err(RuntimeError::InvalidInput(format!(
3276 "edge natural key is soft-deleted; pass resurrect=true to link explicitly: {error}"
3277 )))
3278 }
3279 _ => Err(RuntimeError::Internal(
3280 "link: unexpected composition outcome".into(),
3281 )),
3282 }
3283 }
3284
3285 fn link_mutation_event(
3286 attribution: &crate::EventAttribution,
3287 result: &EdgeUpsertResult,
3288 source_kind: EdgeEndpointKind,
3289 target_kind: EdgeEndpointKind,
3290 ) -> Event {
3291 let kind = match result.disposition {
3292 EdgeUpsertDisposition::Created => EventKind::LinkCreated,
3293 EdgeUpsertDisposition::Updated | EdgeUpsertDisposition::Resurrected => {
3294 EventKind::EdgeUpdated
3295 }
3296 };
3297 let edge_id = Uuid::from(result.edge.id);
3298 let mut payload = serde_json::json!({
3299 "id": edge_id,
3300 "namespace": result.edge.namespace,
3301 "mutation": result.disposition.name(),
3302 "source_id": result.edge.source_id,
3303 "target_id": result.edge.target_id,
3304 "relation": result.edge.relation,
3305 "weight": result.edge.weight,
3306 "metadata": result.edge.metadata,
3307 "previous": result.previous,
3308 });
3309 if kind == EventKind::LinkCreated {
3310 payload["source_kind"] = serde_json::json!(source_kind.name());
3311 payload["target_kind"] = serde_json::json!(target_kind.name());
3312 }
3313 attribution.stamp(
3314 Event::new(
3315 result.edge.namespace.clone(),
3316 "link",
3317 kind,
3318 SubstrateKind::Entity,
3319 "",
3320 )
3321 .with_target(edge_id)
3322 .with_payload(payload),
3323 )
3324 }
3325
3326 fn link_composition_shape_error(message: &'static str) -> khive_storage::StorageError {
3327 khive_storage::StorageError::InvalidInput {
3328 capability: khive_storage::StorageCapability::Graph,
3329 operation: "link_mutation_event".into(),
3330 message: message.into(),
3331 }
3332 }
3333
3334 pub(crate) async fn substrate_exists_in_ns(
3341 &self,
3342 token: &NamespaceToken,
3343 id: Uuid,
3344 ) -> RuntimeResult<bool> {
3345 if self.resolve(token, id).await?.is_some() {
3346 return Ok(true);
3347 }
3348 match self.get_edge_visible(token, id).await {
3349 Ok(Some(_)) => Ok(true),
3350 Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
3351 Err(err) => Err(err),
3352 }
3353 }
3354
3355 pub(crate) async fn substrate_exists_by_id(
3362 &self,
3363 token: &NamespaceToken,
3364 id: Uuid,
3365 ) -> RuntimeResult<bool> {
3366 if self.resolve_edge_endpoint(token, id).await?.is_some() {
3367 return Ok(true);
3368 }
3369 match self.get_edge(token, id).await {
3370 Ok(Some(_)) => Ok(true),
3371 Ok(None) | Err(RuntimeError::NotFound(_)) => Ok(false),
3372 Err(err) => Err(err),
3373 }
3374 }
3375
3376 pub async fn latest_annotating_note(
3381 &self,
3382 token: &NamespaceToken,
3383 node_id: Uuid,
3384 kind: &str,
3385 tag: &str,
3386 ) -> RuntimeResult<Option<Uuid>> {
3387 self.latest_annotating_note_inner(token, node_id, kind, tag, None)
3388 .await
3389 }
3390
3391 pub async fn latest_annotating_note_with_property(
3394 &self,
3395 token: &NamespaceToken,
3396 node_id: Uuid,
3397 kind: &str,
3398 tag: &str,
3399 property_key: &str,
3400 property_value: &str,
3401 ) -> RuntimeResult<Option<Uuid>> {
3402 self.latest_annotating_note_inner(
3403 token,
3404 node_id,
3405 kind,
3406 tag,
3407 Some((property_key, property_value)),
3408 )
3409 .await
3410 }
3411
3412 async fn latest_annotating_note_inner(
3413 &self,
3414 token: &NamespaceToken,
3415 node_id: Uuid,
3416 kind: &str,
3417 tag: &str,
3418 required_property: Option<(&str, &str)>,
3419 ) -> RuntimeResult<Option<Uuid>> {
3420 if !self.substrate_exists_in_ns(token, node_id).await? {
3421 return Ok(None);
3422 }
3423 let mut latest: Option<(Uuid, i64)> = None;
3424 for namespace in token.visible_namespaces() {
3425 let scoped = NamespaceToken::for_namespace(namespace.clone());
3426 let graph = self.graph(&scoped)?;
3427 let candidate = match required_property {
3428 Some((key, value)) => {
3429 graph
3430 .latest_annotating_note_with_property(node_id, kind, tag, key, value)
3431 .await?
3432 }
3433 None => graph.latest_annotating_note(node_id, kind, tag).await?,
3434 };
3435 if let Some(candidate) = candidate {
3436 if latest.is_none_or(|(id, created_at)| {
3437 candidate.1 > created_at || (candidate.1 == created_at && candidate.0 < id)
3438 }) {
3439 latest = Some(candidate);
3440 }
3441 }
3442 }
3443 Ok(latest.map(|(id, _)| id))
3444 }
3445
3446 pub async fn neighbors(
3455 &self,
3456 token: &NamespaceToken,
3457 node_id: Uuid,
3458 direction: Direction,
3459 limit: Option<u32>,
3460 relations: Option<Vec<EdgeRelation>>,
3461 ) -> RuntimeResult<Vec<NeighborHit>> {
3462 self.neighbors_with_query(
3463 token,
3464 node_id,
3465 NeighborQuery {
3466 direction,
3467 relations,
3468 limit,
3469 min_weight: None,
3470 },
3471 )
3472 .await
3473 }
3474
3475 pub async fn neighbors_with_query(
3485 &self,
3486 token: &NamespaceToken,
3487 node_id: Uuid,
3488 query: NeighborQuery,
3489 ) -> RuntimeResult<Vec<NeighborHit>> {
3490 self.neighbors_with_query_page(token, node_id, query, None, None, true)
3491 .await
3492 }
3493
3494 pub async fn neighbors_with_query_page(
3498 &self,
3499 token: &NamespaceToken,
3500 node_id: Uuid,
3501 query: NeighborQuery,
3502 after: Option<NeighborCursor>,
3503 neighbor_kinds: Option<Vec<String>>,
3504 enrich: bool,
3505 ) -> RuntimeResult<Vec<NeighborHit>> {
3506 if !self.substrate_exists_by_id(token, node_id).await? {
3509 return Err(RuntimeError::NotFound(format!(
3510 "neighbor anchor {node_id} not found"
3511 )));
3512 }
3513
3514 self.neighbors_for_resolved_kg_read(
3515 token,
3516 node_id,
3517 crate::KgNeighborRead {
3518 query,
3519 after,
3520 neighbor_kinds,
3521 enrich,
3522 namespace: None,
3523 },
3524 )
3525 .await
3526 }
3527
3528 pub async fn neighbors_for_resolved_kg_read(
3536 &self,
3537 token: &NamespaceToken,
3538 node_id: Uuid,
3539 options: crate::KgNeighborRead,
3540 ) -> RuntimeResult<Vec<NeighborHit>> {
3541 self.neighbors_for_resolved_kg_read_inner(token, node_id, options, false)
3542 .await
3543 .map(|(hits, _)| hits)
3544 }
3545
3546 pub async fn neighbors_for_resolved_kg_read_with_entity_kinds(
3551 &self,
3552 token: &NamespaceToken,
3553 node_id: Uuid,
3554 options: crate::KgNeighborRead,
3555 ) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
3556 self.neighbors_for_resolved_kg_read_inner(token, node_id, options, true)
3557 .await
3558 }
3559
3560 async fn neighbors_for_resolved_kg_read_inner(
3561 &self,
3562 token: &NamespaceToken,
3563 node_id: Uuid,
3564 options: crate::KgNeighborRead,
3565 with_entity_kinds: bool,
3566 ) -> RuntimeResult<(Vec<NeighborHit>, HashMap<Uuid, String>)> {
3567 let crate::KgNeighborRead {
3568 mut query,
3569 after,
3570 neighbor_kinds,
3571 enrich,
3572 namespace,
3573 } = options;
3574 let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace.as_ref())?;
3575 query.direction =
3576 normalize_symmetric_direction(query.direction, query.relations.as_deref());
3577 let mut hits = Vec::new();
3578 for ns in namespaces {
3579 let temp = NamespaceToken::for_namespace(ns.clone());
3580 let mut ns_hits = self
3581 .graph(&temp)?
3582 .neighbors_page(node_id, query.clone(), after, neighbor_kinds.clone())
3583 .await?;
3584 hits.append(&mut ns_hits);
3585 }
3586 hits.sort_by_key(|h| (h.node_id, h.edge_id));
3587 hits.dedup_by_key(|h| (h.node_id, h.edge_id));
3588 if enrich {
3589 self.enrich_neighbor_hits(token, &mut hits).await;
3590 }
3591 let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.node_id).collect();
3593 let (deleted, entity_kinds) = self
3594 .neighbor_node_screen(
3595 candidate_ids,
3596 (with_entity_kinds && !enrich).then_some(token),
3597 )
3598 .await?;
3599 if !deleted.is_empty() {
3600 hits.retain(|h| !deleted.contains(&h.node_id));
3601 }
3602 hits.sort_by(|a, b| {
3609 b.weight
3610 .partial_cmp(&a.weight)
3611 .unwrap_or(std::cmp::Ordering::Equal)
3612 .then(a.node_id.cmp(&b.node_id))
3613 .then(a.edge_id.cmp(&b.edge_id))
3614 });
3615 Ok((hits, entity_kinds))
3616 }
3617
3618 pub async fn annotation_neighbors_by_target_id(
3626 &self,
3627 target_id: Uuid,
3628 ) -> RuntimeResult<Vec<NeighborHit>> {
3629 let mut reader = self.sql().reader().await?;
3630 let rows = reader
3631 .query_all(SqlStatement {
3632 sql: "SELECT source_id, id, weight FROM graph_edges \
3633 WHERE target_id = ?1 AND relation = ?2 AND deleted_at IS NULL \
3634 ORDER BY weight DESC, source_id ASC"
3635 .to_string(),
3636 params: vec![
3637 SqlValue::Text(target_id.to_string()),
3638 SqlValue::Text(EdgeRelation::Annotates.to_string()),
3639 ],
3640 label: Some("annotations.by_target_id_unfiltered".into()),
3641 })
3642 .await?;
3643
3644 rows.into_iter()
3645 .map(|row| {
3646 let parse_uuid = |name: &str| match row.get(name) {
3647 Some(SqlValue::Text(value)) => Uuid::from_str(value).map_err(|error| {
3648 RuntimeError::Internal(format!("graph_edges.{name} is not a UUID: {error}"))
3649 }),
3650 Some(value) => Err(RuntimeError::Internal(format!(
3651 "graph_edges.{name} has unexpected SQL value {value:?}"
3652 ))),
3653 None => Err(RuntimeError::Internal(format!(
3654 "graph_edges row missing {name}"
3655 ))),
3656 };
3657 let weight = match row.get("weight") {
3658 Some(SqlValue::Float(value)) => Ok(*value),
3659 Some(value) => Err(RuntimeError::Internal(format!(
3660 "graph_edges.weight has unexpected SQL value {value:?}"
3661 ))),
3662 None => Err(RuntimeError::Internal(
3663 "graph_edges row missing weight".into(),
3664 )),
3665 }?;
3666
3667 Ok(NeighborHit {
3668 node_id: parse_uuid("source_id")?,
3669 edge_id: parse_uuid("id")?,
3670 relation: EdgeRelation::Annotates,
3671 weight,
3672 name: None,
3673 kind: None,
3674 entity_type: None,
3675 })
3676 })
3677 .collect()
3678 }
3679
3680 pub async fn neighbors_with_query_directed(
3689 &self,
3690 token: &NamespaceToken,
3691 node_id: Uuid,
3692 query: NeighborQuery,
3693 ) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
3694 if !self.substrate_exists_by_id(token, node_id).await? {
3695 return Err(RuntimeError::NotFound(format!(
3696 "neighbor anchor {node_id} not found"
3697 )));
3698 }
3699
3700 self.directed_neighbors_for_resolved_kg_read(token, node_id, query, None)
3701 .await
3702 }
3703
3704 pub(crate) async fn directed_neighbors_for_resolved_kg_read(
3705 &self,
3706 token: &NamespaceToken,
3707 node_id: Uuid,
3708 query: NeighborQuery,
3709 namespace: Option<&crate::Namespace>,
3710 ) -> RuntimeResult<Vec<(NeighborHit, Direction)>> {
3711 let namespaces = crate::kg_read::neighbor_read_namespaces(token, namespace)?;
3712 let mut hits: Vec<DirectedNeighborHit> = Vec::new();
3713 for ns in namespaces {
3714 let temp = NamespaceToken::for_namespace(ns.clone());
3715 let mut ns_hits = self
3716 .graph(&temp)?
3717 .neighbors_both_directions(node_id, query.clone())
3718 .await?;
3719 hits.append(&mut ns_hits);
3720 }
3721 hits.sort_by_key(|h| {
3725 (
3726 h.hit.node_id,
3727 h.hit.edge_id,
3728 direction_sort_rank(&h.direction),
3729 )
3730 });
3731 hits.dedup_by_key(|h| {
3732 (
3733 h.hit.node_id,
3734 h.hit.edge_id,
3735 direction_sort_rank(&h.direction),
3736 )
3737 });
3738
3739 let mut plain_hits: Vec<NeighborHit> = hits.iter().map(|h| h.hit.clone()).collect();
3740 self.enrich_neighbor_hits(token, &mut plain_hits).await;
3741 for (dh, enriched) in hits.iter_mut().zip(plain_hits) {
3742 dh.hit = enriched;
3743 }
3744
3745 let candidate_ids: Vec<Uuid> = hits.iter().map(|h| h.hit.node_id).collect();
3747 let deleted = self.deleted_entity_ids(candidate_ids).await?;
3748 if !deleted.is_empty() {
3749 hits.retain(|h| !deleted.contains(&h.hit.node_id));
3750 }
3751 hits.sort_by(|a, b| {
3755 b.hit
3756 .weight
3757 .partial_cmp(&a.hit.weight)
3758 .unwrap_or(std::cmp::Ordering::Equal)
3759 .then(a.hit.node_id.cmp(&b.hit.node_id))
3760 .then(a.hit.edge_id.cmp(&b.hit.edge_id))
3761 });
3762 Ok(hits.into_iter().map(|h| (h.hit, h.direction)).collect())
3763 }
3764
3765 pub async fn traverse(
3771 &self,
3772 token: &NamespaceToken,
3773 request: TraversalRequest,
3774 ) -> RuntimeResult<Vec<GraphPath>> {
3775 let mut request = request;
3776 request.validate().map_err(RuntimeError::InvalidInput)?;
3777 let mut roots = Vec::with_capacity(request.roots.len());
3778 let mut seen_roots = std::collections::HashSet::with_capacity(request.roots.len());
3779 for root in request.roots.drain(..) {
3780 if seen_roots.insert(root) {
3781 if !self.substrate_exists_by_id(token, root).await? {
3782 return Err(RuntimeError::NotFound(format!(
3783 "traverse root {root} not found"
3784 )));
3785 }
3786 roots.push(root);
3787 }
3788 }
3789 request.roots = roots;
3790 if request.roots.is_empty() {
3791 return Ok(Vec::new());
3792 }
3793
3794 let mut paths = Vec::new();
3795 for ns in token.visible_namespaces() {
3796 let temp = NamespaceToken::for_namespace(ns.clone());
3797 let mut ns_paths = self.graph(&temp)?.traverse(request.clone()).await?;
3798 paths.append(&mut ns_paths);
3799 }
3800 let mut paths =
3804 merge_traversal_paths_by_root(paths, Some(request.options.effective_limit()));
3805 self.enrich_path_nodes(token, &mut paths, request.include_properties)
3806 .await;
3807 let all_node_ids: Vec<Uuid> = paths
3809 .iter()
3810 .flat_map(|p| p.nodes.iter().map(|n| n.node_id))
3811 .collect();
3812 let deleted = self.deleted_entity_ids(all_node_ids).await?;
3813 if !deleted.is_empty() {
3814 for path in paths.iter_mut() {
3815 path.nodes.retain(|n| !deleted.contains(&n.node_id));
3816 recompute_total_weight(path);
3817 }
3818 paths.retain(|p| !p.nodes.is_empty());
3819 }
3820 Ok(paths)
3821 }
3822
3823 async fn deleted_entity_ids(
3843 &self,
3844 ids: Vec<Uuid>,
3845 ) -> RuntimeResult<std::collections::HashSet<Uuid>> {
3846 self.neighbor_node_screen(ids, None)
3847 .await
3848 .map(|(deleted, _)| deleted)
3849 }
3850
3851 async fn neighbor_node_screen(
3855 &self,
3856 ids: Vec<Uuid>,
3857 kind_token: Option<&NamespaceToken>,
3858 ) -> RuntimeResult<(std::collections::HashSet<Uuid>, HashMap<Uuid, String>)> {
3859 if ids.is_empty() {
3860 return Ok((std::collections::HashSet::new(), HashMap::new()));
3861 }
3862 let id_strs: Vec<String> = ids.iter().map(|u| u.to_string()).collect();
3863 let n = id_strs.len();
3864 let entities_placeholders = (0..n)
3871 .map(|i| format!("?{}", i + 1))
3872 .collect::<Vec<_>>()
3873 .join(",");
3874 let notes_placeholders = (0..n)
3875 .map(|i| format!("?{}", n + i + 1))
3876 .collect::<Vec<_>>()
3877 .join(",");
3878 let sql_str = if kind_token.is_some() {
3879 format!(
3880 "SELECT id, kind, namespace, deleted_at IS NOT NULL AS is_deleted \
3881 FROM entities WHERE id IN ({entities_placeholders}) \
3882 UNION ALL \
3883 SELECT id, NULL, NULL, 1 FROM notes \
3884 WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
3885 )
3886 } else {
3887 format!(
3888 "SELECT id FROM entities WHERE id IN ({entities_placeholders}) AND deleted_at IS NOT NULL \
3889 UNION \
3890 SELECT id FROM notes WHERE id IN ({notes_placeholders}) AND deleted_at IS NOT NULL"
3891 )
3892 };
3893 let params: Vec<SqlValue> = id_strs
3895 .iter()
3896 .chain(id_strs.iter())
3897 .cloned()
3898 .map(SqlValue::Text)
3899 .collect();
3900 let stmt = SqlStatement {
3901 sql: sql_str,
3902 params,
3903 label: Some("deleted_entity_ids".into()),
3904 };
3905 let mut out = std::collections::HashSet::new();
3906 let mut entity_kinds = HashMap::new();
3907 let sql = self.sql();
3908 let mut reader = sql.reader().await?;
3909 let rows = reader.query_all(stmt).await?;
3910 for row in rows {
3911 if let Some(col) = row.columns.first() {
3912 if let SqlValue::Text(s) = &col.value {
3913 if let Ok(u) = s.parse::<Uuid>() {
3914 if kind_token.is_none()
3915 || matches!(
3916 row.columns.get(3).map(|col| &col.value),
3917 Some(SqlValue::Integer(1))
3918 )
3919 {
3920 out.insert(u);
3921 } else if let (
3922 Some(token),
3923 Some(SqlValue::Text(kind)),
3924 Some(SqlValue::Text(namespace)),
3925 Some(SqlValue::Integer(0)),
3926 ) = (
3927 kind_token,
3928 row.columns.get(1).map(|col| &col.value),
3929 row.columns.get(2).map(|col| &col.value),
3930 row.columns.get(3).map(|col| &col.value),
3931 ) {
3932 if token
3933 .visible_namespaces()
3934 .iter()
3935 .any(|ns| ns.as_str() == namespace.as_str())
3936 {
3937 entity_kinds.insert(u, kind.clone());
3938 }
3939 }
3940 }
3941 }
3942 }
3943 }
3944 Ok((out, entity_kinds))
3945 }
3946
3947 async fn enrich_neighbor_hits(&self, token: &NamespaceToken, hits: &mut [NeighborHit]) {
3956 if hits.is_empty() {
3957 return;
3958 }
3959
3960 let unique_ids: Vec<Uuid> = {
3962 let mut seen = std::collections::HashSet::new();
3963 hits.iter()
3964 .filter_map(|h| {
3965 if seen.insert(h.node_id) {
3966 Some(h.node_id)
3967 } else {
3968 None
3969 }
3970 })
3971 .collect()
3972 };
3973
3974 let entity_map: HashMap<Uuid, Entity> = self
3975 .get_entities_by_ids_visible(token, &unique_ids)
3976 .await
3977 .unwrap_or_default()
3978 .into_iter()
3979 .map(|e| (e.id, e))
3980 .collect();
3981
3982 let residual_ids: Vec<Uuid> = unique_ids
3984 .iter()
3985 .filter(|id| !entity_map.contains_key(id))
3986 .copied()
3987 .collect();
3988
3989 let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
3990 if let Ok(store) = self.notes(token) {
3991 store
3992 .get_notes_batch(&residual_ids)
3993 .await
3994 .unwrap_or_default()
3995 .into_iter()
3996 .map(|n| (n.id, n))
3997 .collect()
3998 } else {
3999 HashMap::new()
4000 }
4001 } else {
4002 HashMap::new()
4003 };
4004
4005 for hit in hits.iter_mut() {
4006 if let Some(entity) = entity_map.get(&hit.node_id) {
4007 hit.name = Some(entity.name.clone());
4008 hit.kind = Some(entity.kind.clone());
4009 hit.entity_type = entity.entity_type.clone();
4010 } else if let Some(note) = note_map.get(&hit.node_id) {
4011 hit.name = Some(note_graph_name(note));
4012 hit.kind = Some(note.kind.clone());
4013 }
4014 }
4015 }
4016
4017 async fn enrich_path_nodes(
4028 &self,
4029 token: &NamespaceToken,
4030 paths: &mut [GraphPath],
4031 include_properties: bool,
4032 ) {
4033 if paths.is_empty() {
4034 return;
4035 }
4036
4037 let unique_ids: Vec<Uuid> = {
4039 let mut seen = std::collections::HashSet::new();
4040 paths
4041 .iter()
4042 .flat_map(|p| p.nodes.iter())
4043 .filter_map(|n| {
4044 if seen.insert(n.node_id) {
4045 Some(n.node_id)
4046 } else {
4047 None
4048 }
4049 })
4050 .collect()
4051 };
4052
4053 let entity_map: HashMap<Uuid, Entity> = self
4054 .get_entities_by_ids_visible(token, &unique_ids)
4055 .await
4056 .unwrap_or_default()
4057 .into_iter()
4058 .map(|e| (e.id, e))
4059 .collect();
4060
4061 let residual_ids: Vec<Uuid> = unique_ids
4062 .iter()
4063 .filter(|id| !entity_map.contains_key(id))
4064 .copied()
4065 .collect();
4066
4067 let note_map: HashMap<Uuid, Note> = if !residual_ids.is_empty() {
4068 if let Ok(store) = self.notes(token) {
4069 store
4070 .get_notes_batch(&residual_ids)
4071 .await
4072 .unwrap_or_default()
4073 .into_iter()
4074 .map(|n| (n.id, n))
4075 .collect()
4076 } else {
4077 HashMap::new()
4078 }
4079 } else {
4080 HashMap::new()
4081 };
4082
4083 for path in paths.iter_mut() {
4084 for node in path.nodes.iter_mut() {
4085 if let Some(entity) = entity_map.get(&node.node_id) {
4086 node.name = Some(entity.name.clone());
4087 node.kind = Some(entity.kind.clone());
4088 if include_properties {
4089 node.properties = entity.properties.clone();
4090 }
4091 } else if let Some(note) = note_map.get(&node.node_id) {
4092 node.name = Some(note_graph_name(note));
4093 node.kind = Some(note.kind.clone());
4094 }
4095 }
4096 }
4097 }
4098
4099 #[allow(clippy::too_many_arguments)]
4117 pub async fn create_note(
4118 &self,
4119 token: &NamespaceToken,
4120 kind: &str,
4121 name: Option<&str>,
4122 content: &str,
4123 salience: Option<f64>,
4124 properties: Option<serde_json::Value>,
4125 annotates: Vec<Uuid>,
4126 ) -> RuntimeResult<Note> {
4127 let (note, embedding, degradations, _) = self
4128 .create_note_inner(
4129 token, kind, name, content, None, salience, None, properties, annotates, None,
4130 false, false,
4131 )
4132 .await?;
4133 legacy_post_commit_result_with_embedding(
4134 "create_note",
4135 note.id,
4136 note,
4137 embedding,
4138 degradations,
4139 )
4140 }
4141
4142 pub async fn create_web_receipt_note(
4147 &self,
4148 token: &NamespaceToken,
4149 summary: &str,
4150 request: serde_json::Value,
4151 annotates: Vec<Uuid>,
4152 ) -> RuntimeResult<Note> {
4153 let properties = serde_json::json!({
4154 "tags": ["web.receipt"],
4155 "request": request,
4156 });
4157 let (note, _, degradations, _) = self
4158 .create_note_inner(
4159 token,
4160 "observation",
4161 None,
4162 summary,
4163 None,
4164 None,
4165 None,
4166 Some(properties),
4167 annotates,
4168 None,
4169 false,
4170 true,
4171 )
4172 .await?;
4173 legacy_post_commit_result("create_web_receipt_note", note.id, note, degradations)
4174 }
4175
4176 #[allow(clippy::too_many_arguments)]
4187 pub async fn create_note_with_embedding_content(
4188 &self,
4189 token: &NamespaceToken,
4190 kind: &str,
4191 name: Option<&str>,
4192 content: &str,
4193 embedding_content: Option<&str>,
4194 salience: Option<f64>,
4195 properties: Option<serde_json::Value>,
4196 annotates: Vec<Uuid>,
4197 ) -> RuntimeResult<Note> {
4198 let (note, embedding, degradations, _) = self
4199 .create_note_inner(
4200 token,
4201 kind,
4202 name,
4203 content,
4204 embedding_content,
4205 salience,
4206 None,
4207 properties,
4208 annotates,
4209 None,
4210 false,
4211 false,
4212 )
4213 .await?;
4214 legacy_post_commit_result_with_embedding(
4215 "create_note_with_embedding_content",
4216 note.id,
4217 note,
4218 embedding,
4219 degradations,
4220 )
4221 }
4222
4223 #[allow(clippy::too_many_arguments)]
4224 pub async fn create_note_with_embedding_content_and_report(
4225 &self,
4226 token: &NamespaceToken,
4227 kind: &str,
4228 name: Option<&str>,
4229 content: &str,
4230 embedding_content: Option<&str>,
4231 salience: Option<f64>,
4232 properties: Option<serde_json::Value>,
4233 annotates: Vec<Uuid>,
4234 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4235 let (note, embedding, degradations, _) = self
4236 .create_note_inner(
4237 token,
4238 kind,
4239 name,
4240 content,
4241 embedding_content,
4242 salience,
4243 None,
4244 properties,
4245 annotates,
4246 None,
4247 false,
4248 false,
4249 )
4250 .await?;
4251 legacy_post_commit_result(
4252 "create_note_with_embedding_content_and_report",
4253 note.id,
4254 (note, embedding),
4255 degradations,
4256 )
4257 }
4258
4259 #[allow(clippy::too_many_arguments)]
4261 pub async fn create_note_with_embedding_content_and_post_commit_report(
4262 &self,
4263 token: &NamespaceToken,
4264 kind: &str,
4265 name: Option<&str>,
4266 content: &str,
4267 embedding_content: Option<&str>,
4268 salience: Option<f64>,
4269 properties: Option<serde_json::Value>,
4270 annotates: Vec<Uuid>,
4271 ) -> RuntimeResult<(
4272 Note,
4273 crate::retrieval::EmbeddingTruncationReport,
4274 Vec<PostCommitDegradation>,
4275 )> {
4276 let (note, embedding, degradations, _) = self
4277 .create_note_inner(
4278 token,
4279 kind,
4280 name,
4281 content,
4282 embedding_content,
4283 salience,
4284 None,
4285 properties,
4286 annotates,
4287 None,
4288 false,
4289 false,
4290 )
4291 .await?;
4292 Ok((note, embedding, degradations))
4293 }
4294
4295 #[allow(clippy::too_many_arguments)]
4299 pub async fn create_note_with_decay(
4300 &self,
4301 token: &NamespaceToken,
4302 kind: &str,
4303 name: Option<&str>,
4304 content: &str,
4305 salience: Option<f64>,
4306 decay_factor: f64,
4307 properties: Option<serde_json::Value>,
4308 annotates: Vec<Uuid>,
4309 ) -> RuntimeResult<Note> {
4310 self.create_note_with_decay_for_embedding_model(
4311 token,
4312 kind,
4313 name,
4314 content,
4315 salience,
4316 decay_factor,
4317 properties,
4318 annotates,
4319 None,
4320 )
4321 .await
4322 }
4323
4324 #[allow(clippy::too_many_arguments)]
4326 pub async fn create_note_with_decay_and_report(
4327 &self,
4328 token: &NamespaceToken,
4329 kind: &str,
4330 name: Option<&str>,
4331 content: &str,
4332 salience: Option<f64>,
4333 decay_factor: f64,
4334 properties: Option<serde_json::Value>,
4335 annotates: Vec<Uuid>,
4336 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4337 self.create_note_with_decay_for_embedding_model_and_report(
4338 token,
4339 kind,
4340 name,
4341 content,
4342 salience,
4343 decay_factor,
4344 properties,
4345 annotates,
4346 None,
4347 )
4348 .await
4349 }
4350
4351 #[allow(clippy::too_many_arguments)]
4356 pub async fn create_note_with_decay_for_embedding_model(
4357 &self,
4358 token: &NamespaceToken,
4359 kind: &str,
4360 name: Option<&str>,
4361 content: &str,
4362 salience: Option<f64>,
4363 decay_factor: f64,
4364 properties: Option<serde_json::Value>,
4365 annotates: Vec<Uuid>,
4366 embedding_model: Option<&str>,
4367 ) -> RuntimeResult<Note> {
4368 let (note, embedding, degradations, _) = self
4369 .create_note_inner(
4370 token,
4371 kind,
4372 name,
4373 content,
4374 None,
4375 salience,
4376 Some(decay_factor),
4377 properties,
4378 annotates,
4379 embedding_model,
4380 false,
4381 false,
4382 )
4383 .await?;
4384 legacy_post_commit_result_with_embedding(
4385 "create_note_with_decay_for_embedding_model",
4386 note.id,
4387 note,
4388 embedding,
4389 degradations,
4390 )
4391 }
4392
4393 #[allow(clippy::too_many_arguments)]
4395 pub async fn create_note_with_decay_for_embedding_model_and_report(
4396 &self,
4397 token: &NamespaceToken,
4398 kind: &str,
4399 name: Option<&str>,
4400 content: &str,
4401 salience: Option<f64>,
4402 decay_factor: f64,
4403 properties: Option<serde_json::Value>,
4404 annotates: Vec<Uuid>,
4405 embedding_model: Option<&str>,
4406 ) -> RuntimeResult<(Note, crate::retrieval::EmbeddingTruncationReport)> {
4407 let (note, embedding, degradations, _) = self
4408 .create_note_inner(
4409 token,
4410 kind,
4411 name,
4412 content,
4413 None,
4414 salience,
4415 Some(decay_factor),
4416 properties,
4417 annotates,
4418 embedding_model,
4419 false,
4420 false,
4421 )
4422 .await?;
4423 legacy_post_commit_result(
4424 "create_note_with_decay_for_embedding_model_and_report",
4425 note.id,
4426 (note, embedding),
4427 degradations,
4428 )
4429 }
4430
4431 #[allow(clippy::too_many_arguments)]
4435 pub async fn create_note_with_decay_for_embedding_model_with_visibility(
4436 &self,
4437 token: &NamespaceToken,
4438 kind: &str,
4439 name: Option<&str>,
4440 content: &str,
4441 salience: Option<f64>,
4442 decay_factor: f64,
4443 properties: Option<serde_json::Value>,
4444 annotates: Vec<Uuid>,
4445 embedding_model: Option<&str>,
4446 ) -> RuntimeResult<(Note, Vec<(String, u64)>)> {
4447 let (note, _, degradations, fences) = self
4448 .create_note_inner(
4449 token,
4450 kind,
4451 name,
4452 content,
4453 None,
4454 salience,
4455 Some(decay_factor),
4456 properties,
4457 annotates,
4458 embedding_model,
4459 true,
4460 false,
4461 )
4462 .await?;
4463 legacy_post_commit_result(
4464 "create_note_with_decay_for_embedding_model_with_visibility",
4465 note.id,
4466 (note, fences),
4467 degradations,
4468 )
4469 }
4470
4471 #[allow(clippy::too_many_arguments)]
4475 pub async fn create_note_with_decay_for_embedding_model_with_visibility_and_report(
4476 &self,
4477 token: &NamespaceToken,
4478 kind: &str,
4479 name: Option<&str>,
4480 content: &str,
4481 salience: Option<f64>,
4482 decay_factor: f64,
4483 properties: Option<serde_json::Value>,
4484 annotates: Vec<Uuid>,
4485 embedding_model: Option<&str>,
4486 ) -> RuntimeResult<(
4487 Note,
4488 Vec<(String, u64)>,
4489 crate::retrieval::EmbeddingTruncationReport,
4490 )> {
4491 let (note, embedding, degradations, fences) = self
4492 .create_note_inner(
4493 token,
4494 kind,
4495 name,
4496 content,
4497 None,
4498 salience,
4499 Some(decay_factor),
4500 properties,
4501 annotates,
4502 embedding_model,
4503 true,
4504 false,
4505 )
4506 .await?;
4507 legacy_post_commit_result(
4508 "create_note_with_decay_for_embedding_model_with_visibility_and_report",
4509 note.id,
4510 (note, fences, embedding),
4511 degradations,
4512 )
4513 }
4514
4515 pub async fn try_create_note(
4538 &self,
4539 token: &NamespaceToken,
4540 kind: &str,
4541 name: Option<&str>,
4542 content: &str,
4543 properties: Option<serde_json::Value>,
4544 ) -> RuntimeResult<Option<Note>> {
4545 self.try_create_note_impl(token, kind, name, content, properties, false, None, None)
4546 .await
4547 }
4548
4549 #[allow(clippy::too_many_arguments)]
4564 pub async fn try_create_note_as_trusted_ingest(
4565 &self,
4566 _capability: &crate::pack::ChannelIngestCapability,
4567 token: &NamespaceToken,
4568 kind: &str,
4569 name: Option<&str>,
4570 content: &str,
4571 properties: Option<serde_json::Value>,
4572 expires_after: Option<std::time::Duration>,
4573 ) -> RuntimeResult<Option<Note>> {
4574 self.try_create_note_impl(
4575 token,
4576 kind,
4577 name,
4578 content,
4579 properties,
4580 true,
4581 None,
4582 expires_after,
4583 )
4584 .await
4585 }
4586
4587 #[allow(clippy::too_many_arguments)]
4591 pub async fn try_create_note_as_trusted_ingest_with_attachment(
4592 &self,
4593 _capability: &crate::pack::ChannelIngestCapability,
4594 token: &NamespaceToken,
4595 kind: &str,
4596 name: Option<&str>,
4597 content: &str,
4598 properties: Option<serde_json::Value>,
4599 attachment: NewAttachment,
4600 expires_after: Option<std::time::Duration>,
4601 ) -> RuntimeResult<Option<Note>> {
4602 self.try_create_note_impl(
4603 token,
4604 kind,
4605 name,
4606 content,
4607 properties,
4608 true,
4609 Some(attachment),
4610 expires_after,
4611 )
4612 .await
4613 }
4614
4615 #[allow(clippy::too_many_arguments)]
4616 async fn try_create_note_impl(
4617 &self,
4618 token: &NamespaceToken,
4619 kind: &str,
4620 name: Option<&str>,
4621 content: &str,
4622 properties: Option<serde_json::Value>,
4623 allow_transport_owned_message_properties: bool,
4624 attachment: Option<NewAttachment>,
4625 expires_after: Option<std::time::Duration>,
4626 ) -> RuntimeResult<Option<Note>> {
4627 self.validate_note_kind(kind)?;
4628 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
4629 crate::secret_gate::check_at(content, "note", "content")?;
4630 if let Some(n) = name {
4631 crate::secret_gate::check_at(n, "note", "name")?;
4632 }
4633 if let Some(ref p) = properties {
4634 crate::secret_gate::check_json_at(p, "note", "properties")?;
4635 }
4636 if let Some(ref attachment) = attachment {
4637 drop(self.attachments()?);
4640 attachment.validate()?;
4641 let blob_store = self.blob_store().ok_or_else(|| {
4642 RuntimeError::Unconfigured(
4643 "trusted ingest attachment requires an installed BlobStore".to_string(),
4644 )
4645 })?;
4646 if !blob_store.exists(&attachment.content_ref).await? {
4647 return Err(RuntimeError::InvalidInput(format!(
4648 "trusted ingest attachment refers to an unpublished blob: {}",
4649 attachment.content_ref
4650 )));
4651 }
4652 }
4653 if !allow_transport_owned_message_properties && kind == "message" {
4654 if let Some(key) = properties
4655 .as_ref()
4656 .and_then(serde_json::Value::as_object)
4657 .and_then(|supplied| {
4658 crate::curation::kind_owned_properties("message")
4659 .iter()
4660 .copied()
4661 .find(|key| supplied.contains_key(*key))
4662 })
4663 {
4664 return Err(RuntimeError::InvalidInput(format!(
4665 "`{key}` is transport-owned on a `message` note and cannot be supplied \
4666 through `try_create_note`; only the trusted channel-ingest path may \
4667 establish quarantine disposition and channel provenance"
4668 )));
4669 }
4670 }
4671
4672 let ns = token.namespace().as_str();
4673 let mut note = Note::new(ns, kind, content);
4674 if let Some(retention) = expires_after {
4675 let duration_us = i64::try_from(retention.as_micros()).map_err(|_| {
4676 RuntimeError::InvalidInput(
4677 "trusted ingest retention exceeds i64 microseconds".into(),
4678 )
4679 })?;
4680 note.expires_at = Some(note.created_at.checked_add(duration_us).ok_or_else(|| {
4681 RuntimeError::InvalidInput("trusted ingest expiry exceeds i64 microseconds".into())
4682 })?);
4683 }
4684 if let Some(n) = name {
4685 note = note.with_name(n);
4686 }
4687 if let Some(p) = properties {
4688 note = note.with_properties(p);
4689 }
4690
4691 let inserted = if let Some(attachment) = attachment {
4698 self.raw_notes(token)?
4699 .try_insert_note_with_attachments(
4700 note.clone(),
4701 vec![Attachment::from_new(
4702 note.id,
4703 AttachmentSubstrate::Note,
4704 attachment,
4705 note.created_at,
4706 )],
4707 )
4708 .await?
4709 } else {
4710 self.raw_notes(token)?.try_insert_note(note.clone()).await?
4711 };
4712 if !inserted {
4713 return Ok(None);
4714 }
4715
4716 let mut degradations = Vec::new();
4717 match self.text_for_notes(token) {
4718 Ok(fts) => {
4719 if let Err(error) = fts.upsert_document(note_fts_document(¬e)).await {
4720 record_conditional_insert_degradation(
4721 &mut degradations,
4722 note.id,
4723 ConditionalInsertStage::FtsUpsert,
4724 error,
4725 );
4726 }
4727 }
4728 Err(error) => record_conditional_insert_degradation(
4729 &mut degradations,
4730 note.id,
4731 ConditionalInsertStage::FtsAcquisition,
4732 error,
4733 ),
4734 }
4735
4736 let embed_model_names = self.embedding_models_for_note_kind(kind);
4737 for model_name in &embed_model_names {
4738 match self
4739 .embed_document_with_model_outcome_for_token(
4740 token,
4741 model_name,
4742 note_embedding_text_ref(¬e),
4743 )
4744 .await
4745 {
4746 Ok(outcome) => {
4747 if outcome.truncated {
4748 tracing::warn!(
4749 note_id = %note.id,
4750 model = %outcome.model_name,
4751 source_bytes = outcome.source_bytes,
4752 embedded_bytes = outcome.embedded_bytes,
4753 "try_create_note: embedding input truncated; full content stored unchanged"
4754 );
4755 }
4756 match self.vectors_for_model(token, model_name) {
4757 Ok(vs) => {
4758 if let Err(error) = vs
4759 .insert(
4760 note.id,
4761 SubstrateKind::Note,
4762 ns,
4763 "note.content",
4764 vec![outcome.vector],
4765 )
4766 .await
4767 {
4768 record_conditional_insert_degradation(
4769 &mut degradations,
4770 note.id,
4771 ConditionalInsertStage::VectorInsert,
4772 format!("model {model_name}: {error}"),
4773 );
4774 }
4775 }
4776 Err(error) => record_conditional_insert_degradation(
4777 &mut degradations,
4778 note.id,
4779 ConditionalInsertStage::VectorAcquisition,
4780 format!("model {model_name}: {error}"),
4781 ),
4782 }
4783 }
4784 Err(error) => record_conditional_insert_degradation(
4785 &mut degradations,
4786 note.id,
4787 ConditionalInsertStage::Embedding,
4788 format!("model {model_name}: {error}"),
4789 ),
4790 }
4791 }
4792
4793 legacy_post_commit_result("try_create_note", note.id, Some(note), degradations)
4794 }
4795
4796 #[allow(clippy::too_many_arguments)]
4800 async fn create_note_inner(
4801 &self,
4802 token: &NamespaceToken,
4803 kind: &str,
4804 name: Option<&str>,
4805 content: &str,
4806 embedding_content: Option<&str>,
4807 salience: Option<f64>,
4808 decay_factor: Option<f64>,
4809 properties: Option<serde_json::Value>,
4810 annotates: Vec<Uuid>,
4811 embedding_model: Option<&str>,
4812 capture_visibility: bool,
4813 web_receipt: bool,
4814 ) -> RuntimeResult<(
4815 Note,
4816 crate::retrieval::EmbeddingTruncationReport,
4817 Vec<PostCommitDegradation>,
4818 Vec<(String, u64)>,
4819 )> {
4820 self.validate_note_kind(kind)?;
4821 let mut properties = self.derive_note_write_properties(kind, token, properties)?;
4827 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
4828 if web_receipt {
4829 let map = properties
4830 .as_mut()
4831 .and_then(serde_json::Value::as_object_mut)
4832 .expect("web receipt properties are constructed as an object");
4833 map.insert(
4834 crate::secret_gate::RESERVED_WEB_RECEIPT_KEY.to_string(),
4835 serde_json::Value::String(
4836 crate::secret_gate::WEB_RECEIPT_PROVENANCE_VALUE.to_string(),
4837 ),
4838 );
4839 }
4840 crate::secret_gate::check_at(content, "note", "content")?;
4842 if let Some(n) = name {
4843 crate::secret_gate::check_at(n, "note", "name")?;
4844 }
4845 if let Some(ref p) = properties {
4846 crate::secret_gate::check_json_at(p, "note", "properties")?;
4847 }
4848 if let Some(ec) = embedding_content {
4854 if ec.is_empty() {
4855 return Err(RuntimeError::InvalidInput(
4856 "embedding_content must not be empty".into(),
4857 ));
4858 }
4859 if ec.len() >= content.len() || !content.starts_with(ec) {
4860 return Err(RuntimeError::InvalidInput(
4861 "embedding_content must be a proper prefix of content".into(),
4862 ));
4863 }
4864 crate::secret_gate::check_at(ec, "note", "embedding_content")?;
4865 }
4866 let ns = token.namespace().as_str();
4867
4868 for &target_id in &annotates {
4871 if !self.substrate_exists_by_id(token, target_id).await? {
4872 return Err(RuntimeError::NotFound(format!(
4873 "create_note annotates target {target_id} not found"
4874 )));
4875 }
4876 }
4877
4878 if let Some(s) = salience {
4881 if !s.is_finite() || !(0.0..=1.0).contains(&s) {
4882 return Err(RuntimeError::InvalidInput(format!(
4883 "salience must be a finite value in [0.0, 1.0]; got {s}"
4884 )));
4885 }
4886 }
4887 if let Some(d) = decay_factor {
4888 if !d.is_finite() || d < 0.0 {
4889 return Err(RuntimeError::InvalidInput(format!(
4890 "decay_factor must be a finite value >= 0.0; got {d}"
4891 )));
4892 }
4893 }
4894
4895 if let Some(model_name) = embedding_model {
4899 self.resolve_embedding_model(Some(model_name))?;
4900 }
4901
4902 let mut note = Note::new(ns, kind, content);
4903 if let Some(s) = salience {
4904 note = note.with_salience(s);
4905 }
4906 if let Some(df) = decay_factor {
4907 note = note.with_decay(df);
4908 }
4909 if let Some(n) = name {
4910 note = note.with_name(n);
4911 }
4912 if let Some(p) = properties {
4913 note = note.with_properties(p);
4914 }
4915 let notes = if web_receipt {
4916 self.raw_notes(token)?
4917 } else {
4918 self.notes(token)?
4919 };
4920 notes.upsert_note(note.clone()).await?;
4921
4922 let embed_model_names: Vec<String> = if let Some(m) = embedding_model {
4928 vec![m.to_string()]
4929 } else {
4930 self.embedding_models_for_note_kind(kind)
4931 };
4932
4933 {
4935 #[cfg(any(test, feature = "fault-injection"))]
4941 let fts_inject = consume_fault(&FTS_FAIL_NS, ns);
4942 #[cfg(not(any(test, feature = "fault-injection")))]
4943 let fts_inject = false;
4944 let fts_result: RuntimeResult<()> = if fts_inject {
4945 Err(RuntimeError::Internal("injected FTS failure".to_string()))
4946 } else {
4947 let statements =
4948 khive_db::stores::text::delete_document_statements("fts_notes", ns, note.id)
4949 .into_iter()
4950 .chain(khive_db::stores::text::insert_document_statements(
4951 "fts_notes",
4952 ¬e_fts_document(¬e),
4953 ))
4954 .collect();
4955 self.apply_note_index_revision(¬e, statements)
4956 .await
4957 .map(|_| ())
4958 };
4959
4960 if let Err(e) = fts_result {
4961 self.compensate_note_creation(¬e).await;
4962 return Err(e);
4963 }
4964 }
4965
4966 let canonical_embed_text = note_embedding_text_ref(¬e);
4976 let embed_text = embedding_content.unwrap_or(canonical_embed_text);
4977
4978 let mut embedding_report = crate::retrieval::EmbeddingTruncationReport::default();
4979 let mut vector_fences = Vec::with_capacity(embed_model_names.len());
4980 if embed_model_names.len() == 1 {
4981 let model_name = &embed_model_names[0];
4983 let vec_result = self
4984 .embed_document_with_model_outcome_for_token(token, model_name, embed_text)
4985 .await;
4986
4987 #[cfg(any(test, feature = "fault-injection"))]
4997 let vec_inject = {
4998 let ns_inject = consume_fault(&VECTOR_FAIL_NS, ns);
4999 let count_inject = VECTOR_FAIL_AFTER.with(|cell| match cell.get() {
5000 Some(0) => {
5001 cell.set(None);
5002 true
5003 }
5004 Some(n) => {
5005 cell.set(Some(n - 1));
5006 false
5007 }
5008 None => false,
5009 });
5010 ns_inject || count_inject
5011 };
5012 #[cfg(not(any(test, feature = "fault-injection")))]
5013 let vec_inject = false;
5014 let vec_result: RuntimeResult<crate::retrieval::DocumentEmbeddingOutcome> =
5015 if vec_inject {
5016 Err(RuntimeError::Internal(
5017 "injected vector failure".to_string(),
5018 ))
5019 } else {
5020 vec_result
5021 };
5022
5023 let single_model_result: RuntimeResult<()> = match vec_result {
5024 Ok(outcome) => {
5025 embedding_report.observe(&outcome);
5026 if capture_visibility {
5027 match self
5028 .publish_note_vector_revision_with_seq(
5029 token,
5030 ¬e,
5031 model_name,
5032 &outcome.vector,
5033 )
5034 .await
5035 {
5036 Ok(Some(seq)) => {
5037 vector_fences.push((model_name.clone(), seq));
5038 Ok(())
5039 }
5040 Ok(None) => Ok(()),
5041 Err(error) => Err(error),
5042 }
5043 } else {
5044 self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
5045 .await
5046 .map(|_| ())
5047 }
5048 }
5049 Err(e) => Err(e),
5050 };
5051 if let Err(e) = single_model_result {
5052 self.compensate_note_creation(¬e).await;
5053 return Err(e);
5054 }
5055 } else if !embed_model_names.is_empty() {
5056 let rt_clone = self.clone();
5059 let content_owned: std::sync::Arc<str> = std::sync::Arc::from(embed_text);
5062 let usage_ctx = crate::usage::current();
5063 let mut join_set = tokio::task::JoinSet::new();
5064 for (idx, model_name) in embed_model_names.iter().enumerate() {
5065 let rt = rt_clone.clone();
5066 let text = std::sync::Arc::clone(&content_owned);
5067 let name = model_name.clone();
5068 let ctx = usage_ctx.clone();
5069 let token = (*token).clone();
5070 join_set.spawn(crate::runtime::inherit_request_embedder_scope(async move {
5071 let fut = rt.embed_document_with_model_outcome_for_token(
5072 &token,
5073 &name,
5074 text.as_ref(),
5075 );
5076 let result = match ctx {
5077 Some(ctx) => crate::usage::scope(ctx, fut).await,
5078 None => fut.await,
5079 };
5080 (idx, result)
5081 }));
5082 }
5083 let outcomes = match drain_embed_join_set(join_set, embed_model_names.len()).await {
5087 Ok(outcomes) => outcomes,
5088 Err(e) => {
5089 self.compensate_note_creation(¬e).await;
5090 return Err(e);
5091 }
5092 };
5093 for (model_name, outcome) in embed_model_names.iter().zip(outcomes) {
5095 embedding_report.observe(&outcome);
5096 let insert_result = if capture_visibility {
5097 self.publish_note_vector_revision_with_seq(
5098 token,
5099 ¬e,
5100 model_name,
5101 &outcome.vector,
5102 )
5103 .await
5104 .map(|seq| {
5105 if let Some(seq) = seq {
5106 vector_fences.push((model_name.clone(), seq));
5107 }
5108 })
5109 } else {
5110 self.publish_note_vector_revision(token, ¬e, model_name, &outcome.vector)
5111 .await
5112 .map(|_| ())
5113 };
5114 if let Err(e) = insert_result {
5115 self.compensate_note_creation(¬e).await;
5116 return Err(e);
5117 }
5118 }
5119 }
5120
5121 let mut created_edges: Vec<Uuid> = Vec::with_capacity(annotates.len());
5126
5127 #[cfg(test)]
5130 let annotates_iter: Vec<(usize, Uuid)> = annotates
5131 .iter()
5132 .enumerate()
5133 .map(|(i, &id)| (i, id))
5134 .collect();
5135 #[cfg(test)]
5136 macro_rules! next_target {
5137 ($pair:expr) => {
5138 $pair.1
5139 };
5140 }
5141 #[cfg(not(test))]
5142 let annotates_iter: Vec<Uuid> = annotates.to_vec();
5143 #[cfg(not(test))]
5144 macro_rules! next_target {
5145 ($pair:expr) => {
5146 $pair
5147 };
5148 }
5149
5150 for pair in annotates_iter {
5151 let target_id = next_target!(pair);
5152
5153 #[cfg(test)]
5155 let injected_err: Option<RuntimeError> = {
5156 let call_idx = pair.0;
5157 LINK_FAIL_AFTER.with(|cell| {
5158 let n = cell.get();
5159 if n > 0 && call_idx + 1 == n {
5160 cell.set(0); Some(RuntimeError::Internal("injected link failure".to_string()))
5162 } else {
5163 None
5164 }
5165 })
5166 };
5167 #[cfg(not(test))]
5168 let injected_err: Option<RuntimeError> = None;
5169
5170 let link_result = if let Some(e) = injected_err {
5171 Err(e)
5172 } else {
5173 self.link(
5174 token,
5175 note.id,
5176 target_id,
5177 EdgeRelation::Annotates,
5178 1.0,
5179 None,
5180 )
5181 .await
5182 };
5183
5184 match link_result {
5185 Ok(edge) => created_edges.push(edge.id.into()),
5186 Err(e) => {
5187 let edge_ids = created_edges
5192 .iter()
5193 .map(Uuid::to_string)
5194 .collect::<Vec<_>>()
5195 .join(", ");
5196 match self.compensate_note_creation_with_edges(¬e).await {
5197 Ok(true) => return Err(e),
5198 Ok(false) => {
5199 return Err(RuntimeError::Internal(format!(
5200 "create_note: annotates link failed: {e}; note {} changed before \
5201 compensation, retaining its incident edges [{edge_ids}]",
5202 note.id
5203 )));
5204 }
5205 Err(cleanup_error) => {
5206 return Err(RuntimeError::Internal(format!(
5207 "create_note: annotates link failed: {e}; compensation failed \
5208 for note {} and retained edges [{edge_ids}]: {cleanup_error}; \
5209 note and edges remain for reconciliation",
5210 note.id
5211 )));
5212 }
5213 }
5214 }
5215 }
5216 }
5217
5218 let created_event = khive_storage::event::Event::new(
5223 note.namespace.clone(),
5224 "create",
5225 EventKind::NoteCreated,
5226 SubstrateKind::Note,
5227 "",
5228 )
5229 .with_target(note.id)
5230 .with_payload(serde_json::json!({
5231 "id": note.id,
5232 "namespace": note.namespace,
5233 "kind": note.kind,
5234 "salience": note.salience,
5235 }));
5236 let event_result = match self.events(token) {
5237 Ok(store) => store
5238 .append_event(created_event)
5239 .await
5240 .map_err(RuntimeError::from),
5241 Err(error) => Err(error),
5242 };
5243 let mut degradations = Vec::new();
5244 if let Err(error) = event_result {
5245 record_post_commit_degradation(
5246 &mut degradations,
5247 "create_note",
5248 note.id,
5249 "event_append",
5250 error,
5251 );
5252 }
5253
5254 vector_fences.sort_by(|a, b| a.0.cmp(&b.0));
5255 Ok((note, embedding_report, degradations, vector_fences))
5256 }
5257
5258 pub async fn list_notes(
5264 &self,
5265 token: &NamespaceToken,
5266 kind: Option<&str>,
5267 limit: u32,
5268 offset: u32,
5269 ) -> RuntimeResult<Vec<Note>> {
5270 let visible = token.visible_namespaces();
5271 if visible.len() == 1 {
5272 let page = self
5274 .notes(token)?
5275 .query_notes_count_free(
5276 token.namespace().as_str(),
5277 kind,
5278 PageRequest {
5279 offset: offset.into(),
5280 limit,
5281 },
5282 )
5283 .await?;
5284 return Ok(page.items);
5285 }
5286 use khive_storage::note::NoteFilter;
5288 let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
5289 let filter = NoteFilter {
5290 kind: kind.map(|k| k.to_string()),
5291 namespaces: ns_strs,
5292 ..Default::default()
5293 };
5294 let page = self
5295 .notes(token)?
5296 .query_notes_filtered_count_free(
5297 token.namespace().as_str(),
5298 &filter,
5299 PageRequest {
5300 offset: offset.into(),
5301 limit,
5302 },
5303 )
5304 .await?;
5305 Ok(page.items)
5306 }
5307
5308 pub async fn list_notes_after(
5314 &self,
5315 token: &NamespaceToken,
5316 kind: Option<&str>,
5317 after: Option<Uuid>,
5318 limit: u32,
5319 ) -> RuntimeResult<(Vec<Note>, Option<Uuid>)> {
5320 let store = self.notes(token)?;
5321 let after = match after {
5322 Some(id) => {
5323 let note = self
5324 .get_note_including_deleted(token, id)
5325 .await?
5326 .ok_or_else(|| RuntimeError::NotFound(format!("note cursor {id}")))?;
5327 Self::ensure_namespace_visible(¬e.namespace, token)?;
5328 let sequence = store.note_sequence(id).await?.ok_or_else(|| {
5329 RuntimeError::Internal(format!(
5330 "note cursor {id} has no insertion-sequence ledger row"
5331 ))
5332 })?;
5333 Some(SeekCursor { sequence, id })
5334 }
5335 None => None,
5336 };
5337 let filter = khive_storage::note::NoteFilter {
5338 kind: kind.map(str::to_string),
5339 namespaces: token
5340 .visible_namespaces()
5341 .iter()
5342 .map(|namespace| namespace.as_str().to_owned())
5343 .collect(),
5344 ..Default::default()
5345 };
5346 let page = store
5347 .query_notes_filtered_after(token.namespace().as_str(), &filter, after, limit)
5348 .await?;
5349 Ok((page.items, page.next_after.map(|cursor| cursor.id)))
5350 }
5351
5352 pub async fn count_notes(
5354 &self,
5355 token: &NamespaceToken,
5356 kind: Option<&str>,
5357 ) -> RuntimeResult<u64> {
5358 let namespaces: Vec<String> = token
5359 .visible_namespaces()
5360 .iter()
5361 .map(|namespace| namespace.as_str().to_owned())
5362 .collect();
5363 Ok(self
5364 .notes(token)?
5365 .count_notes_in_namespaces(&namespaces, kind)
5366 .await?)
5367 }
5368
5369 #[allow(clippy::too_many_arguments)]
5389 pub async fn search_notes(
5390 &self,
5391 token: &NamespaceToken,
5392 query_text: &str,
5393 query_vector: Option<Vec<f32>>,
5394 limit: u32,
5395 note_kind: Option<&str>,
5396 include_superseded: bool,
5397 tags_any: &[String],
5398 properties_filter: Option<&serde_json::Value>,
5399 ) -> RuntimeResult<Vec<NoteSearchHit>> {
5400 self.search_notes_with_text_mode(
5401 token,
5402 query_text,
5403 query_vector,
5404 limit,
5405 note_kind,
5406 include_superseded,
5407 tags_any,
5408 properties_filter,
5409 TextQueryMode::Plain,
5410 )
5411 .await
5412 }
5413
5414 #[allow(clippy::too_many_arguments)]
5416 pub async fn search_notes_with_text_mode(
5417 &self,
5418 token: &NamespaceToken,
5419 query_text: &str,
5420 query_vector: Option<Vec<f32>>,
5421 limit: u32,
5422 note_kind: Option<&str>,
5423 include_superseded: bool,
5424 tags_any: &[String],
5425 properties_filter: Option<&serde_json::Value>,
5426 text_mode: TextQueryMode,
5427 ) -> RuntimeResult<Vec<NoteSearchHit>> {
5428 let (hits, _vector_error) = self
5429 .search_notes_inner(
5430 token,
5431 query_text,
5432 query_vector,
5433 limit,
5434 note_kind,
5435 include_superseded,
5436 tags_any,
5437 properties_filter,
5438 text_mode,
5439 false,
5440 )
5441 .await?;
5442 Ok(hits)
5443 }
5444
5445 #[allow(clippy::too_many_arguments)]
5452 pub async fn search_notes_outcome(
5453 &self,
5454 token: &NamespaceToken,
5455 query_text: &str,
5456 limit: u32,
5457 note_kind: Option<&str>,
5458 include_superseded: bool,
5459 tags_any: &[String],
5460 properties_filter: Option<&serde_json::Value>,
5461 ) -> RuntimeResult<NoteSearchOutcome> {
5462 self.search_notes_outcome_with_text_mode(
5463 token,
5464 query_text,
5465 limit,
5466 note_kind,
5467 include_superseded,
5468 tags_any,
5469 properties_filter,
5470 TextQueryMode::Plain,
5471 )
5472 .await
5473 }
5474
5475 #[allow(clippy::too_many_arguments)]
5477 pub async fn search_notes_outcome_with_text_mode(
5478 &self,
5479 token: &NamespaceToken,
5480 query_text: &str,
5481 limit: u32,
5482 note_kind: Option<&str>,
5483 include_superseded: bool,
5484 tags_any: &[String],
5485 properties_filter: Option<&serde_json::Value>,
5486 text_mode: TextQueryMode,
5487 ) -> RuntimeResult<NoteSearchOutcome> {
5488 let (hits, vector_error) = self
5489 .search_notes_inner(
5490 token,
5491 query_text,
5492 None,
5493 limit,
5494 note_kind,
5495 include_superseded,
5496 tags_any,
5497 properties_filter,
5498 text_mode,
5499 true,
5500 )
5501 .await?;
5502 Ok(NoteSearchOutcome { hits, vector_error })
5503 }
5504
5505 #[allow(clippy::too_many_arguments)]
5506 async fn search_notes_inner(
5507 &self,
5508 token: &NamespaceToken,
5509 query_text: &str,
5510 query_vector: Option<Vec<f32>>,
5511 limit: u32,
5512 note_kind: Option<&str>,
5513 include_superseded: bool,
5514 tags_any: &[String],
5515 properties_filter: Option<&serde_json::Value>,
5516 text_mode: TextQueryMode,
5517 tolerate_vector_error: bool,
5518 ) -> RuntimeResult<(Vec<NoteSearchHit>, Option<String>)> {
5519 const RRF_K: usize = 60;
5520 let candidates = limit.saturating_mul(4).max(limit);
5521 let visible_ns: Vec<String> = token
5522 .visible_namespaces()
5523 .iter()
5524 .map(|ns| ns.as_str().to_owned())
5525 .collect();
5526
5527 #[cfg(any(test, feature = "fault-injection"))]
5540 let fts_search_inject = {
5541 let mut g = FTS_SEARCH_FAIL_NS.lock().unwrap();
5542 match g.as_deref() {
5543 Some(armed) if visible_ns.iter().any(|ns| ns == armed) => {
5544 *g = None;
5545 true
5546 }
5547 _ => false,
5548 }
5549 };
5550 #[cfg(not(any(test, feature = "fault-injection")))]
5551 let fts_search_inject = false;
5552
5553 let text_store = self.text_for_notes(token)?;
5554 let text_fut = async {
5555 if fts_search_inject {
5556 return Err(khive_storage::StorageError::Timeout {
5557 operation: "fts_search".into(),
5558 });
5559 }
5560 text_store
5561 .search(TextSearchRequest {
5562 query: query_text.to_string(),
5563 mode: text_mode,
5564 filter: Some(TextFilter {
5565 namespaces: visible_ns.clone(),
5566 record_kinds: note_kind
5573 .map(|kind| vec![kind.to_string()])
5574 .unwrap_or_default(),
5575 ..TextFilter::default()
5576 }),
5577 top_k: candidates,
5578 snippet_chars: 200,
5579 })
5580 .await
5581 };
5582 let text_fut = crate::stage_seam::text_stage(text_fut);
5583
5584 let vector_fut = async {
5586 if query_vector.is_some() || self.config().embedding_model.is_some() {
5587 self.note_search_vector_search(token, query_vector, query_text, candidates)
5588 .await
5589 } else {
5590 Ok(vec![])
5591 }
5592 };
5593 let (text_search_result, vector_result) = tokio::join!(text_fut, vector_fut);
5594
5595 let text_hits = crate::error::fts_text_leg_or_err(
5601 text_search_result.map_err(RuntimeError::from),
5602 "search_notes",
5603 query_text,
5604 )?;
5605
5606 let mut vector_error: Option<String> = None;
5607 let vector_hits = match vector_result {
5608 Ok(hits) => hits,
5609 Err(e) if tolerate_vector_error => {
5610 vector_error = Some(e.to_string());
5611 Vec::new()
5612 }
5613 Err(e) => return Err(e),
5614 };
5615
5616 let fuse_k = text_hits.len() + vector_hits.len();
5622 let fused = crate::fusion::rrf_fuse_k(self, text_hits, vector_hits, RRF_K, fuse_k).await?;
5623
5624 let candidate_ids: Vec<Uuid> = fused.iter().map(|hit| hit.entity_id).collect();
5625 if candidate_ids.is_empty() {
5626 return Ok((vec![], vector_error));
5627 }
5628
5629 let note_store = self.notes(token)?;
5636 let search_pool = self.backend().pool_arc();
5637 let mailbox_view = crate::MailboxView {
5638 actor_id: token.actor().id.clone(),
5639 delegated: false,
5640 };
5641 let mut alive_notes: HashMap<Uuid, Note> = HashMap::new();
5642 for note in note_store.get_notes_batch(&candidate_ids).await? {
5643 search_pool.record_note_candidate_hydration_row();
5644 if note.deleted_at.is_some() {
5645 continue;
5646 }
5647 if !mailbox_view.permits_message_note(token, ¬e) {
5648 continue;
5649 }
5650 if let Some(want_kind) = note_kind {
5651 if note.kind != want_kind {
5652 continue;
5653 }
5654 }
5655 if !tags_any.is_empty() {
5660 let note_tags: Vec<String> = note
5661 .properties
5662 .as_ref()
5663 .and_then(|p| p.get("tags"))
5664 .and_then(serde_json::Value::as_array)
5665 .map(|arr| {
5666 arr.iter()
5667 .filter_map(serde_json::Value::as_str)
5668 .map(str::to_owned)
5669 .collect()
5670 })
5671 .unwrap_or_default();
5672 if !note_tags
5673 .iter()
5674 .any(|t| tags_any.iter().any(|w| t.eq_ignore_ascii_case(w)))
5675 {
5676 continue;
5677 }
5678 }
5679 if let Some(pf) = properties_filter {
5681 if !crate::retrieval::properties_match(note.properties.as_ref(), pf) {
5682 continue;
5683 }
5684 }
5685 alive_notes.insert(note.id, note);
5686 }
5687
5688 if !include_superseded && !alive_notes.is_empty() {
5691 let graph = self.graph(token)?;
5692 let note_ids: Vec<Uuid> = alive_notes.keys().copied().collect();
5693 let superseded: std::collections::HashSet<Uuid> = graph
5694 .batch_neighbors(
5695 ¬e_ids,
5696 NeighborQuery {
5697 direction: Direction::In,
5698 relations: Some(vec![EdgeRelation::Supersedes]),
5699 limit: Some(1),
5700 min_weight: None,
5701 },
5702 )
5703 .await?
5704 .into_iter()
5705 .map(|(note_id, _)| note_id)
5706 .collect();
5707 alive_notes.retain(|id, _| !superseded.contains(id));
5708 }
5709
5710 let mut hits: Vec<NoteSearchHit> = fused
5712 .into_iter()
5713 .filter_map(|hit| {
5714 let note = alive_notes.get(&hit.entity_id)?;
5715 let weighted = salience_weighted_rank(hit.score, note.salience);
5716 Some(NoteSearchHit {
5717 note_id: hit.entity_id,
5718 score: weighted,
5719 rank_score_kind: hit.rank_score_kind,
5720 signals: hit.signals,
5721 source: hit.source,
5722 title: hit.title.or_else(|| note_title(note)),
5723 snippet: hit.snippet.or_else(|| note_snippet(note)),
5724 })
5725 })
5726 .collect();
5727
5728 hits.sort_by(|a, b| b.score.cmp(&a.score).then(a.note_id.cmp(&b.note_id)));
5729 hits.truncate(limit as usize);
5730 Ok((hits, vector_error))
5731 }
5732
5733 pub async fn resolve_prefix(
5740 &self,
5741 token: &NamespaceToken,
5742 prefix: &str,
5743 ) -> RuntimeResult<Option<Uuid>> {
5744 let namespaces = [token.namespace().as_str().to_owned()];
5745 self.resolve_prefix_inner(Some(&namespaces), prefix, false, false)
5746 .await
5747 }
5748
5749 pub async fn resolve_prefix_including_deleted(
5750 &self,
5751 token: &NamespaceToken,
5752 prefix: &str,
5753 ) -> RuntimeResult<Option<Uuid>> {
5754 let namespaces = [token.namespace().as_str().to_owned()];
5755 self.resolve_prefix_inner(Some(&namespaces), prefix, true, false)
5756 .await
5757 }
5758
5759 pub async fn resolve_prefix_unfiltered(&self, prefix: &str) -> RuntimeResult<Option<Uuid>> {
5767 self.resolve_prefix_inner(None, prefix, false, false).await
5768 }
5769
5770 pub async fn resolve_prefix_unfiltered_including_deleted(
5773 &self,
5774 prefix: &str,
5775 ) -> RuntimeResult<Option<Uuid>> {
5776 self.resolve_prefix_inner(None, prefix, true, false).await
5777 }
5778
5779 pub(crate) async fn resolve_prefix_for_kg_read(
5782 &self,
5783 prefix: &str,
5784 include_deleted: bool,
5785 ) -> RuntimeResult<Option<Uuid>> {
5786 self.resolve_prefix_inner(None, prefix, include_deleted, true)
5787 .await
5788 }
5789
5790 async fn resolve_prefix_inner(
5802 &self,
5803 namespaces: Option<&[String]>,
5804 prefix: &str,
5805 include_deleted: bool,
5806 require_tables: bool,
5807 ) -> RuntimeResult<Option<Uuid>> {
5808 if !prefix.chars().all(|c| c.is_ascii_hexdigit() || c == '-') {
5815 return Ok(None);
5816 }
5817
5818 #[cfg(any(test, feature = "fault-injection"))]
5822 if consume_fault(&PREFIX_RESOLVE_FAIL_NS, prefix) {
5823 return Err(RuntimeError::Storage(
5824 khive_storage::StorageError::Timeout {
5825 operation: "resolve_prefix".into(),
5826 },
5827 ));
5828 }
5829
5830 let Some((lower, upper)) = uuid_prefix_bounds(prefix) else {
5831 return Ok(None);
5832 };
5833
5834 let tables = [
5835 ("entities", true),
5836 ("notes", true),
5837 ("events", false),
5838 ("graph_edges", false),
5839 ];
5840
5841 let mut matches: Vec<String> = Vec::new();
5850 let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
5851 let mut reader = self.sql().reader().await.map_err(RuntimeError::Storage)?;
5852
5853 for (table, has_deleted_at) in tables {
5854 let sql = resolve_prefix_statement(
5855 table,
5856 has_deleted_at,
5857 include_deleted,
5858 namespaces,
5859 &lower,
5860 &upper,
5861 );
5862 match reader.query_all(sql).await {
5863 Ok(rows) => {
5864 for row in rows {
5865 if let Some(col) = row.columns.first() {
5866 if let SqlValue::Text(s) = &col.value {
5867 if seen.insert(s.clone()) {
5868 matches.push(s.clone());
5869 }
5870 }
5871 }
5872 }
5873 }
5874 Err(e) => {
5875 let msg = e.to_string();
5876 if !require_tables && msg.contains("no such table") {
5877 continue;
5878 }
5879 return Err(RuntimeError::Storage(e));
5880 }
5881 }
5882 if matches.len() > 1 {
5883 break;
5884 }
5885 }
5886
5887 if matches.len() <= 1 {
5892 if let Some(sidecar_sql) = self.events_sidecar_sql_read_only()? {
5893 let mut sql = resolve_prefix_statement(
5896 "events",
5897 false,
5898 include_deleted,
5899 namespaces,
5900 &lower,
5901 &upper,
5902 );
5903 sql.label = Some("resolve_prefix.events_sidecar".into());
5904 let mut sidecar_reader =
5905 sidecar_sql.reader().await.map_err(RuntimeError::Storage)?;
5906 match sidecar_reader.query_all(sql).await {
5907 Ok(rows) => {
5908 for row in rows {
5909 if let Some(col) = row.columns.first() {
5910 if let SqlValue::Text(s) = &col.value {
5911 if seen.insert(s.clone()) {
5912 matches.push(s.clone());
5913 }
5914 }
5915 }
5916 }
5917 }
5918 Err(e) => {
5919 let msg = e.to_string();
5920 if require_tables || !msg.contains("no such table") {
5921 return Err(RuntimeError::Storage(e));
5922 }
5923 }
5924 }
5925 }
5926 }
5927
5928 match matches.len() {
5929 0 => Ok(None),
5930 1 => {
5931 let uuid = Uuid::from_str(&matches[0])
5932 .map_err(|e| RuntimeError::Internal(format!("stored UUID is invalid: {e}")))?;
5933 Ok(Some(uuid))
5934 }
5935 _ => {
5936 let uuids: Vec<uuid::Uuid> = matches
5937 .iter()
5938 .filter_map(|s| Uuid::from_str(s).ok())
5939 .collect();
5940 Err(RuntimeError::AmbiguousPrefix {
5941 prefix: prefix.to_string(),
5942 matches: uuids,
5943 })
5944 }
5945 }
5946 }
5947
5948 pub async fn resolve_by_id(
5959 &self,
5960 token: &NamespaceToken,
5961 id: Uuid,
5962 ) -> RuntimeResult<Option<Resolved>> {
5963 if let Some(entity) = self.entities(token)?.get_entity(id).await? {
5965 return Ok(Some(Resolved::Entity(entity)));
5966 }
5967
5968 if let Some(note) = self.notes(token)?.get_note(id).await? {
5970 return Ok(Some(Resolved::Note(note)));
5971 }
5972
5973 Ok(None)
5976 }
5977
5978 pub async fn resolve_by_id_including_deleted(
5984 &self,
5985 token: &NamespaceToken,
5986 id: Uuid,
5987 ) -> RuntimeResult<Option<Resolved>> {
5988 if let Some(entity) = self
5990 .entities(token)?
5991 .get_entity_including_deleted(id)
5992 .await?
5993 {
5994 return Ok(Some(Resolved::Entity(entity)));
5995 }
5996
5997 if let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? {
5999 return Ok(Some(Resolved::Note(note)));
6000 }
6001
6002 Ok(None)
6005 }
6006
6007 pub async fn resolve(
6012 &self,
6013 token: &NamespaceToken,
6014 id: Uuid,
6015 ) -> RuntimeResult<Option<Resolved>> {
6016 match self.get_entity(token, id).await {
6018 Ok(entity) => return Ok(Some(Resolved::Entity(entity))),
6019 Err(RuntimeError::NotFound(_) | RuntimeError::NamespaceMismatch { .. }) => {}
6020 Err(e) => return Err(e),
6021 }
6022
6023 if let Some(note) = self.notes(token)?.get_note(id).await? {
6025 if Self::ensure_namespace_visible(¬e.namespace, token).is_ok() {
6026 return Ok(Some(Resolved::Note(note)));
6027 }
6028 }
6029
6030 if let Some(event) = self.events(token)?.get_event(id).await? {
6032 if Self::ensure_namespace_visible(&event.namespace, token).is_ok() {
6033 return Ok(Some(Resolved::Event(event)));
6034 }
6035 }
6036
6037 Ok(None)
6038 }
6039
6040 pub async fn resolve_edge_endpoint(
6051 &self,
6052 token: &NamespaceToken,
6053 id: Uuid,
6054 ) -> RuntimeResult<Option<Resolved>> {
6055 if let Some(resolved) = self.resolve_by_id(token, id).await? {
6056 return Ok(Some(resolved));
6057 }
6058 if let Some(event) = self.events(token)?.get_event(id).await? {
6059 return Ok(Some(Resolved::Event(event)));
6060 }
6061 Ok(None)
6062 }
6063
6064 pub async fn resolve_primary(
6069 &self,
6070 token: &NamespaceToken,
6071 id: Uuid,
6072 ) -> RuntimeResult<Option<Resolved>> {
6073 self.resolve_primary_inner(token, id, false).await
6074 }
6075
6076 pub async fn resolve_including_deleted(
6081 &self,
6082 token: &NamespaceToken,
6083 id: Uuid,
6084 ) -> RuntimeResult<Option<Resolved>> {
6085 self.resolve_primary_inner(token, id, true).await
6086 }
6087
6088 async fn resolve_primary_inner(
6092 &self,
6093 token: &NamespaceToken,
6094 id: Uuid,
6095 include_deleted: bool,
6096 ) -> RuntimeResult<Option<Resolved>> {
6097 let ns = token.namespace().as_str();
6098
6099 let entity = if include_deleted {
6100 self.entities(token)?
6101 .get_entity_including_deleted(id)
6102 .await?
6103 } else {
6104 self.entities(token)?.get_entity(id).await?
6105 };
6106 if let Some(entity) = entity {
6107 if Self::ensure_namespace(&entity.namespace, ns).is_ok() {
6108 return Ok(Some(Resolved::Entity(entity)));
6109 }
6110 }
6111
6112 let note = if include_deleted {
6113 self.notes(token)?.get_note_including_deleted(id).await?
6114 } else {
6115 self.notes(token)?.get_note(id).await?
6116 };
6117 if let Some(note) = note {
6118 if Self::ensure_namespace(¬e.namespace, ns).is_ok() {
6119 return Ok(Some(Resolved::Note(note)));
6120 }
6121 }
6122
6123 if let Some(event) = self.events(token)?.get_event(id).await? {
6124 if Self::ensure_namespace(&event.namespace, ns).is_ok() {
6125 return Ok(Some(Resolved::Event(event)));
6126 }
6127 }
6128
6129 Ok(None)
6130 }
6131
6132 async fn atomic_hard_delete_with_edge_purge(
6144 &self,
6145 row_statement: SqlStatement,
6146 node_id: Uuid,
6147 namespace: &str,
6148 actor: &str,
6149 substrate: SubstrateKind,
6150 ) -> RuntimeResult<bool> {
6151 let mut statements = vec![PlanStatement {
6152 statement: row_statement,
6153 guard: Some(AffectedRowGuard::exactly(1)),
6154 }];
6155 if matches!(substrate, SubstrateKind::Entity | SubstrateKind::Note) {
6156 statements.push(PlanStatement {
6157 statement: khive_db::stores::attachment::delete_record_attachments_statement(
6158 node_id,
6159 if substrate == SubstrateKind::Entity {
6160 AttachmentSubstrate::Entity
6161 } else {
6162 AttachmentSubstrate::Note
6163 },
6164 ),
6165 guard: None,
6166 });
6167 }
6168 statements.extend(
6169 hard_delete_lineage_warning_statements(namespace, actor, node_id, substrate)
6170 .into_iter()
6171 .map(|statement| PlanStatement {
6172 statement,
6173 guard: None,
6174 }),
6175 );
6176 statements.push(PlanStatement {
6177 statement: purge_incident_edges_statement(node_id),
6178 guard: None,
6179 });
6180 let plan = AtomicOpPlan::Delete(DeletePlan {
6181 target_id: node_id,
6182 statements,
6183 post_commit: PostCommitEffect::None,
6184 });
6185 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6186 Ok(AtomicRunOutcome::Committed { .. }) => Ok(true),
6187 Ok(AtomicRunOutcome::RolledBack {
6188 failure: AtomicOpFailure::NoteConflict(conflict),
6189 ..
6190 }) => Err(conflict.into_error().into()),
6191 Ok(AtomicRunOutcome::RolledBack {
6192 failure: AtomicOpFailure::EntityConflict(conflict),
6193 ..
6194 }) => Err(conflict.into_error().into()),
6195 Ok(AtomicRunOutcome::RolledBack {
6196 failure: AtomicOpFailure::GuardFailed { .. },
6197 ..
6198 }) => Ok(false),
6199 Ok(AtomicRunOutcome::RolledBack {
6200 failure: AtomicOpFailure::SqlError { message, .. },
6201 ..
6202 }) => Err(RuntimeError::Internal(format!(
6203 "hard delete + edge purge for {node_id} failed: {message}"
6204 ))),
6205 Err(e) => Err(RuntimeError::Internal(format!(
6206 "hard delete + edge purge for {node_id}: atomic unit seam failure: {}",
6207 e.0
6208 ))),
6209 }
6210 }
6211
6212 pub async fn restore_entity(
6218 &self,
6219 token: &NamespaceToken,
6220 id: Uuid,
6221 ) -> RuntimeResult<Option<(Entity, bool)>> {
6222 let Some(entity) = self
6223 .entities(token)?
6224 .get_entity_including_deleted(id)
6225 .await?
6226 else {
6227 return Ok(None);
6228 };
6229 if entity.namespace != token.namespace().as_str() {
6230 return Ok(None);
6231 }
6232 if let Some(kept_id) = entity.merged_into {
6244 if entity.deleted_at.is_none() {
6245 return Err(live_merged_entity_refused(id, kept_id));
6246 }
6247 return Err(merge_tombstone_restore_refused(id, kept_id));
6248 }
6249 if entity.deleted_at.is_none() {
6250 return Ok(Some((entity, false)));
6251 }
6252 let updated_at =
6253 Utc::now()
6254 .timestamp_micros()
6255 .max(entity.updated_at.checked_add(1).ok_or_else(|| {
6256 RuntimeError::Internal(format!(
6257 "entity {id} updated_at is already at i64::MAX and cannot advance"
6258 ))
6259 })?);
6260 let mut restored = entity;
6261 restored.deleted_at = None;
6262 restored.updated_at = updated_at;
6263 restored.version = restored
6264 .version
6265 .checked_add(1)
6266 .ok_or_else(|| RuntimeError::InvalidInput("entity version overflow".into()))?;
6267 let mut statements = vec![PlanStatement {
6268 statement: SqlStatement {
6269 sql: "UPDATE entities SET deleted_at=NULL, updated_at=?1, version=version+1 \
6270 WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL AND version=?4"
6271 .into(),
6272 params: vec![
6273 SqlValue::Integer(updated_at),
6274 SqlValue::Text(id.to_string()),
6275 SqlValue::Text(token.namespace().as_str().to_owned()),
6276 SqlValue::Integer(restored.version - 1),
6277 ],
6278 label: Some("entity-restore".into()),
6279 },
6280 guard: Some(AffectedRowGuard::exactly(1)),
6281 }];
6282 for statement in khive_db::stores::text::delete_document_statements(
6287 "fts_entities",
6288 &restored.namespace,
6289 id,
6290 )
6291 .into_iter()
6292 .chain(insert_document_statements(
6293 "fts_entities",
6294 &entity_fts_document(&restored),
6295 )) {
6296 statements.push(PlanStatement {
6297 statement,
6298 guard: None,
6299 });
6300 }
6301 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6302 graph_effects: Vec::new(),
6303 target_id: id,
6304 statements,
6305 post_commit: PostCommitEffect::None,
6306 edge_natural_key: None,
6307 idempotent_noop: false,
6308 entity_guard: None,
6309 note_guard: None,
6310 note_vector_purge: None,
6311 note_embedding_inheritance: None,
6312 }));
6313 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6314 Ok(AtomicRunOutcome::Committed { .. }) => {
6315 #[cfg(any(test, feature = "fault-injection"))]
6318 if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
6319 return Err(restore_reindex_failed(
6320 "entity",
6321 id,
6322 RuntimeError::Internal("injected FTS failure".to_string()),
6323 ));
6324 }
6325 self.reindex_entity(token, &restored)
6326 .await
6327 .map_err(|e| restore_reindex_failed("entity", id, e))?;
6328 Ok(Some((restored, true)))
6329 }
6330 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6331 format!("entity restore rolled back: {failure:?}"),
6332 )),
6333 Err(error) => Err(RuntimeError::Storage(error.0)),
6334 }
6335 }
6336
6337 pub async fn restore_note(
6343 &self,
6344 token: &NamespaceToken,
6345 id: Uuid,
6346 ) -> RuntimeResult<Option<(Note, bool)>> {
6347 let Some(note) = self.notes(token)?.get_note_including_deleted(id).await? else {
6348 return Ok(None);
6349 };
6350 if note.namespace != token.namespace().as_str() {
6351 return Ok(None);
6352 }
6353 if note.deleted_at.is_none() {
6354 return Ok(Some((note, false)));
6355 }
6356 if let Some(key) = note.key.as_deref() {
6357 if let Some(holder) = self
6358 .notes(token)?
6359 .get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
6360 .await?
6361 .into_iter()
6362 .find(|holder| holder.id != note.id)
6363 {
6364 return Err(restore_key_conflict(key, &holder));
6365 }
6366 }
6367 let updated_at =
6368 Utc::now()
6369 .timestamp_micros()
6370 .max(note.updated_at.checked_add(1).ok_or_else(|| {
6371 RuntimeError::Internal(format!(
6372 "note {id} updated_at is already at i64::MAX and cannot advance"
6373 ))
6374 })?);
6375 let mut params = vec![
6376 SqlValue::Text("active".into()),
6377 SqlValue::Integer(updated_at),
6378 SqlValue::Text(id.to_string()),
6379 SqlValue::Text(note.namespace.clone()),
6380 SqlValue::Text(note.kind.clone()),
6381 ];
6382 let key_clause = if let Some(key) = note.key.as_deref() {
6383 params.push(SqlValue::Text(key.to_owned()));
6384 format!(
6385 " AND (key IS NULL OR NOT EXISTS (SELECT 1 FROM notes live \
6386 WHERE live.namespace=?4 AND live.kind=?5 AND live.key=?{} \
6387 AND live.deleted_at IS NULL AND live.id != notes.id))",
6388 params.len()
6389 )
6390 } else {
6391 String::new()
6392 };
6393 let mut restored = note.clone();
6394 restored.status = "active".into();
6395 restored.deleted_at = None;
6396 restored.updated_at = updated_at;
6397 restored.version = restored
6398 .version
6399 .checked_add(1)
6400 .ok_or_else(|| RuntimeError::Internal(format!("note {id} version is exhausted")))?;
6401 let mut statements = vec![PlanStatement {
6402 statement: SqlStatement {
6403 sql: format!(
6404 "UPDATE notes SET status=?1, deleted_at=NULL, updated_at=?2 \
6405 WHERE id=?3 AND namespace=?4 AND kind=?5 AND deleted_at IS NOT NULL{key_clause}"
6406 ),
6407 params,
6408 label: Some("note-restore".into()),
6409 },
6410 guard: Some(AffectedRowGuard::exactly(1)),
6411 }];
6412 for statement in
6414 khive_db::stores::text::delete_document_statements("fts_notes", &restored.namespace, id)
6415 .into_iter()
6416 .chain(insert_document_statements(
6417 "fts_notes",
6418 ¬e_fts_document(&restored),
6419 ))
6420 {
6421 statements.push(PlanStatement {
6422 statement,
6423 guard: None,
6424 });
6425 }
6426 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6427 graph_effects: Vec::new(),
6428 target_id: id,
6429 statements,
6430 post_commit: PostCommitEffect::None,
6431 edge_natural_key: None,
6432 idempotent_noop: false,
6433 entity_guard: None,
6434 note_guard: None,
6435 note_vector_purge: None,
6436 note_embedding_inheritance: None,
6437 }));
6438 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6439 Ok(AtomicRunOutcome::Committed { .. }) => {
6440 #[cfg(any(test, feature = "fault-injection"))]
6441 if consume_fault(&FTS_FAIL_NS, &restored.namespace) {
6442 return Err(restore_reindex_failed(
6443 "note",
6444 id,
6445 RuntimeError::Internal("injected FTS failure".to_string()),
6446 ));
6447 }
6448 let reindexed = self.reindex_note_with_report(token, &restored).await;
6449 let report = reindexed.map_err(|e| restore_reindex_failed("note", id, e))?;
6450 let degradations = report.post_commit_degradations();
6451 legacy_post_commit_result("restore_note", id, Some((restored, true)), degradations)
6452 }
6453 Ok(AtomicRunOutcome::RolledBack {
6454 failure: AtomicOpFailure::GuardFailed { .. },
6455 ..
6456 }) => {
6457 if let Some(key) = note.key.as_deref() {
6458 if let Some(holder) = self
6459 .notes(token)?
6460 .get_live_notes_by_key(¬e.namespace, key, Some(¬e.kind))
6461 .await?
6462 .into_iter()
6463 .find(|holder| holder.id != note.id)
6464 {
6465 return Err(restore_key_conflict(key, &holder));
6466 }
6467 }
6468 Err(RuntimeError::NotFound(format!(
6469 "note {id} is no longer a caller-owned tombstone"
6470 )))
6471 }
6472 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6473 format!("note restore rolled back: {failure:?}"),
6474 )),
6475 Err(error) => Err(RuntimeError::Storage(error.0)),
6476 }
6477 }
6478
6479 pub async fn restore_edge(
6481 &self,
6482 token: &NamespaceToken,
6483 id: Uuid,
6484 ) -> RuntimeResult<Option<(Edge, bool)>> {
6485 let Some(edge) = self.get_edge_including_deleted(token, id).await? else {
6486 return Ok(None);
6487 };
6488 if edge.namespace != token.namespace().as_str() {
6489 return Ok(None);
6490 }
6491 if edge.deleted_at.is_none() {
6492 return Ok(Some((edge, false)));
6493 }
6494 let updated_at = Utc::now();
6495 let plan = AtomicOpPlan::Update(Box::new(UpdatePlan {
6496 graph_effects: Vec::new(),
6497 target_id: id,
6498 statements: vec![PlanStatement {
6499 statement: SqlStatement {
6500 sql: "UPDATE graph_edges SET deleted_at=NULL, updated_at=?1 \
6501 WHERE id=?2 AND namespace=?3 AND deleted_at IS NOT NULL"
6502 .into(),
6503 params: vec![
6504 SqlValue::Integer(updated_at.timestamp_micros()),
6505 SqlValue::Text(id.to_string()),
6506 SqlValue::Text(edge.namespace.clone()),
6507 ],
6508 label: Some("edge-restore".into()),
6509 },
6510 guard: Some(AffectedRowGuard::exactly(1)),
6511 }],
6512 post_commit: PostCommitEffect::None,
6513 edge_natural_key: None,
6514 idempotent_noop: false,
6515 entity_guard: None,
6516 note_guard: None,
6517 note_vector_purge: None,
6518 note_embedding_inheritance: None,
6519 }));
6520 match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
6521 Ok(AtomicRunOutcome::Committed { .. }) => {
6522 let mut restored = edge;
6523 restored.deleted_at = None;
6524 restored.updated_at = updated_at;
6525 Ok(Some((restored, true)))
6526 }
6527 Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
6528 format!("edge restore rolled back: {failure:?}"),
6529 )),
6530 Err(error) => Err(RuntimeError::Storage(error.0)),
6531 }
6532 }
6533
6534 pub async fn delete_note(
6545 &self,
6546 token: &NamespaceToken,
6547 id: Uuid,
6548 hard: bool,
6549 ) -> RuntimeResult<bool> {
6550 let (deleted, degradations) = self
6551 .delete_note_with_post_commit_report(token, id, hard)
6552 .await?;
6553 legacy_post_commit_result("delete_note", id, deleted, degradations)
6554 }
6555
6556 pub async fn delete_note_with_post_commit_report(
6560 &self,
6561 token: &NamespaceToken,
6562 id: Uuid,
6563 hard: bool,
6564 ) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
6565 let note_store = self.notes(token)?;
6566 let note = if hard {
6567 match note_store.get_note_including_deleted(id).await? {
6568 Some(n) => n,
6569 None => return Ok((false, Vec::new())),
6570 }
6571 } else {
6572 match note_store.get_note(id).await? {
6573 Some(n) => n,
6574 None => return Ok((false, Vec::new())),
6575 }
6576 };
6577 if let Some(error) = self.stream_member_error(¬e).await? {
6578 return Err(error);
6579 }
6580 let mode = if hard {
6581 DeleteMode::Hard
6582 } else {
6583 DeleteMode::Soft
6584 };
6585
6586 let record_tok = token.with_namespace(
6588 khive_types::Namespace::parse(¬e.namespace)
6589 .map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
6590 );
6591 let record_ns = note.namespace.clone();
6592 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
6593
6594 let deleted = if hard {
6599 self.atomic_hard_delete_with_edge_purge(
6600 note_hard_delete_statement(id),
6601 id,
6602 &record_ns,
6603 &actor,
6604 SubstrateKind::Note,
6605 )
6606 .await?
6607 } else {
6608 note_store.delete_note(id, mode).await?
6609 };
6610 let mut degradations = Vec::new();
6611 if deleted {
6612 let fts_result = match self.text_for_notes(&record_tok) {
6613 Ok(store) => store
6614 .delete_document(&record_ns, id)
6615 .await
6616 .map_err(RuntimeError::from),
6617 Err(error) => Err(error),
6618 };
6619 if let Err(error) = fts_result {
6620 record_post_commit_degradation(
6621 &mut degradations,
6622 "delete_note",
6623 id,
6624 "fts_cleanup",
6625 error,
6626 );
6627 }
6628 for model_name in self.registered_embedding_model_names() {
6630 let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
6631 Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
6632 Err(error) => Err(error),
6633 };
6634 if let Err(error) = vector_result {
6635 record_post_commit_degradation(
6636 &mut degradations,
6637 "delete_note",
6638 id,
6639 "vector_cleanup",
6640 format!("{model_name}: {error}"),
6641 );
6642 }
6643 }
6644 let event = khive_storage::event::Event::new(
6645 record_ns.clone(),
6646 "delete",
6647 EventKind::NoteDeleted,
6648 SubstrateKind::Note,
6649 "",
6650 )
6651 .with_target(id)
6652 .with_payload(serde_json::json!({"id": id, "namespace": record_ns, "hard": hard}));
6653 let event_result = match self.events(&record_tok) {
6654 Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
6655 Err(error) => Err(error),
6656 };
6657 if let Err(error) = event_result {
6658 record_post_commit_degradation(
6659 &mut degradations,
6660 "delete_note",
6661 id,
6662 "event_append",
6663 error,
6664 );
6665 }
6666 self.fire_note_mutation_hook(¬e.kind, id).await;
6673 }
6674 Ok((deleted, degradations))
6675 }
6676
6677 pub async fn delete_note_row_first_for_compensation(
6696 &self,
6697 token: &NamespaceToken,
6698 id: Uuid,
6699 ) -> RuntimeResult<()> {
6700 let note_store = self.notes(token)?;
6701 let Some(note) = note_store.get_note_including_deleted(id).await? else {
6702 return Ok(());
6703 };
6704 let record_tok = NamespaceToken::for_namespace(
6705 khive_types::Namespace::parse(¬e.namespace)
6706 .map_err(|e| RuntimeError::Internal(format!("note namespace invalid: {e}")))?,
6707 );
6708 let record_ns = note.namespace.clone();
6709
6710 note_store.delete_note(id, DeleteMode::Hard).await?;
6712
6713 #[cfg(any(test, feature = "fault-injection"))]
6714 {
6715 let armed = ROLLBACK_CLEANUP_FAIL_NS.lock().unwrap().take();
6716 if armed.as_deref() == Some(record_ns.as_str()) {
6717 return Err(RuntimeError::Internal(
6718 "row removed but compensation cleanup failed: injected=true".to_string(),
6719 ));
6720 }
6721 }
6722
6723 let mut cleanup_errors = Vec::new();
6724 if let Err(e) = self.graph(&record_tok)?.purge_incident_edges(id).await {
6725 cleanup_errors.push(format!("graph={e}"));
6726 }
6727 if let Err(e) = self
6728 .text_for_notes(&record_tok)?
6729 .delete_document(&record_ns, id)
6730 .await
6731 {
6732 cleanup_errors.push(format!("fts={e}"));
6733 }
6734 for model_name in self.registered_embedding_model_names() {
6735 if let Err(e) = self
6736 .vectors_for_model(&record_tok, &model_name)?
6737 .delete(id)
6738 .await
6739 {
6740 cleanup_errors.push(format!("vector[{model_name}]={e}"));
6741 }
6742 }
6743 if cleanup_errors.is_empty() {
6744 Ok(())
6745 } else {
6746 Err(RuntimeError::Internal(format!(
6747 "row removed but compensation cleanup failed: {}",
6748 cleanup_errors.join("; ")
6749 )))
6750 }
6751 }
6752}
6753
6754#[derive(Clone, Debug, Serialize)]
6756pub struct QueryResult {
6757 pub rows: Vec<SqlRow>,
6758 #[serde(skip_serializing_if = "Vec::is_empty")]
6759 pub warnings: Vec<String>,
6760 pub offset: usize,
6762 pub page_size: usize,
6764 pub has_more: bool,
6766 #[serde(skip_serializing_if = "Option::is_none")]
6768 pub next_offset: Option<usize>,
6769 pub truncated: bool,
6771}
6772
6773#[derive(Debug)]
6775enum SymmetricEdgeUpdateOutcome {
6776 Absorbed(String),
6780 Updated,
6782 Stale,
6787}
6788
6789impl KhiveRuntime {
6790 pub async fn query(&self, token: &NamespaceToken, query: &str) -> RuntimeResult<Vec<SqlRow>> {
6798 Ok(self
6799 .query_with_metadata(token, query, khive_query::CompileOptions::default())
6800 .await?
6801 .rows)
6802 }
6803
6804 pub async fn query_with_metadata(
6806 &self,
6807 token: &NamespaceToken,
6808 query: &str,
6809 mut opts: khive_query::CompileOptions,
6810 ) -> RuntimeResult<QueryResult> {
6811 use khive_query::QueryValue;
6812 use khive_storage::types::SqlValue;
6813
6814 let (language, ast) = khive_query::language::parse_auto_with_language(query)?;
6815 if opts.max_limit == 0 {
6816 return Err(RuntimeError::InvalidInput(
6817 "query page size must be at least 1".into(),
6818 ));
6819 }
6820 let offset = ast.offset;
6821 let page_size = ast.limit.unwrap_or(opts.max_limit).min(opts.max_limit);
6822 opts.scopes = token
6823 .visible_namespaces()
6824 .iter()
6825 .map(|ns| ns.as_str().to_string())
6826 .collect();
6827 let compiled = khive_query::compile(&ast, &opts)?;
6828 let mut warnings = compiled.warnings;
6829 let truncation_check = compiled.truncation_check;
6830
6831 warnings.extend(self.with_pack_edge_rules(|pack_rules| {
6832 static_impossible_edge_pattern_warnings(language, &ast.pattern, pack_rules)
6833 }));
6834
6835 let params: Vec<SqlValue> = compiled
6838 .params
6839 .into_iter()
6840 .map(|qv| match qv {
6841 QueryValue::Null => SqlValue::Null,
6842 QueryValue::Integer(n) => SqlValue::Integer(n),
6843 QueryValue::Float(f) => SqlValue::Float(f),
6844 QueryValue::Text(s) => SqlValue::Text(s),
6845 QueryValue::Blob(b) => SqlValue::Blob(b),
6846 })
6847 .collect();
6848
6849 let mut reader = self.sql().reader().await?;
6850 let stmt = SqlStatement {
6851 sql: compiled.sql,
6852 params,
6853 label: None,
6854 };
6855 let mut rows = reader.query_all(stmt).await?;
6856
6857 let mut truncated = false;
6861 if let Some(check) = truncation_check {
6862 if rows.len() > check.max_limit {
6863 rows.truncate(check.max_limit);
6864 truncated = true;
6865 }
6866 }
6867
6868 let next_offset = if truncated && language == khive_query::QueryLanguage::Gql {
6869 let next = offset.checked_add(rows.len()).ok_or_else(|| {
6870 RuntimeError::InvalidInput("GQL next_offset exceeds usize::MAX".into())
6871 })?;
6872 if next == offset {
6873 return Err(RuntimeError::InvalidInput(
6874 "query page did not advance; page size must be at least 1".into(),
6875 ));
6876 }
6877 i64::try_from(next).map_err(|_| {
6878 RuntimeError::InvalidInput("GQL next_offset exceeds i64::MAX".into())
6879 })?;
6880 Some(next)
6881 } else {
6882 None
6883 };
6884
6885 if truncated {
6886 let Some(check) = truncation_check else {
6887 return Err(RuntimeError::Internal(
6888 "truncated query result is missing sentinel metadata".into(),
6889 ));
6890 };
6891 let bound = match check.requested_limit {
6892 Some(requested) => {
6893 format!("requested query LIMIT {requested} exceeds the effective page size")
6894 }
6895 None => "the query has no explicit LIMIT".to_string(),
6896 };
6897 let warning = match language {
6898 khive_query::QueryLanguage::Gql => {
6899 let Some(next) = next_offset else {
6900 return Err(RuntimeError::Internal(
6901 "truncated GQL result is missing its continuation offset".into(),
6902 ));
6903 };
6904 format!(
6905 "result page capped at {} rows because {bound}; more matches exist. \
6906 Continue the same GQL query with `SKIP {next}` (the machine-readable \
6907 `next_offset`) and keep the same page size.",
6908 check.max_limit
6909 )
6910 }
6911 khive_query::QueryLanguage::Sparql => format!(
6912 "result page capped at {} rows because {bound}; more matches exist. \
6913 SPARQL OFFSET paging is not part of the supported dialect.",
6914 check.max_limit
6915 ),
6916 };
6917 warnings.push(warning);
6918 }
6919
6920 Ok(QueryResult {
6921 rows,
6922 warnings,
6923 offset,
6924 page_size,
6925 has_more: truncated,
6926 next_offset,
6927 truncated,
6928 })
6929 }
6930
6931 pub async fn delete_entity(
6940 &self,
6941 token: &NamespaceToken,
6942 id: Uuid,
6943 hard: bool,
6944 ) -> RuntimeResult<bool> {
6945 let (deleted, degradations) = self
6946 .delete_entity_with_post_commit_report(token, id, hard)
6947 .await?;
6948 legacy_post_commit_result("delete_entity", id, deleted, degradations)
6949 }
6950
6951 pub async fn delete_entity_with_post_commit_report(
6953 &self,
6954 token: &NamespaceToken,
6955 id: Uuid,
6956 hard: bool,
6957 ) -> RuntimeResult<(bool, Vec<PostCommitDegradation>)> {
6958 let entity = if hard {
6959 match self
6960 .entities(token)?
6961 .get_entity_including_deleted(id)
6962 .await?
6963 {
6964 Some(e) => e,
6965 None => return Ok((false, Vec::new())),
6966 }
6967 } else {
6968 match self.entities(token)?.get_entity(id).await? {
6969 Some(e) => e,
6970 None => return Ok((false, Vec::new())),
6971 }
6972 };
6973 let mode = if hard {
6974 DeleteMode::Hard
6975 } else {
6976 DeleteMode::Soft
6977 };
6978
6979 let record_tok = token.with_namespace(
6981 khive_types::Namespace::parse(&entity.namespace)
6982 .map_err(|e| RuntimeError::Internal(format!("entity namespace invalid: {e}")))?,
6983 );
6984 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
6985
6986 let deleted = if hard {
6991 self.atomic_hard_delete_with_edge_purge(
6993 entity_hard_delete_statement(id),
6994 id,
6995 &entity.namespace,
6996 &actor,
6997 SubstrateKind::Entity,
6998 )
6999 .await?
7000 } else {
7001 self.entities(token)?.delete_entity(id, mode).await?
7002 };
7003 let mut degradations = Vec::new();
7004 if deleted {
7005 let ns = entity.namespace.clone();
7006 let fts_result = match self.text(&record_tok) {
7007 Ok(store) => store
7008 .delete_document(&ns, id)
7009 .await
7010 .map_err(RuntimeError::from),
7011 Err(error) => Err(error),
7012 };
7013 if let Err(error) = fts_result {
7014 record_post_commit_degradation(
7015 &mut degradations,
7016 "delete_entity",
7017 id,
7018 "fts_cleanup",
7019 error,
7020 );
7021 }
7022 for model_name in self.registered_embedding_model_names() {
7023 let vector_result = match self.vectors_for_model(&record_tok, &model_name) {
7024 Ok(store) => store.delete(id).await.map_err(RuntimeError::from),
7025 Err(error) => Err(error),
7026 };
7027 if let Err(error) = vector_result {
7028 record_post_commit_degradation(
7029 &mut degradations,
7030 "delete_entity",
7031 id,
7032 "vector_cleanup",
7033 format!("{model_name}: {error}"),
7034 );
7035 }
7036 }
7037 let event = khive_storage::event::Event::new(
7038 ns.clone(),
7039 "delete",
7040 EventKind::EntityDeleted,
7041 SubstrateKind::Entity,
7042 "",
7043 )
7044 .with_target(id)
7045 .with_payload(serde_json::json!({"id": id, "namespace": ns, "hard": hard}));
7046 let event_result = match self.events(&record_tok) {
7047 Ok(store) => store.append_event(event).await.map_err(RuntimeError::from),
7048 Err(error) => Err(error),
7049 };
7050 if let Err(error) = event_result {
7051 record_post_commit_degradation(
7052 &mut degradations,
7053 "delete_entity",
7054 id,
7055 "event_append",
7056 error,
7057 );
7058 }
7059 }
7060 Ok((deleted, degradations))
7061 }
7062
7063 pub(crate) async fn delete_entity_attachments_on_core(&self, id: Uuid) -> RuntimeResult<bool> {
7064 let core = self.core();
7065 drop(core.attachments()?);
7066 let statement = khive_db::stores::attachment::delete_record_attachments_statement(
7067 id,
7068 AttachmentSubstrate::Entity,
7069 );
7070 Ok(core.sql().writer().await?.execute(statement).await? > 0)
7071 }
7072
7073 pub async fn count_entities(
7075 &self,
7076 token: &NamespaceToken,
7077 kind: Option<&str>,
7078 ) -> RuntimeResult<u64> {
7079 let ns_strs: Vec<String> = token
7080 .visible_namespaces()
7081 .iter()
7082 .map(|ns| ns.as_str().to_owned())
7083 .collect();
7084 let filter = EntityFilter {
7085 kinds: match kind {
7086 Some(k) => vec![k.to_string()],
7087 None => vec![],
7088 },
7089 namespaces: ns_strs,
7090 ..Default::default()
7091 };
7092 Ok(self
7093 .entities(token)?
7094 .count_entities(token.namespace().as_str(), filter)
7095 .await?)
7096 }
7097
7098 pub async fn entity_stats_counts(
7100 &self,
7101 token: &NamespaceToken,
7102 ) -> RuntimeResult<EntityStatsCounts> {
7103 entity_stats_counts(self.entities(token)?.as_ref(), token).await
7104 }
7105
7106 pub async fn get_edge(
7113 &self,
7114 _token: &NamespaceToken,
7115 edge_id: Uuid,
7116 ) -> RuntimeResult<Option<Edge>> {
7117 let mut reader = self.sql().reader().await?;
7118 let record_ns = reader
7119 .query_scalar(SqlStatement {
7120 sql: "SELECT namespace FROM graph_edges \
7121 WHERE id = ?1 AND deleted_at IS NULL LIMIT 1"
7122 .into(),
7123 params: vec![SqlValue::Text(edge_id.to_string())],
7124 label: Some("get_edge_namespace".into()),
7125 })
7126 .await?;
7127
7128 let Some(SqlValue::Text(record_ns)) = record_ns else {
7129 return Ok(None);
7130 };
7131 let record_tok = NamespaceToken::for_namespace(
7134 khive_types::Namespace::parse(&record_ns)
7135 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
7136 );
7137 Ok(self
7138 .graph(&record_tok)?
7139 .get_edge(LinkId::from(edge_id))
7140 .await?)
7141 }
7142
7143 pub async fn get_edges_by_id(
7152 &self,
7153 _token: &NamespaceToken,
7154 ids: &[Uuid],
7155 ) -> RuntimeResult<Vec<Option<Edge>>> {
7156 let mut edges = Vec::with_capacity(ids.len());
7157 for chunk in ids.chunks(900) {
7158 let window = self.prepare_edge_read_window(chunk).await?;
7159 edges.extend(
7160 Self::hydrate_edge_read_window(chunk, window, |record_token| {
7161 self.graph(record_token)
7162 })
7163 .await?,
7164 );
7165 }
7166 Ok(edges)
7167 }
7168
7169 async fn prepare_edge_read_window(&self, ids: &[Uuid]) -> RuntimeResult<EdgeReadWindow> {
7170 let placeholders = (1..=ids.len())
7171 .map(|index| format!("?{index}"))
7172 .collect::<Vec<_>>()
7173 .join(",");
7174 let mut reader = self.sql().reader().await?;
7175 let rows = reader
7176 .query_all(SqlStatement {
7177 sql: format!(
7178 "SELECT id, namespace FROM graph_edges WHERE id IN ({placeholders}) AND deleted_at IS NULL"
7179 ),
7180 params: ids.iter().map(|id| SqlValue::Text(id.to_string())).collect(),
7181 label: Some("get_edge_namespace".into()),
7182 })
7183 .await?;
7184 let mut namespaces = HashMap::with_capacity(rows.len());
7185 for row in rows {
7186 let Some(SqlValue::Text(id)) = row.columns.first().map(|column| &column.value) else {
7187 return Err(RuntimeError::Internal(
7188 "edge namespace lookup returned an invalid id".into(),
7189 ));
7190 };
7191 let id = Uuid::parse_str(id).map_err(|e| {
7192 RuntimeError::Internal(format!("edge namespace lookup returned an invalid id: {e}"))
7193 })?;
7194 let value = row
7195 .columns
7196 .get(1)
7197 .map(|column| column.value.clone())
7198 .unwrap_or(SqlValue::Null);
7199 namespaces.insert(id, value);
7200 }
7201 let mut window = EdgeReadWindow {
7202 outcomes: (0..ids.len()).map(|_| Some(Ok(None))).collect(),
7203 groups: Vec::new(),
7204 };
7205 let mut group_indices = HashMap::new();
7206 for (index, id) in ids.iter().enumerate() {
7207 let Some(SqlValue::Text(record_ns)) = namespaces.get(id) else {
7208 continue;
7209 };
7210 match khive_types::Namespace::parse(record_ns) {
7211 Ok(namespace) => {
7212 let next_group = window.groups.len();
7213 let group = *group_indices.entry(record_ns.clone()).or_insert(next_group);
7214 if group == next_group {
7215 window.groups.push((namespace, Vec::new()));
7216 }
7217 window.groups[group].1.push(index);
7218 window.outcomes[index] = None;
7219 }
7220 Err(error) => {
7221 window.outcomes[index] = Some(Err(RuntimeError::Internal(format!(
7222 "edge namespace invalid: {error}"
7223 ))));
7224 }
7225 }
7226 }
7227 Ok(window)
7228 }
7229
7230 async fn hydrate_edge_read_window<F>(
7231 ids: &[Uuid],
7232 mut window: EdgeReadWindow,
7233 mut graph: F,
7234 ) -> RuntimeResult<Vec<Option<Edge>>>
7235 where
7236 F: FnMut(&NamespaceToken) -> RuntimeResult<std::sync::Arc<dyn khive_storage::GraphStore>>,
7237 {
7238 for (namespace, indices) in window.groups {
7239 let record_token = NamespaceToken::for_namespace(namespace);
7240 let group_ids: Vec<LinkId> = indices
7241 .iter()
7242 .map(|&index| LinkId::from(ids[index]))
7243 .collect();
7244 let outcomes = match graph(&record_token) {
7245 Ok(store) => store
7246 .get_edge_read_outcomes(&group_ids)
7247 .await
7248 .map_err(RuntimeError::from),
7249 Err(error) => Err(error),
7250 };
7251 match outcomes {
7252 Ok(outcomes) if outcomes.len() == indices.len() => {
7253 for (index, outcome) in indices.into_iter().zip(outcomes) {
7254 window.outcomes[index] = Some(outcome.map_err(RuntimeError::from));
7255 }
7256 }
7257 Ok(_) => {
7258 window.outcomes[indices[0]] = Some(Err(RuntimeError::Internal(
7259 "edge batch returned an invalid outcome count".into(),
7260 )));
7261 }
7262 Err(error) => {
7263 window.outcomes[indices[0]] = Some(Err(error));
7264 }
7265 }
7266 }
7267 window
7268 .outcomes
7269 .into_iter()
7270 .map(|outcome| {
7271 outcome.unwrap_or_else(|| {
7272 Err(RuntimeError::Internal(
7273 "edge batch omitted an input outcome".into(),
7274 ))
7275 })
7276 })
7277 .collect()
7278 }
7279
7280 pub async fn get_edge_visible(
7285 &self,
7286 token: &NamespaceToken,
7287 edge_id: Uuid,
7288 ) -> RuntimeResult<Option<Edge>> {
7289 self.get_edge(token, edge_id).await
7290 }
7291
7292 pub async fn get_edge_including_deleted(
7298 &self,
7299 _token: &NamespaceToken,
7300 edge_id: Uuid,
7301 ) -> RuntimeResult<Option<Edge>> {
7302 let mut reader = self.sql().reader().await?;
7303 let record_ns = reader
7304 .query_scalar(SqlStatement {
7305 sql: "SELECT namespace FROM graph_edges WHERE id = ?1 LIMIT 1".into(),
7306 params: vec![SqlValue::Text(edge_id.to_string())],
7307 label: Some("get_edge_including_deleted_namespace".into()),
7308 })
7309 .await?;
7310
7311 let Some(SqlValue::Text(record_ns)) = record_ns else {
7312 return Ok(None);
7313 };
7314 let record_tok = NamespaceToken::for_namespace(
7316 khive_types::Namespace::parse(&record_ns)
7317 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
7318 );
7319 Ok(self
7320 .graph(&record_tok)?
7321 .get_edge_including_deleted(LinkId::from(edge_id))
7322 .await?)
7323 }
7324
7325 pub async fn get_edge_by_natural_key_including_deleted(
7339 &self,
7340 token: &NamespaceToken,
7341 namespace: &str,
7342 source_id: Uuid,
7343 target_id: Uuid,
7344 relation: EdgeRelation,
7345 ) -> RuntimeResult<Option<Edge>> {
7346 Ok(self
7347 .graph(token)?
7348 .get_edge_by_natural_key_including_deleted(namespace, source_id, target_id, relation)
7349 .await?)
7350 }
7351
7352 pub const EDGE_LIST_MAX_LIMIT: u32 = 1000;
7357
7358 pub async fn list_edges(
7363 &self,
7364 token: &NamespaceToken,
7365 filter: crate::curation::EdgeListFilter,
7366 limit: u32,
7367 offset: u32,
7368 ) -> RuntimeResult<Vec<Edge>> {
7369 let limit = limit.min(Self::EDGE_LIST_MAX_LIMIT);
7370 let visible = token.visible_namespaces();
7371
7372 if let [ns] = visible {
7375 let temp = NamespaceToken::for_namespace(ns.clone());
7376 let page = self
7377 .graph(&temp)?
7378 .query_edges(
7379 filter.into(),
7380 vec![SortOrder {
7381 field: EdgeSortField::CreatedAt,
7382 direction: khive_storage::types::SortDirection::Asc,
7383 }],
7384 PageRequest {
7385 offset: offset.into(),
7386 limit,
7387 },
7388 )
7389 .await?;
7390 return Ok(page.items);
7391 }
7392
7393 let ns_strs: Vec<String> = visible.iter().map(|ns| ns.as_str().to_owned()).collect();
7400 let sort = vec![SortOrder {
7401 field: EdgeSortField::CreatedAt,
7402 direction: khive_storage::types::SortDirection::Asc,
7403 }];
7404 let graph = self.graph(token)?;
7405 match graph
7406 .query_edges_in_namespaces(
7407 &ns_strs,
7408 filter.clone().into(),
7409 sort.clone(),
7410 PageRequest {
7411 offset: offset.into(),
7412 limit,
7413 },
7414 )
7415 .await
7416 {
7417 Ok(page) => Ok(page.items),
7418 Err(khive_storage::StorageError::Unsupported { operation, .. })
7419 if operation == "query_edges_in_namespaces" =>
7420 {
7421 let fetch_limit = offset.saturating_add(limit);
7434 let mut namespace_prefixes = Vec::new();
7435 for ns in visible {
7436 let temp = NamespaceToken::for_namespace(ns.clone());
7437 let page = self
7438 .graph(&temp)?
7439 .query_edges(
7440 filter.clone().into(),
7441 sort.clone(),
7442 PageRequest {
7443 offset: 0,
7444 limit: fetch_limit,
7445 },
7446 )
7447 .await?;
7448 namespace_prefixes.push(page.items);
7449 }
7450 Ok(Self::merge_paged_namespace_edges(
7451 namespace_prefixes,
7452 offset,
7453 limit,
7454 ))
7455 }
7456 Err(error) => Err(error.into()),
7457 }
7458 }
7459
7460 fn merge_paged_namespace_edges(
7475 namespace_prefixes: Vec<Vec<Edge>>,
7476 offset: u32,
7477 limit: u32,
7478 ) -> Vec<Edge> {
7479 let mut results: Vec<Edge> = namespace_prefixes.into_iter().flatten().collect();
7480 results.sort_by_key(|e| (e.created_at, Uuid::from(e.id)));
7481 let start = (offset as usize).min(results.len());
7482 let end = (start + limit as usize).min(results.len());
7483 results[start..end].to_vec()
7484 }
7485
7486 pub async fn list_edges_after(
7498 &self,
7499 token: &NamespaceToken,
7500 filter: crate::curation::EdgeListFilter,
7501 after: Option<Uuid>,
7502 limit: u32,
7503 ) -> RuntimeResult<(Vec<Edge>, Option<Uuid>)> {
7504 let limit = limit.clamp(1, Self::EDGE_LIST_MAX_LIMIT);
7505 let visible = token.visible_namespaces();
7506 let limit_usize = limit as usize;
7507 let cursor_store = self.graph(token)?;
7508 let after = match after {
7509 Some(id) => {
7510 let edge = self
7511 .get_edge_including_deleted(token, id)
7512 .await?
7513 .ok_or_else(|| RuntimeError::NotFound(format!("edge cursor {id}")))?;
7514 Self::ensure_namespace_visible(&edge.namespace, token)?;
7515 let sequence = cursor_store.edge_sequence(id).await?.ok_or_else(|| {
7516 RuntimeError::Internal(format!(
7517 "edge cursor {id} has no insertion-sequence ledger row"
7518 ))
7519 })?;
7520 Some(SeekCursor { sequence, id })
7521 }
7522 None => None,
7523 };
7524
7525 if let [ns] = visible {
7526 let temp = NamespaceToken::for_namespace(ns.clone());
7527 let page = self
7528 .graph(&temp)?
7529 .query_edges_sequence_after(filter.into(), after, limit)
7530 .await?;
7531 return Ok((page.items, page.next_after.map(|cursor| cursor.id)));
7532 }
7533
7534 let probe_limit = limit.saturating_add(1);
7538 let mut results = Vec::new();
7539 for ns in visible {
7540 let temp = NamespaceToken::for_namespace(ns.clone());
7541 let page = self
7542 .graph(&temp)?
7543 .query_edges_sequence_after(filter.clone().into(), after, probe_limit)
7544 .await?;
7545 results.extend(page.items);
7546 }
7547 let ids = results
7548 .iter()
7549 .map(|edge| Uuid::from(edge.id))
7550 .collect::<Vec<_>>();
7551 let sequences = cursor_store
7552 .edge_sequences(&ids)
7553 .await?
7554 .into_iter()
7555 .collect::<HashMap<_, _>>();
7556 if let Some(missing) = ids.iter().find(|id| !sequences.contains_key(id)) {
7557 return Err(RuntimeError::Internal(format!(
7558 "edge {missing} has no insertion-sequence ledger row"
7559 )));
7560 }
7561 results.sort_by_key(|edge| {
7562 let id = Uuid::from(edge.id);
7563 (sequences[&id], id)
7564 });
7565 results.dedup_by_key(|e| Uuid::from(e.id));
7566 let has_more = results.len() > limit_usize;
7567 if has_more {
7568 results.truncate(limit_usize);
7569 }
7570 let next_after = if has_more {
7571 results.last().map(|e| Uuid::from(e.id))
7572 } else {
7573 None
7574 };
7575 Ok((results, next_after))
7576 }
7577
7578 pub async fn count_edges_by_relation(
7582 &self,
7583 token: &NamespaceToken,
7584 ) -> RuntimeResult<std::collections::HashMap<String, u64>> {
7585 let namespaces: Vec<String> = token
7586 .visible_namespaces()
7587 .iter()
7588 .map(|namespace| namespace.as_str().to_owned())
7589 .collect();
7590 let graph = self.graph(token)?;
7591 let counts = match graph
7592 .count_edges_by_relation_in_namespaces(&namespaces)
7593 .await
7594 {
7595 Ok(counts) => counts,
7596 Err(khive_storage::StorageError::Unsupported { operation, .. })
7597 if operation == "count_edges_by_relation_in_namespaces" =>
7598 {
7599 let mut totals = HashMap::new();
7600 for namespace in token.visible_namespaces() {
7601 let scoped = NamespaceToken::for_namespace(namespace.clone());
7602 for (relation, count) in self.graph(&scoped)?.count_edges_by_relation().await? {
7603 *totals.entry(relation).or_insert(0) += count;
7604 }
7605 }
7606 return Ok(totals
7607 .into_iter()
7608 .map(|(relation, count)| (relation.to_string(), count))
7609 .collect());
7610 }
7611 Err(error) => return Err(error.into()),
7612 };
7613 Ok(counts
7614 .into_iter()
7615 .map(|(relation, count)| (relation.to_string(), count))
7616 .collect())
7617 }
7618
7619 pub async fn count_edges_by_endpoint_base(
7628 &self,
7629 token: &NamespaceToken,
7630 ) -> RuntimeResult<khive_storage::types::EdgeEndpointBaseCounts> {
7631 use khive_storage::types::EdgeEndpointBaseCounts;
7632
7633 let namespaces: Vec<String> = token
7634 .visible_namespaces()
7635 .iter()
7636 .map(|namespace| namespace.as_str().to_owned())
7637 .collect();
7638 let graph = self.graph(token)?;
7639 match graph
7640 .count_edges_by_endpoint_base_in_namespaces(&namespaces)
7641 .await
7642 {
7643 Ok(counts) => Ok(counts),
7644 Err(khive_storage::StorageError::Unsupported { operation, .. })
7645 if operation == "count_edges_by_endpoint_base_in_namespaces"
7646 || operation == "count_edges_by_endpoint_base" =>
7647 {
7648 let mut totals = EdgeEndpointBaseCounts::default();
7649 for namespace in token.visible_namespaces() {
7650 let scoped = NamespaceToken::for_namespace(namespace.clone());
7651 let counts = self.graph(&scoped)?.count_edges_by_endpoint_base().await?;
7652 totals.entity_entity =
7653 totals.entity_entity.saturating_add(counts.entity_entity);
7654 totals.entity_note = totals.entity_note.saturating_add(counts.entity_note);
7655 totals.note_entity = totals.note_entity.saturating_add(counts.note_entity);
7656 totals.note_note = totals.note_note.saturating_add(counts.note_note);
7657 totals.unresolved = totals.unresolved.saturating_add(counts.unresolved);
7658 }
7659 Ok(totals)
7660 }
7661 Err(error) => Err(error.into()),
7662 }
7663 }
7664
7665 #[allow(clippy::too_many_arguments)]
7677 fn update_edge_symmetric_dml(
7678 conn: &rusqlite::Connection,
7679 ns: &str,
7680 edge_id_str: &str,
7681 canon_src_str: &str,
7682 canon_tgt_str: &str,
7683 relation_str: &str,
7684 weight: f64,
7685 metadata: Option<String>,
7686 expected_updated_at_micros: i64,
7687 expected_deleted_at_micros: Option<i64>,
7688 ) -> Result<SymmetricEdgeUpdateOutcome, SqliteError> {
7689 let minimum_updated_at_micros =
7703 expected_updated_at_micros.checked_add(1).ok_or_else(|| {
7704 SqliteError::InvalidData(format!(
7705 "update_edge: edge {edge_id_str} updated_at is already at i64::MAX \
7706 and cannot advance"
7707 ))
7708 })?;
7709 let now_ts = chrono::Utc::now()
7710 .timestamp_micros()
7711 .max(minimum_updated_at_micros);
7712
7713 let conflict_id: Option<String> = conn
7716 .query_row(
7717 khive_db::stores::graph::EDGE_SYMMETRIC_CONFLICT_PROBE_SQL,
7718 rusqlite::params![
7719 &ns,
7720 &canon_src_str,
7721 &canon_tgt_str,
7722 &relation_str,
7723 &edge_id_str
7724 ],
7725 |row| row.get(0),
7726 )
7727 .optional()
7728 .map_err(SqliteError::Rusqlite)?;
7729
7730 if let Some(existing_id) = conflict_id {
7731 let affected = conn
7750 .execute(
7751 khive_db::stores::graph::EDGE_SYMMETRIC_DELETE_NONCANONICAL_GUARDED_SQL,
7752 rusqlite::params![
7753 &ns,
7754 &edge_id_str,
7755 expected_updated_at_micros,
7756 expected_deleted_at_micros,
7757 ],
7758 )
7759 .map_err(SqliteError::Rusqlite)?;
7760 if affected == 0 {
7761 return Ok(SymmetricEdgeUpdateOutcome::Stale);
7762 }
7763 Ok(SymmetricEdgeUpdateOutcome::Absorbed(existing_id))
7764 } else {
7765 let affected = conn
7771 .execute(
7772 khive_db::stores::graph::EDGE_SYMMETRIC_UPDATE_INPLACE_SQL,
7773 rusqlite::params![
7774 &canon_src_str,
7775 &canon_tgt_str,
7776 &relation_str,
7777 weight,
7778 now_ts,
7779 metadata,
7780 &ns,
7781 &edge_id_str,
7782 expected_updated_at_micros,
7783 expected_deleted_at_micros,
7784 ],
7785 )
7786 .map_err(SqliteError::Rusqlite)?;
7787 if affected == 0 {
7788 return Ok(SymmetricEdgeUpdateOutcome::Stale);
7789 }
7790 Ok(SymmetricEdgeUpdateOutcome::Updated)
7791 }
7792 }
7793
7794 pub async fn update_edge(
7808 &self,
7809 token: &NamespaceToken,
7810 edge_id: Uuid,
7811 patch: crate::curation::EdgePatch,
7812 ) -> RuntimeResult<Edge> {
7813 let graph_for_fetch = self.graph(token)?;
7816 let mut edge = graph_for_fetch
7817 .get_edge(LinkId::from(edge_id))
7818 .await?
7819 .ok_or_else(|| crate::RuntimeError::NotFound(format!("edge {edge_id}")))?;
7820 let expected_updated_at = edge.updated_at;
7821 let expected_deleted_at = edge.deleted_at;
7822 #[cfg(test)]
7823 crate::curation::race_seam::pause_after_read().await;
7824
7825 let record_ns: String = edge.namespace.clone();
7830 let record_tok = token.with_namespace(
7831 khive_types::Namespace::parse(&record_ns)
7832 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
7833 );
7834 let graph = self.graph(&record_tok)?;
7835
7836 let mut changed_fields: Vec<&'static str> = Vec::new();
7837 if let Some(r) = patch.relation {
7838 self.validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, r)
7841 .await?;
7842 edge.relation = r;
7843 changed_fields.push("relation");
7844 }
7845 if let Some(w) = patch.weight {
7846 if !khive_types::validate_edge_weight(w) {
7849 return Err(RuntimeError::InvalidInput(format!(
7850 "edge weight must be a finite value in [0.0, 1.0]; got {w}"
7851 )));
7852 }
7853 edge.weight = w;
7854 changed_fields.push("weight");
7855 }
7856 if let Some(props) = patch.properties {
7857 crate::secret_gate::reject_reserved_secret_gate_property(Some(&props))?;
7858 edge.metadata = Some(props);
7859 }
7860
7861 let (canon_src, canon_tgt) =
7871 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
7872
7873 if edge.relation.is_symmetric() {
7874 let ns = record_ns.clone();
7878 let edge_id_str = edge_id.to_string();
7879 let relation_str = edge.relation.to_string();
7880 let canon_src_str = canon_src.to_string();
7881 let canon_tgt_str = canon_tgt.to_string();
7882 let weight = edge.weight;
7883 let metadata = edge
7884 .metadata
7885 .as_ref()
7886 .map(|v| serde_json::to_string(v).unwrap_or_default());
7887
7888 let expected_updated_at_micros = expected_updated_at.timestamp_micros();
7889 let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
7890
7891 let pool = self.backend().pool_arc();
7892 let writer_task = pool
7893 .writer_task_for_runtime_write(RuntimeWriteOperation::UpdateSymmetricEdge)
7894 .map_err(RuntimeError::Storage)?;
7895
7896 let outcome: SymmetricEdgeUpdateOutcome = if let Some(writer_task) = writer_task {
7897 writer_task
7898 .send(move |conn| {
7899 Self::update_edge_symmetric_dml(
7900 conn,
7901 &ns,
7902 &edge_id_str,
7903 &canon_src_str,
7904 &canon_tgt_str,
7905 &relation_str,
7906 weight,
7907 metadata,
7908 expected_updated_at_micros,
7909 expected_deleted_at_micros,
7910 )
7911 .map_err(|e| {
7912 khive_storage::StorageError::driver(
7913 khive_storage::StorageCapability::Graph,
7914 "update_edge",
7915 e,
7916 )
7917 })
7918 })
7919 .await
7920 .map_err(RuntimeError::Storage)?
7921 } else {
7922 tokio::task::spawn_blocking(move || {
7923 let guard = pool.writer()?;
7924 guard.transaction(|conn| {
7925 Self::update_edge_symmetric_dml(
7926 conn,
7927 &ns,
7928 &edge_id_str,
7929 &canon_src_str,
7930 &canon_tgt_str,
7931 &relation_str,
7932 weight,
7933 metadata,
7934 expected_updated_at_micros,
7935 expected_deleted_at_micros,
7936 )
7937 })
7938 })
7939 .await
7940 .map_err(|e| {
7941 RuntimeError::Internal(format!("update_edge: spawn_blocking join: {e}"))
7942 })?
7943 .map_err(RuntimeError::Sqlite)?
7944 };
7945
7946 match outcome {
7947 SymmetricEdgeUpdateOutcome::Absorbed(sid) => {
7948 let surviving_uuid = Uuid::parse_str(&sid).map_err(|e| {
7955 RuntimeError::Internal(format!(
7956 "update_edge: surviving id parse failed: {e}"
7957 ))
7958 })?;
7959 edge = self
7960 .get_edge_including_deleted(&record_tok, surviving_uuid)
7961 .await?
7962 .ok_or_else(|| {
7963 RuntimeError::Internal(format!(
7964 "update_edge: surviving canonical row {surviving_uuid} vanished after update"
7965 ))
7966 })?;
7967 }
7968 SymmetricEdgeUpdateOutcome::Updated => {
7969 edge.source_id = canon_src;
7971 edge.target_id = canon_tgt;
7972 }
7973 SymmetricEdgeUpdateOutcome::Stale => {
7974 return Err(crate::curation::stale_edge_snapshot_error(edge_id));
7975 }
7976 }
7977 } else {
7978 let minimum_updated_at_micros = expected_updated_at
7991 .timestamp_micros()
7992 .checked_add(1)
7993 .ok_or_else(|| {
7994 RuntimeError::Internal(format!(
7995 "edge {edge_id} updated_at is already at i64::MAX and cannot advance"
7996 ))
7997 })?;
7998 let now_micros = chrono::Utc::now()
7999 .timestamp_micros()
8000 .max(minimum_updated_at_micros);
8001 edge.updated_at =
8002 chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
8003 RuntimeError::Internal(format!(
8004 "edge {edge_id}: computed updated_at {now_micros} is not a valid timestamp"
8005 ))
8006 })?;
8007 let persisted = graph
8008 .replace_edge_if_unchanged(edge.clone(), expected_updated_at, expected_deleted_at)
8009 .await?;
8010 if !persisted {
8011 return Err(crate::curation::stale_edge_snapshot_error(edge_id));
8012 }
8013 }
8014
8015 let event_store = self.events(&record_tok)?;
8017 let event = khive_storage::event::Event::new(
8018 record_ns.clone(),
8019 "update",
8020 EventKind::EdgeUpdated,
8021 SubstrateKind::Entity,
8022 "",
8023 )
8024 .with_target(edge_id)
8025 .with_payload(
8026 serde_json::json!({"id": edge_id, "namespace": record_ns, "changed_fields": changed_fields}),
8027 );
8028 event_store.append_event(event).await.map_err(|e| {
8029 RuntimeError::Internal(format!("update_edge: event store write failed: {e}"))
8030 })?;
8031
8032 Ok(edge)
8033 }
8034
8035 pub async fn delete_edge(
8046 &self,
8047 token: &NamespaceToken,
8048 edge_id: Uuid,
8049 hard: bool,
8050 ) -> RuntimeResult<bool> {
8051 let mode = if hard {
8052 DeleteMode::Hard
8053 } else {
8054 DeleteMode::Soft
8055 };
8056
8057 let edge = if hard {
8063 self.get_edge_including_deleted(token, edge_id).await?
8064 } else {
8065 self.get_edge(token, edge_id).await?
8066 };
8067 let Some(edge) = edge else {
8068 return Ok(false);
8069 };
8070
8071 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 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
8079
8080 let deleted = if hard {
8087 self.atomic_hard_delete_with_edge_purge(
8088 edge_hard_delete_statement(edge_id),
8089 edge_id,
8090 &record_ns,
8091 &actor,
8092 SubstrateKind::Entity,
8093 )
8094 .await?
8095 } else {
8096 graph.delete_edge(LinkId::from(edge_id), mode).await?
8097 };
8098 if deleted {
8099 let event_store = self.events(&record_tok)?;
8101 let event = khive_storage::event::Event::new(
8102 record_ns.clone(),
8103 "delete",
8104 EventKind::EdgeDeleted,
8105 SubstrateKind::Entity,
8106 "",
8107 )
8108 .with_target(edge_id)
8109 .with_payload(serde_json::json!({"id": edge_id, "namespace": record_ns, "hard": hard}));
8110 event_store.append_event(event).await.map_err(|e| {
8111 RuntimeError::Internal(format!("delete_edge: event store write failed: {e}"))
8112 })?;
8113 }
8114 Ok(deleted)
8115 }
8116
8117 pub async fn count_edges(
8119 &self,
8120 token: &NamespaceToken,
8121 filter: crate::curation::EdgeListFilter,
8122 ) -> RuntimeResult<u64> {
8123 let namespaces: Vec<String> = token
8124 .visible_namespaces()
8125 .iter()
8126 .map(|namespace| namespace.as_str().to_owned())
8127 .collect();
8128 let graph = self.graph(token)?;
8129 match graph
8130 .count_edges_in_namespaces(&namespaces, filter.clone().into())
8131 .await
8132 {
8133 Ok(count) => Ok(count),
8134 Err(khive_storage::StorageError::Unsupported { operation, .. })
8135 if operation == "count_edges_in_namespaces" =>
8136 {
8137 let mut total = 0;
8138 for namespace in token.visible_namespaces() {
8139 let scoped = NamespaceToken::for_namespace(namespace.clone());
8140 total += self
8141 .graph(&scoped)?
8142 .count_edges(filter.clone().into())
8143 .await?;
8144 }
8145 Ok(total)
8146 }
8147 Err(error) => Err(error.into()),
8148 }
8149 }
8150
8151 pub async fn build_edge(&self, token: &NamespaceToken, spec: &LinkSpec) -> RuntimeResult<Edge> {
8162 self.build_edge_with_endpoint_kinds(token, spec)
8163 .await
8164 .map(|(edge, _)| edge)
8165 }
8166
8167 async fn build_edge_with_endpoint_kinds(
8168 &self,
8169 token: &NamespaceToken,
8170 spec: &LinkSpec,
8171 ) -> RuntimeResult<(Edge, (EdgeEndpointKind, EdgeEndpointKind))> {
8172 validate_edge_metadata(spec.relation, spec.metadata.as_ref())?;
8173 let ns_str = match &spec.namespace {
8174 Some(s) => {
8175 let spec_ns = crate::Namespace::parse(s)
8176 .map_err(|e| RuntimeError::InvalidInput(format!("invalid namespace: {e}")))?;
8177 if &spec_ns != token.namespace() {
8178 return Err(RuntimeError::InvalidInput(
8179 "LinkSpec namespace does not match token namespace".into(),
8180 ));
8181 }
8182 s.as_str()
8183 }
8184 None => token.namespace().as_str(),
8185 };
8186 let endpoint_kinds = self
8187 .validate_edge_relation_endpoints(token, spec.source_id, spec.target_id, spec.relation)
8188 .await?;
8189 let (source_id, target_id) =
8190 canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
8191 let endpoint_kinds = canonical_edge_endpoint_kinds(
8192 spec.source_id,
8193 source_id,
8194 endpoint_kinds.0,
8195 endpoint_kinds.1,
8196 );
8197 let metadata = if spec.relation == EdgeRelation::DependsOn {
8198 match (
8203 self.resolve_edge_endpoint(token, source_id).await?,
8204 self.resolve_edge_endpoint(token, target_id).await?,
8205 ) {
8206 (Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
8207 merge_dependency_kind(&src_e.kind, &tgt_e.kind, spec.metadata.clone())
8208 }
8209 _ => spec.metadata.clone(),
8210 }
8211 } else {
8212 spec.metadata.clone()
8213 };
8214 validate_edge_metadata(spec.relation, metadata.as_ref())?;
8215 let now = chrono::Utc::now();
8216 Ok((
8217 Edge {
8218 id: LinkId::from(Uuid::new_v4()),
8219 namespace: ns_str.to_string(),
8220 source_id,
8221 target_id,
8222 relation: spec.relation,
8223 weight: spec.weight,
8224 created_at: now,
8225 updated_at: now,
8226 deleted_at: None,
8227 metadata,
8228 target_backend: None,
8229 },
8230 endpoint_kinds,
8231 ))
8232 }
8233
8234 pub async fn link_many(
8251 &self,
8252 token: &NamespaceToken,
8253 specs: Vec<LinkSpec>,
8254 ) -> RuntimeResult<Vec<Edge>> {
8255 self.link_many_observed(token, specs)
8256 .await
8257 .map(|rows| rows.into_iter().map(|row| row.edge).collect())
8258 }
8259
8260 pub async fn link_many_observed(
8264 &self,
8265 token: &NamespaceToken,
8266 specs: Vec<LinkSpec>,
8267 ) -> RuntimeResult<Vec<EdgeUpsertResult>> {
8268 self.link_many_guarded_observed(
8269 token,
8270 specs,
8271 GraphMutationPreconditions::default(),
8272 Vec::new(),
8273 )
8274 .await
8275 .map(|(rows, _)| rows)
8276 }
8277
8278 #[doc(hidden)]
8282 pub async fn link_many_guarded_observed(
8283 &self,
8284 token: &NamespaceToken,
8285 specs: Vec<LinkSpec>,
8286 preconditions: GraphMutationPreconditions,
8287 retirements: Vec<Edge>,
8288 ) -> RuntimeResult<(Vec<EdgeUpsertResult>, Vec<LinkId>)> {
8289 let namespace = token.namespace().as_str();
8290 let foreign_document = preconditions
8291 .document
8292 .as_ref()
8293 .is_some_and(|guard| guard.namespace != namespace);
8294 let foreign_edge = preconditions.edges.iter().any(|guard| {
8295 guard.namespace != namespace
8296 || guard
8297 .expected
8298 .as_ref()
8299 .is_some_and(|edge| edge.namespace != namespace)
8300 });
8301 if foreign_document
8302 || foreign_edge
8303 || retirements.iter().any(|edge| edge.namespace != namespace)
8304 {
8305 return Err(RuntimeError::InvalidInput(
8306 "guarded link namespace does not match token namespace".into(),
8307 ));
8308 }
8309 if specs.is_empty()
8310 && preconditions.document.is_none()
8311 && preconditions.edges.is_empty()
8312 && retirements.is_empty()
8313 {
8314 return Ok((Vec::new(), Vec::new()));
8315 }
8316 let mut edges = Vec::with_capacity(specs.len());
8317 let mut endpoint_kinds = Vec::with_capacity(specs.len());
8318 for spec in &specs {
8319 let (edge, kinds) = self.build_edge_with_endpoint_kinds(token, spec).await?;
8320 edges.push(edge);
8321 endpoint_kinds.push(kinds);
8322 }
8323 let requests = edges
8332 .into_iter()
8333 .zip(specs.iter())
8334 .map(|(edge, spec)| EdgeUpsertRequest {
8335 edge,
8336 resurrect: spec.resurrect,
8337 })
8338 .collect();
8339 let attribution = crate::EventAttribution::from_token(token);
8340 let outcome = compose_graph_mutation_events(
8341 self.backend(),
8342 GraphMutationRequest::Batch {
8343 requests,
8344 guard_endpoints: true,
8345 },
8346 preconditions,
8347 retirements,
8348 move |outcome| {
8349 let GraphMutationOutcome::Batch(batch) = &outcome.mutation else {
8350 return Err(Self::link_composition_shape_error(
8351 "expected a written batch",
8352 ));
8353 };
8354 if batch.rows.len() != endpoint_kinds.len() {
8355 return Err(Self::link_composition_shape_error(
8356 "edge result count differs from validated endpoint count",
8357 ));
8358 }
8359 let mut events = Vec::with_capacity(batch.rows.len() + outcome.retired.len());
8360 for (row, (source_kind, target_kind)) in batch.rows.iter().zip(endpoint_kinds) {
8361 events.push(Self::link_mutation_event(
8362 &attribution,
8363 row,
8364 source_kind,
8365 target_kind,
8366 ));
8367 }
8368 for edge in &outcome.retired {
8369 let edge_id = Uuid::from(edge.id);
8370 events.push(
8371 attribution.stamp(
8372 Event::new(
8373 edge.namespace.clone(),
8374 "delete",
8375 EventKind::EdgeDeleted,
8376 SubstrateKind::Entity,
8377 "",
8378 )
8379 .with_target(edge_id)
8380 .with_payload(serde_json::json!({
8381 "id": edge_id, "namespace": edge.namespace, "hard": false,
8382 })),
8383 ),
8384 );
8385 }
8386 Ok(events)
8387 },
8388 )
8389 .await?;
8390 let retired = outcome.retired.into_iter().map(|edge| edge.id).collect();
8391 let GraphMutationOutcome::Batch(outcome) = outcome.mutation else {
8392 return Err(RuntimeError::Internal(
8393 "link_many: unexpected composition outcome".into(),
8394 ));
8395 };
8396 if let Some(refusal) = outcome.refusal {
8397 return match refusal.reason {
8398 EdgeUpsertRefusal::MissingEndpoints(missing) => {
8399 Err(RuntimeError::GuardedWriteFailed(guarded_link_batch_failure(
8400 &specs[refusal.entry_index],
8401 refusal.entry_index,
8402 missing,
8403 )))
8404 }
8405 EdgeUpsertRefusal::ResurrectionRequired { edge } => {
8406 Err(RuntimeError::InvalidInput(format!(
8407 "batch entry {} targets soft-deleted edge {}; pass resurrect=true for that link",
8408 refusal.entry_index, edge.id
8409 )))
8410 }
8411 };
8412 }
8413 Ok((outcome.rows, retired))
8414 }
8415
8416 pub async fn link_commit_annotation_if_absent(
8421 &self,
8422 token: &NamespaceToken,
8423 commit_id: Uuid,
8424 project_id: Uuid,
8425 guard: CommitAnnotationGuard,
8426 ) -> RuntimeResult<CommitAnnotationInsertOutcome> {
8427 if !matches!(guard.expected_sha.len(), 40 | 64)
8428 || !guard
8429 .expected_sha
8430 .bytes()
8431 .all(|byte| byte.is_ascii_hexdigit())
8432 {
8433 return Err(RuntimeError::InvalidInput(
8434 "expected full commit SHA".into(),
8435 ));
8436 }
8437 let edge = self
8438 .build_edge(
8439 token,
8440 &LinkSpec {
8441 namespace: None,
8442 source_id: commit_id,
8443 target_id: project_id,
8444 relation: EdgeRelation::Annotates,
8445 weight: 1.0,
8446 metadata: None,
8447 resurrect: false,
8448 },
8449 )
8450 .await?;
8451 let attribution = crate::EventAttribution::from_token(token);
8452 let outcome = compose_graph_mutation_events(
8453 self.backend(),
8454 GraphMutationRequest::CommitAnnotation { edge, guard },
8455 GraphMutationPreconditions::default(),
8456 Vec::new(),
8457 move |outcome| match &outcome.mutation {
8458 GraphMutationOutcome::CommitAnnotation(CommitAnnotationInsertOutcome::Created(
8459 edge,
8460 )) => Ok(vec![Self::link_mutation_event(
8461 &attribution,
8462 &EdgeUpsertResult {
8463 edge: edge.clone(),
8464 disposition: EdgeUpsertDisposition::Created,
8465 previous: None,
8466 },
8467 EdgeEndpointKind::Note,
8468 EdgeEndpointKind::Entity,
8469 )]),
8470 _ => Err(Self::link_composition_shape_error(
8471 "expected a created annotation",
8472 )),
8473 },
8474 )
8475 .await?;
8476 let GraphMutationOutcome::CommitAnnotation(result) = outcome.mutation else {
8477 return Err(RuntimeError::Internal(
8478 "link annotation: unexpected composition outcome".into(),
8479 ));
8480 };
8481 Ok(result)
8482 }
8483
8484 pub async fn create_many(
8495 &self,
8496 token: &NamespaceToken,
8497 specs: Vec<EntityCreateSpec>,
8498 ) -> RuntimeResult<Vec<Entity>> {
8499 if specs.is_empty() {
8500 return Ok(vec![]);
8501 }
8502 let ns = token.namespace().as_str();
8503
8504 let mut entities = Vec::with_capacity(specs.len());
8508 for (index, spec) in specs.iter().enumerate() {
8509 entities.push(self.validate_bulk_entity(ns, spec, &format!("entity[{index}]"))?);
8510 }
8511
8512 #[cfg(any(test, feature = "fault-injection"))]
8513 let fts_many_inject = consume_fault(&FTS_FAIL_MANY_NS, ns);
8514 #[cfg(not(any(test, feature = "fault-injection")))]
8515 let fts_many_inject = false;
8516
8517 #[cfg(any(test, feature = "fault-injection"))]
8518 let fts_many_inject_partial = consume_fault(&FTS_FAIL_MANY_PARTIAL_NS, ns);
8519 #[cfg(not(any(test, feature = "fault-injection")))]
8520 let fts_many_inject_partial = false;
8521
8522 let injected_failure_index = if fts_many_inject {
8523 Some(0)
8524 } else if fts_many_inject_partial {
8525 Some(usize::from(entities.len() > 1))
8526 } else {
8527 None
8528 };
8529
8530 let _ = self.entities(token)?;
8531 let _ = self.text(token)?;
8532
8533 let plans = entities
8534 .iter()
8535 .enumerate()
8536 .map(|(index, entity)| {
8537 let mut plan = bulk_entity_plan(entity)?;
8538 if injected_failure_index == Some(index) {
8539 plan.statements.truncate(1);
8541 plan.statements.push(PlanStatement {
8542 statement: SqlStatement {
8543 sql:
8544 "INSERT INTO __khive_create_many_injected_failure__ DEFAULT VALUES"
8545 .to_string(),
8546 params: vec![],
8547 label: Some("fts-insert-injected-failure".to_string()),
8548 },
8549 guard: None,
8550 });
8551 }
8552 Ok(AtomicOpPlan::AddEntity(plan))
8553 })
8554 .collect::<RuntimeResult<Vec<_>>>()?;
8555
8556 match run_atomic_unit(self.sql().as_ref(), plans).await {
8557 Ok(AtomicRunOutcome::Committed { .. }) => Ok(entities),
8558 Ok(AtomicRunOutcome::RolledBack {
8559 failed_op_index,
8560 failure,
8561 }) => Err(RuntimeError::Internal(format!(
8562 "create_many: atomic batch rolled back at entity index {failed_op_index}: \
8563 {failure:?}"
8564 ))),
8565 Err(e) => Err(RuntimeError::Internal(format!(
8566 "create_many: atomic batch failed: {}",
8567 e.0
8568 ))),
8569 }
8570 }
8571
8572 fn validate_bulk_entity(
8577 &self,
8578 ns: &str,
8579 spec: &EntityCreateSpec,
8580 record: &str,
8581 ) -> RuntimeResult<Entity> {
8582 self.validate_entity_kind(&spec.kind)?;
8583 let validated_type =
8588 self.validate_entity_type_for_kind(&spec.kind, spec.entity_type.as_deref())?;
8589 if spec.name.trim().is_empty() {
8590 return Err(RuntimeError::InvalidInput("name must not be empty".into()));
8591 }
8592 crate::secret_gate::reject_reserved_secret_gate_property(spec.properties.as_ref())?;
8593 crate::secret_gate::check_at(&spec.name, record, "name")?;
8594 if let Some(d) = &spec.description {
8595 crate::secret_gate::check_at(d, record, "description")?;
8596 }
8597 if let Some(ref p) = spec.properties {
8598 crate::secret_gate::check_json_at(p, record, "properties")?;
8599 }
8600 crate::secret_gate::check_tags_at(&spec.tags, record, "tags")?;
8601
8602 let mut entity =
8603 Entity::new(ns, &spec.kind, &spec.name).with_entity_type(validated_type.as_deref());
8604 if let Some(d) = &spec.description {
8605 entity = entity.with_description(d);
8606 }
8607 if let Some(p) = spec.properties.clone() {
8608 entity = entity.with_properties(p);
8609 }
8610 if !spec.tags.is_empty() {
8611 entity = entity.with_tags(spec.tags.clone());
8612 }
8613 Ok(entity)
8614 }
8615
8616 pub async fn prepare_bulk_entity_plan(
8623 &self,
8624 token: &NamespaceToken,
8625 spec: EntityCreateSpec,
8626 ) -> RuntimeResult<(Entity, AtomicOpPlan)> {
8627 let entity = self.validate_bulk_entity(token.namespace().as_str(), &spec, "entity")?;
8628 let _ = self.entities(token)?;
8629 let _ = self.text(token)?;
8630
8631 let plan = AtomicOpPlan::AddEntity(bulk_entity_plan(&entity)?);
8632 Ok((entity, plan))
8633 }
8634
8635 pub async fn prepare_bulk_note_plan(
8644 &self,
8645 token: &NamespaceToken,
8646 spec: NoteCreateSpec,
8647 ) -> RuntimeResult<(Note, AtomicOpPlan)> {
8648 let mut candidate = Note::new(token.namespace().as_str(), &spec.kind, &spec.content);
8649 candidate.name = spec.name.clone();
8650 candidate.properties = spec.properties.clone();
8651 crate::note_write::validate_head(&candidate)?;
8652 let mut prepared = crate::atomic_message::prepare_atomic_notes(
8653 self,
8654 vec![crate::atomic_message::AtomicNoteSpec {
8655 token,
8656 id: None,
8657 kind: &spec.kind,
8658 name: spec.name.as_deref(),
8659 content: &spec.content,
8660 properties: spec.properties,
8661 }],
8662 crate::atomic_message::AtomicNoteOptions {
8663 salience: spec.salience,
8664 embed: Some(false),
8665 ..Default::default()
8666 },
8667 )
8668 .await?;
8669 match (prepared.notes.pop(), prepared.plans.pop()) {
8670 (Some(note), Some(plan)) if prepared.notes.is_empty() && prepared.plans.is_empty() => {
8671 Ok((note, plan))
8672 }
8673 _ => Err(RuntimeError::Internal(
8674 "bulk note preparation must yield exactly one note and one plan".into(),
8675 )),
8676 }
8677 }
8678}
8679
8680#[derive(Clone, Debug)]
8685pub struct NoteCreateSpec {
8686 pub kind: String,
8687 pub name: Option<String>,
8688 pub content: String,
8689 pub salience: Option<f64>,
8690 pub properties: Option<serde_json::Value>,
8691}
8692
8693fn bulk_entity_plan(entity: &Entity) -> RuntimeResult<AddEntityPlan> {
8694 crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
8695 let mut statements = vec![PlanStatement {
8696 statement: entity_upsert_statement(entity),
8697 guard: Some(AffectedRowGuard::exactly(1)),
8698 }];
8699 statements.extend(
8701 insert_document_statements("fts_entities", &entity_fts_document(entity))
8702 .into_iter()
8703 .map(|statement| PlanStatement {
8704 statement,
8705 guard: None,
8706 }),
8707 );
8708 Ok(AddEntityPlan {
8709 entity_id: entity.id,
8710 statements,
8711 post_commit: PostCommitEffect::None,
8712 })
8713}
8714
8715fn guarded_link_batch_failure(
8716 spec: &LinkSpec,
8717 entry_index: usize,
8718 missing: khive_storage::MissingEndpoints,
8719) -> GuardedWriteFailure {
8720 let (source_id, target_id) =
8723 canonical_edge_endpoints(spec.relation, spec.source_id, spec.target_id);
8724 GuardedWriteFailure {
8725 entry_index: Some(entry_index),
8726 missing_source: missing.source.then_some(source_id),
8727 missing_target: missing.target.then_some(target_id),
8728 }
8729}
8730
8731#[derive(Clone, Debug)]
8734pub struct LinkSpec {
8735 pub namespace: Option<String>,
8736 pub source_id: Uuid,
8737 pub target_id: Uuid,
8738 pub relation: EdgeRelation,
8739 pub weight: f64,
8740 pub metadata: Option<serde_json::Value>,
8741 pub resurrect: bool,
8742}
8743
8744#[derive(Clone, Debug)]
8752pub struct EntityCreateSpec {
8753 pub kind: String,
8754 pub entity_type: Option<String>,
8755 pub name: String,
8756 pub description: Option<String>,
8757 pub properties: Option<serde_json::Value>,
8758 pub tags: Vec<String>,
8759}
8760
8761#[cfg(test)]
8767#[path = "operations_tests.rs"]
8768mod tests;