1#[path = "graph/write_transaction.rs"]
4mod write_transaction;
5use write_transaction::run_graph_mutation_transaction;
6
7use std::collections::{HashMap, HashSet, VecDeque};
8use std::sync::Arc;
9
10use async_trait::async_trait;
11use chrono::{DateTime, TimeZone, Utc};
12use rusqlite::OptionalExtension;
13use uuid::Uuid;
14
15use khive_storage::error::StorageError;
16use khive_storage::graph::{CommitAnnotationGuard, CommitAnnotationInsertOutcome};
17use khive_storage::types::{
18 BatchWriteSummary, DeleteMode, DirectedNeighborHit, Direction, Edge, EdgeEndpointBaseCounts,
19 EdgeFilter, EdgeSeekPage, EdgeSortField, EdgeUpsertDisposition, EdgeUpsertRefusal,
20 EdgeUpsertRequest, EdgeUpsertResult, GraphPath, GuardedBatchOutcome, GuardedBatchRefusal,
21 GuardedEdgeBatchRefusal, GuardedEdgeBatchUpsertOutcome, GuardedEdgeUpsertOutcome,
22 GuardedWriteOutcome, MissingEndpoints, NeighborCursor, NeighborHit, NeighborQuery, Page,
23 PageRequest, PathNode, SeekCursor, SeekPage, SortDirection, SortOrder, SqlStatement, SqlValue,
24 TraversalExecutionBudget, TraversalOptions, TraversalRequest,
25};
26use khive_storage::LinkId;
27use khive_storage::StorageCapability;
28use khive_storage::{Event, GraphStore, StorageResult};
29use khive_types::EdgeRelation;
30
31use crate::error::SqliteError;
32use crate::pool::ConnectionPool;
33use crate::sql_bridge::bind_params;
34use crate::writer_task::WriterTaskHandle;
35
36fn map_err(e: rusqlite::Error, op: &'static str) -> StorageError {
38 StorageError::driver(StorageCapability::Graph, op, e)
39}
40
41fn map_sqlite_err(e: SqliteError, op: &'static str) -> StorageError {
42 e.into_storage_error(StorageCapability::Graph, op)
43}
44
45fn resurrection_required_error(operation: &'static str, edge: &Edge) -> StorageError {
46 StorageError::Conflict {
47 capability: StorageCapability::Graph,
48 operation: operation.into(),
49 message: format!(
50 "edge {} is soft-deleted; explicit resurrection is required",
51 edge.id
52 ),
53 }
54}
55
56const NAMESPACE_COUNT_CHUNK_SIZE: usize = 500;
57
58const LATEST_ANNOTATING_NOTE_SQL: &str = r#"WITH incident AS MATERIALIZED (
67 SELECT source_id, deleted_at
68 FROM graph_edges INDEXED BY idx_graph_edges_ns_tgt_rel
69 WHERE namespace = ?1 AND target_id = ?2 AND relation = 'annotates'
70 LIMIT 3
71)
72SELECT result.id, result.created_at
73FROM notes AS result
74WHERE result.id = CASE WHEN (SELECT count(*) FROM incident) <= 2 THEN (
75 SELECT n.id
76 FROM incident AS e CROSS JOIN notes AS n
77 WHERE n.id = e.source_id AND e.deleted_at IS NULL
78 AND n.deleted_at IS NULL AND n.kind = ?3
79 AND EXISTS (SELECT 1 FROM json_each(CASE
80 WHEN json_type(n.properties, '$.tags') = 'array'
81 THEN json_extract(n.properties, '$.tags') ELSE '[]' END) AS tag
82 WHERE tag.type = 'text' AND tag.value = ?4 COLLATE BINARY)
83 ORDER BY n.created_at DESC, n.id ASC LIMIT 1
84) ELSE (
85 SELECT n.id
86 FROM notes AS n INDEXED BY idx_notes_created
87 WHERE n.deleted_at IS NULL AND n.kind = ?3
88 AND EXISTS (SELECT 1 FROM json_each(CASE
89 WHEN json_type(n.properties, '$.tags') = 'array'
90 THEN json_extract(n.properties, '$.tags') ELSE '[]' END) AS tag
91 WHERE tag.type = 'text' AND tag.value = ?4 COLLATE BINARY)
92 AND EXISTS (SELECT 1 FROM graph_edges AS e INDEXED BY idx_graph_edges_unique_triple
93 WHERE e.namespace = ?1 AND e.source_id = n.id AND e.target_id = ?2
94 AND e.relation = 'annotates' AND e.deleted_at IS NULL)
95 ORDER BY n.created_at DESC, n.id ASC LIMIT 1
96) END"#;
97
98const LATEST_ANNOTATING_NOTE_WITH_PROPERTY_SQL: &str = r#"WITH incident AS MATERIALIZED (
102 SELECT source_id, deleted_at
103 FROM graph_edges INDEXED BY idx_graph_edges_ns_tgt_rel
104 WHERE namespace = ?1 AND target_id = ?2 AND relation = 'annotates'
105 LIMIT 3
106)
107SELECT result.id, result.created_at
108FROM notes AS result
109WHERE result.id = CASE WHEN (SELECT count(*) FROM incident) <= 2 THEN (
110 SELECT n.id
111 FROM incident AS e CROSS JOIN notes AS n
112 WHERE n.id = e.source_id AND e.deleted_at IS NULL
113 AND n.deleted_at IS NULL AND n.kind = ?3
114 AND EXISTS (SELECT 1 FROM json_each(CASE
115 WHEN json_type(n.properties, '$.tags') = 'array'
116 THEN json_extract(n.properties, '$.tags') ELSE '[]' END) AS tag
117 WHERE tag.type = 'text' AND tag.value = ?4 COLLATE BINARY)
118 AND EXISTS (SELECT 1 FROM json_each(n.properties) AS property
119 WHERE property.key = ?5 COLLATE BINARY AND property.type = 'text'
120 AND property.value = ?6 COLLATE BINARY)
121 ORDER BY n.created_at DESC, n.id ASC LIMIT 1
122) ELSE (
123 SELECT n.id
124 FROM notes AS n INDEXED BY idx_notes_created
125 WHERE n.deleted_at IS NULL AND n.kind = ?3
126 AND EXISTS (SELECT 1 FROM json_each(CASE
127 WHEN json_type(n.properties, '$.tags') = 'array'
128 THEN json_extract(n.properties, '$.tags') ELSE '[]' END) AS tag
129 WHERE tag.type = 'text' AND tag.value = ?4 COLLATE BINARY)
130 AND EXISTS (SELECT 1 FROM json_each(n.properties) AS property
131 WHERE property.key = ?5 COLLATE BINARY AND property.type = 'text'
132 AND property.value = ?6 COLLATE BINARY)
133 AND EXISTS (SELECT 1 FROM graph_edges AS e INDEXED BY idx_graph_edges_unique_triple
134 WHERE e.namespace = ?1 AND e.source_id = n.id AND e.target_id = ?2
135 AND e.relation = 'annotates' AND e.deleted_at IS NULL)
136 ORDER BY n.created_at DESC, n.id ASC LIMIT 1
137) END"#;
138
139fn edge_conflict_clause(resurrect: bool) -> String {
151 let deleted_at = if resurrect {
152 "NULL"
153 } else {
154 "graph_edges.deleted_at"
155 };
156 let predicate = if resurrect {
157 ""
158 } else {
159 " WHERE graph_edges.deleted_at IS NULL"
160 };
161 format!(
162 "weight = excluded.weight, \
163 updated_at = excluded.updated_at, \
164 deleted_at = {deleted_at}, \
165 metadata = excluded.metadata, \
166 target_backend = excluded.target_backend{predicate}"
167 )
168}
169
170fn endpoint_exists_clause(id_param: &str) -> String {
181 format!(
182 "EXISTS (SELECT 1 FROM entities WHERE id = {id_param} AND deleted_at IS NULL) \
183 OR EXISTS (SELECT 1 FROM notes WHERE id = {id_param} AND deleted_at IS NULL) \
184 OR EXISTS (SELECT 1 FROM events WHERE id = {id_param}) \
185 OR EXISTS (SELECT 1 FROM graph_edges WHERE id = {id_param} AND deleted_at IS NULL)"
186 )
187}
188
189pub fn edge_snapshot_assertion_statement(edge: &Edge, require_endpoints: bool) -> SqlStatement {
193 let mut sql = "SELECT 1 FROM graph_edges WHERE id=?1 AND namespace=?2 \
194 AND source_id=?3 AND target_id=?4 AND relation=?5 \
195 AND updated_at=?6 AND deleted_at IS NULL"
196 .to_string();
197 if require_endpoints {
198 sql.push_str(&format!(
199 " AND ({}) AND ({})",
200 endpoint_exists_clause("?3"),
201 endpoint_exists_clause("?4")
202 ));
203 }
204 SqlStatement {
205 sql,
206 params: vec![
207 SqlValue::Text(Uuid::from(edge.id).to_string()),
208 SqlValue::Text(edge.namespace.clone()),
209 SqlValue::Text(edge.source_id.to_string()),
210 SqlValue::Text(edge.target_id.to_string()),
211 SqlValue::Text(edge.relation.to_string()),
212 SqlValue::Integer(edge.updated_at.timestamp_micros()),
213 ],
214 label: Some("edge-snapshot-assertion".into()),
215 }
216}
217
218pub fn edge_upsert_statement(edge: &Edge) -> SqlStatement {
222 edge_upsert_statement_with_resurrection(edge, false)
223}
224
225pub fn edge_upsert_statement_with_resurrection(edge: &Edge, resurrect: bool) -> SqlStatement {
227 let (source_id, target_id) =
228 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
229 let metadata_str = edge
230 .metadata
231 .as_ref()
232 .map(|v| serde_json::to_string(v).unwrap_or_default());
233 let conflict_clause = edge_conflict_clause(resurrect);
234 SqlStatement {
235 sql: format!(
236 "INSERT INTO graph_edges \
237 (namespace, id, source_id, target_id, relation, weight, \
238 created_at, updated_at, deleted_at, metadata, target_backend) \
239 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) \
240 ON CONFLICT(namespace, id) DO UPDATE SET \
241 source_id = excluded.source_id, \
242 target_id = excluded.target_id, \
243 relation = excluded.relation, \
244 {conflict_clause} \
245 ON CONFLICT(namespace, source_id, target_id, relation) DO UPDATE SET \
246 {conflict_clause}"
247 ),
248 params: vec![
249 SqlValue::Text(edge.namespace.clone()),
250 SqlValue::Text(Uuid::from(edge.id).to_string()),
251 SqlValue::Text(source_id.to_string()),
252 SqlValue::Text(target_id.to_string()),
253 SqlValue::Text(edge.relation.to_string()),
254 SqlValue::Float(edge.weight),
255 SqlValue::Integer(edge.created_at.timestamp_micros()),
256 SqlValue::Integer(edge.updated_at.timestamp_micros()),
257 match edge.deleted_at {
258 Some(t) => SqlValue::Integer(t.timestamp_micros()),
259 None => SqlValue::Null,
260 },
261 match metadata_str {
262 Some(m) => SqlValue::Text(m),
263 None => SqlValue::Null,
264 },
265 match &edge.target_backend {
266 Some(b) => SqlValue::Text(b.clone()),
267 None => SqlValue::Null,
268 },
269 ],
270 label: Some("edge-upsert".to_string()),
271 }
272}
273
274pub fn edge_insert_only_guarded_by_endpoints_statement(edge: &Edge) -> SqlStatement {
279 let mut statement = edge_upsert_statement(edge);
280 let src_exists = endpoint_exists_clause("?3");
281 let tgt_exists = endpoint_exists_clause("?4");
282 statement.sql = format!(
283 "INSERT INTO graph_edges \
284 (namespace, id, source_id, target_id, relation, weight, \
285 created_at, updated_at, deleted_at, metadata, target_backend) \
286 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11 \
287 WHERE ({src_exists}) AND ({tgt_exists})"
288 );
289 statement.label = Some("edge-insert-only-where-endpoints-exist".to_string());
290 statement
291}
292
293pub fn edge_insert_if_absent_statement(edge: &Edge) -> SqlStatement {
297 let mut statement = edge_upsert_statement(edge);
298 statement.sql = "INSERT INTO graph_edges \
299 (namespace, id, source_id, target_id, relation, weight, \
300 created_at, updated_at, deleted_at, metadata, target_backend) \
301 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11) \
302 ON CONFLICT DO NOTHING"
303 .to_string();
304 statement.label = Some("edge-insert-if-absent".to_string());
305 statement
306}
307
308pub fn edge_replace_if_unchanged_statement(
319 edge: &Edge,
320 expected_updated_at: DateTime<Utc>,
321 expected_deleted_at: Option<DateTime<Utc>>,
322) -> SqlStatement {
323 let (source_id, target_id) =
324 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
325 let metadata_str = edge
326 .metadata
327 .as_ref()
328 .map(|v| serde_json::to_string(v).unwrap_or_default());
329 SqlStatement {
330 sql: "UPDATE graph_edges SET \
331 namespace = ?1, source_id = ?2, target_id = ?3, relation = ?4, weight = ?5, \
332 updated_at = ?6, deleted_at = ?7, metadata = ?8, target_backend = ?9 \
333 WHERE id = ?10 AND updated_at = ?11 AND deleted_at IS ?12 \
334 AND ?6 > updated_at"
335 .to_string(),
336 params: vec![
337 SqlValue::Text(edge.namespace.clone()),
338 SqlValue::Text(source_id.to_string()),
339 SqlValue::Text(target_id.to_string()),
340 SqlValue::Text(edge.relation.to_string()),
341 SqlValue::Float(edge.weight),
342 SqlValue::Integer(edge.updated_at.timestamp_micros()),
343 match edge.deleted_at {
344 Some(t) => SqlValue::Integer(t.timestamp_micros()),
345 None => SqlValue::Null,
346 },
347 match metadata_str {
348 Some(m) => SqlValue::Text(m),
349 None => SqlValue::Null,
350 },
351 match &edge.target_backend {
352 Some(b) => SqlValue::Text(b.clone()),
353 None => SqlValue::Null,
354 },
355 SqlValue::Text(Uuid::from(edge.id).to_string()),
356 SqlValue::Integer(expected_updated_at.timestamp_micros()),
357 match expected_deleted_at {
358 Some(value) => SqlValue::Integer(value.timestamp_micros()),
359 None => SqlValue::Null,
360 },
361 ],
362 label: Some("edge-replace-if-unchanged".to_string()),
363 }
364}
365
366#[allow(clippy::too_many_arguments)]
386pub fn edge_insert_guarded_by_endpoints_statement(
387 namespace: &str,
388 edge_id: Uuid,
389 source_id: Uuid,
390 target_id: Uuid,
391 relation: EdgeRelation,
392 weight: f64,
393 now: i64,
394 metadata: Option<&str>,
395) -> SqlStatement {
396 edge_insert_guarded_by_endpoints_with_resurrection_statement(
397 namespace, source_id, target_id, edge_id, relation, weight, now, metadata, false,
398 )
399}
400
401#[allow(clippy::too_many_arguments)]
403pub fn edge_insert_guarded_by_endpoints_with_resurrection_statement(
404 namespace: &str,
405 source_id: Uuid,
406 target_id: Uuid,
407 edge_id: Uuid,
408 relation: EdgeRelation,
409 weight: f64,
410 now: i64,
411 metadata: Option<&str>,
412 resurrect: bool,
413) -> SqlStatement {
414 let src_exists = endpoint_exists_clause("?3");
415 let tgt_exists = endpoint_exists_clause("?4");
416 let conflict_clause = edge_conflict_clause(resurrect);
417 SqlStatement {
418 sql: format!(
419 "INSERT INTO graph_edges \
420 (namespace, id, source_id, target_id, relation, weight, \
421 created_at, updated_at, metadata) \
422 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?7, ?8 \
423 WHERE ({src_exists}) AND ({tgt_exists}) \
424 ON CONFLICT(namespace, source_id, target_id, relation) DO UPDATE SET \
425 {conflict_clause}"
426 ),
427 params: vec![
428 SqlValue::Text(namespace.to_string()),
429 SqlValue::Text(edge_id.to_string()),
430 SqlValue::Text(source_id.to_string()),
431 SqlValue::Text(target_id.to_string()),
432 SqlValue::Text(relation.as_str().to_string()),
433 SqlValue::Float(weight),
434 SqlValue::Integer(now),
435 match metadata {
436 Some(m) => SqlValue::Text(m.to_string()),
437 None => SqlValue::Null,
438 },
439 ],
440 label: Some("atomic-link-insert-edge-where-exists".to_string()),
441 }
442}
443
444#[allow(clippy::too_many_arguments)]
448pub fn edge_insert_new_guarded_by_endpoints_statement(
449 namespace: &str,
450 edge_id: Uuid,
451 source_id: Uuid,
452 target_id: Uuid,
453 relation: EdgeRelation,
454 weight: f64,
455 now: i64,
456 metadata: Option<&str>,
457) -> SqlStatement {
458 let src_exists = endpoint_exists_clause("?3");
459 let tgt_exists = endpoint_exists_clause("?4");
460 SqlStatement {
461 sql: format!(
462 "INSERT INTO graph_edges \
463 (namespace, id, source_id, target_id, relation, weight, \
464 created_at, updated_at, metadata) \
465 SELECT ?1, ?2, ?3, ?4, ?5, ?6, ?7, ?7, ?8 \
466 WHERE ({src_exists}) AND ({tgt_exists}) \
467 ON CONFLICT DO NOTHING"
468 ),
469 params: vec![
470 SqlValue::Text(namespace.to_string()),
471 SqlValue::Text(edge_id.to_string()),
472 SqlValue::Text(source_id.to_string()),
473 SqlValue::Text(target_id.to_string()),
474 SqlValue::Text(relation.to_string()),
475 SqlValue::Float(weight),
476 SqlValue::Integer(now),
477 metadata
478 .map(|value| SqlValue::Text(value.to_string()))
479 .unwrap_or(SqlValue::Null),
480 ],
481 label: Some("edge-link-create-if-absent-and-endpoints-exist".to_string()),
482 }
483}
484
485pub fn edge_link_replace_if_unchanged_and_endpoints_exist_statement(
491 previous: &Edge,
492 weight: f64,
493 now: i64,
494 metadata: Option<&str>,
495) -> SqlStatement {
496 let src_exists = endpoint_exists_clause("?6");
497 let tgt_exists = endpoint_exists_clause("?7");
498 SqlStatement {
499 sql: format!(
500 "UPDATE graph_edges SET \
501 weight = ?1, updated_at = ?2, deleted_at = NULL, \
502 metadata = ?3, target_backend = NULL \
503 WHERE namespace = ?4 AND id = ?5 \
504 AND source_id = ?6 AND target_id = ?7 AND relation = ?8 \
505 AND updated_at = ?9 AND deleted_at IS ?10 AND ?2 > updated_at \
506 AND ({src_exists}) AND ({tgt_exists})"
507 ),
508 params: vec![
509 SqlValue::Float(weight),
510 SqlValue::Integer(now),
511 metadata
512 .map(|value| SqlValue::Text(value.to_string()))
513 .unwrap_or(SqlValue::Null),
514 SqlValue::Text(previous.namespace.clone()),
515 SqlValue::Text(Uuid::from(previous.id).to_string()),
516 SqlValue::Text(previous.source_id.to_string()),
517 SqlValue::Text(previous.target_id.to_string()),
518 SqlValue::Text(previous.relation.to_string()),
519 SqlValue::Integer(previous.updated_at.timestamp_micros()),
520 previous
521 .deleted_at
522 .map(|value| SqlValue::Integer(value.timestamp_micros()))
523 .unwrap_or(SqlValue::Null),
524 ],
525 label: Some("edge-link-replace-if-unchanged-and-endpoints-exist".to_string()),
526 }
527}
528
529pub fn edge_soft_delete_statement(id: Uuid, now: i64) -> SqlStatement {
531 SqlStatement {
532 sql: "UPDATE graph_edges SET deleted_at = ?2, updated_at = ?2 \
533 WHERE id = ?1 AND deleted_at IS NULL"
534 .to_string(),
535 params: vec![SqlValue::Text(id.to_string()), SqlValue::Integer(now)],
536 label: Some("edge-delete-soft".to_string()),
537 }
538}
539
540pub fn edge_hard_delete_statement(id: Uuid) -> SqlStatement {
542 SqlStatement {
543 sql: "DELETE FROM graph_edges WHERE id = ?1".to_string(),
544 params: vec![SqlValue::Text(id.to_string())],
545 label: Some("edge-delete-hard".to_string()),
546 }
547}
548
549pub fn purge_incident_edges_statement(node_id: Uuid) -> SqlStatement {
551 SqlStatement {
552 sql: "DELETE FROM graph_edges WHERE source_id = ?1 OR target_id = ?1".to_string(),
553 params: vec![SqlValue::Text(node_id.to_string())],
554 label: Some("edge-purge-incident".to_string()),
555 }
556}
557
558pub const EDGE_SYMMETRIC_CONFLICT_PROBE_SQL: &str = "SELECT id FROM graph_edges \
579 WHERE namespace = ?1 AND source_id = ?2 AND target_id = ?3 \
580 AND relation = ?4 AND id != ?5";
581
582pub const EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL: &str =
583 "DELETE FROM graph_edges WHERE namespace = ?1 AND id = ?2";
584
585pub const EDGE_SYMMETRIC_DELETE_NONCANONICAL_GUARDED_SQL: &str =
598 "DELETE FROM graph_edges WHERE namespace = ?1 AND id = ?2 \
599 AND updated_at = ?3 AND deleted_at IS ?4";
600
601pub const EDGE_SYMMETRIC_UPDATE_INPLACE_SQL: &str = "UPDATE graph_edges SET \
608 source_id = ?1, target_id = ?2, relation = ?3, \
609 weight = ?4, updated_at = ?5, metadata = ?6 \
610 WHERE namespace = ?7 AND id = ?8 \
611 AND updated_at = ?9 AND deleted_at IS ?10 \
612 AND ?5 > updated_at";
613
614pub fn edge_symmetric_conflict_probe_statement(
617 namespace: &str,
618 canon_src: Uuid,
619 canon_tgt: Uuid,
620 relation: EdgeRelation,
621 exclude_id: Uuid,
622) -> SqlStatement {
623 SqlStatement {
624 sql: EDGE_SYMMETRIC_CONFLICT_PROBE_SQL.to_string(),
625 params: vec![
626 SqlValue::Text(namespace.to_string()),
627 SqlValue::Text(canon_src.to_string()),
628 SqlValue::Text(canon_tgt.to_string()),
629 SqlValue::Text(relation.to_string()),
630 SqlValue::Text(exclude_id.to_string()),
631 ],
632 label: Some("edge-symmetric-conflict-probe".to_string()),
633 }
634}
635
636pub fn edge_symmetric_delete_noncanonical_statement(namespace: &str, id: Uuid) -> SqlStatement {
639 SqlStatement {
640 sql: EDGE_SYMMETRIC_DELETE_NONCANONICAL_SQL.to_string(),
641 params: vec![
642 SqlValue::Text(namespace.to_string()),
643 SqlValue::Text(id.to_string()),
644 ],
645 label: Some("edge-symmetric-delete-noncanonical".to_string()),
646 }
647}
648
649#[allow(clippy::too_many_arguments)]
653pub fn edge_symmetric_update_inplace_statement(
654 namespace: &str,
655 id: Uuid,
656 canon_src: Uuid,
657 canon_tgt: Uuid,
658 relation: EdgeRelation,
659 weight: f64,
660 updated_at_micros: i64,
661 metadata: Option<&str>,
662 expected_updated_at_micros: i64,
663 expected_deleted_at_micros: Option<i64>,
664) -> SqlStatement {
665 SqlStatement {
666 sql: EDGE_SYMMETRIC_UPDATE_INPLACE_SQL.to_string(),
667 params: vec![
668 SqlValue::Text(canon_src.to_string()),
669 SqlValue::Text(canon_tgt.to_string()),
670 SqlValue::Text(relation.to_string()),
671 SqlValue::Float(weight),
672 SqlValue::Integer(updated_at_micros),
673 match metadata {
674 Some(m) => SqlValue::Text(m.to_string()),
675 None => SqlValue::Null,
676 },
677 SqlValue::Text(namespace.to_string()),
678 SqlValue::Text(id.to_string()),
679 SqlValue::Integer(expected_updated_at_micros),
680 match expected_deleted_at_micros {
681 Some(value) => SqlValue::Integer(value),
682 None => SqlValue::Null,
683 },
684 ],
685 label: Some("edge-symmetric-update-inplace".to_string()),
686 }
687}
688
689#[allow(clippy::too_many_arguments)]
783pub fn edge_symmetric_delete_if_conflict_statement(
784 namespace: &str,
785 id: Uuid,
786 canon_src: Uuid,
787 canon_tgt: Uuid,
788 relation: EdgeRelation,
789 expected_updated_at_micros: i64,
790 expected_deleted_at_micros: Option<i64>,
791) -> SqlStatement {
792 SqlStatement {
793 sql: "DELETE FROM graph_edges \
794 WHERE namespace = ?1 AND id = ?2 \
795 AND updated_at = ?6 AND deleted_at IS ?7 \
796 AND EXISTS ( \
797 SELECT 1 FROM graph_edges \
798 WHERE namespace = ?1 AND source_id = ?3 AND target_id = ?4 \
799 AND relation = ?5 AND id != ?2 \
800 )"
801 .to_string(),
802 params: vec![
803 SqlValue::Text(namespace.to_string()),
804 SqlValue::Text(id.to_string()),
805 SqlValue::Text(canon_src.to_string()),
806 SqlValue::Text(canon_tgt.to_string()),
807 SqlValue::Text(relation.to_string()),
808 SqlValue::Integer(expected_updated_at_micros),
809 match expected_deleted_at_micros {
810 Some(value) => SqlValue::Integer(value),
811 None => SqlValue::Null,
812 },
813 ],
814 label: Some("edge-symmetric-delete-if-conflict".to_string()),
815 }
816}
817
818#[allow(clippy::too_many_arguments)]
827pub fn edge_symmetric_absorb_or_update_inplace_statement(
828 namespace: &str,
829 id: Uuid,
830 canon_src: Uuid,
831 canon_tgt: Uuid,
832 relation: EdgeRelation,
833 weight: f64,
834 updated_at_micros: i64,
835 metadata: Option<&str>,
836 target_backend: Option<&str>,
837 expected_updated_at_micros: i64,
838 expected_deleted_at_micros: Option<i64>,
839) -> SqlStatement {
840 SqlStatement {
841 sql: "UPDATE graph_edges SET \
842 source_id = CASE WHEN id = ?2 THEN ?3 ELSE source_id END, \
843 target_id = CASE WHEN id = ?2 THEN ?4 ELSE target_id END, \
844 relation = CASE WHEN id = ?2 THEN ?5 ELSE relation END, \
845 weight = CASE WHEN id = ?2 THEN ?6 ELSE weight END, \
846 updated_at = CASE WHEN id = ?2 THEN ?7 ELSE updated_at END, \
847 deleted_at = CASE WHEN id = ?2 THEN NULL ELSE deleted_at END, \
848 metadata = CASE WHEN id = ?2 THEN ?8 ELSE metadata END, \
849 target_backend = CASE WHEN id = ?2 THEN ?9 ELSE target_backend END \
850 WHERE namespace = ?1 \
851 AND ( \
852 (id = ?2 AND changes() = 0 AND updated_at = ?10 AND deleted_at IS ?11 \
853 AND ?7 > updated_at) \
854 OR (source_id = ?3 AND target_id = ?4 AND relation = ?5 \
855 AND id != ?2 AND changes() = 1) \
856 )"
857 .to_string(),
858 params: vec![
859 SqlValue::Text(namespace.to_string()),
860 SqlValue::Text(id.to_string()),
861 SqlValue::Text(canon_src.to_string()),
862 SqlValue::Text(canon_tgt.to_string()),
863 SqlValue::Text(relation.to_string()),
864 SqlValue::Float(weight),
865 SqlValue::Integer(updated_at_micros),
866 match metadata {
867 Some(m) => SqlValue::Text(m.to_string()),
868 None => SqlValue::Null,
869 },
870 match target_backend {
871 Some(b) => SqlValue::Text(b.to_string()),
872 None => SqlValue::Null,
873 },
874 SqlValue::Integer(expected_updated_at_micros),
875 match expected_deleted_at_micros {
876 Some(value) => SqlValue::Integer(value),
877 None => SqlValue::Null,
878 },
879 ],
880 label: Some("edge-symmetric-absorb-or-update-inplace".to_string()),
881 }
882}
883
884#[doc(hidden)]
888pub struct GraphDocumentGuard {
889 pub namespace: String,
890 pub id: Uuid,
891 pub expected_blob_ref: String,
892}
893
894#[doc(hidden)]
898pub struct GraphEdgeSnapshotGuard {
899 pub namespace: String,
900 pub source_id: Uuid,
901 pub target_id: Uuid,
902 pub relation: EdgeRelation,
903 pub expected: Option<Edge>,
904}
905
906#[doc(hidden)]
910#[derive(Default)]
911pub struct GraphMutationPreconditions {
912 pub document: Option<GraphDocumentGuard>,
913 pub edges: Vec<GraphEdgeSnapshotGuard>,
914}
915
916#[doc(hidden)]
920pub enum GraphMutationRequest {
921 Single {
922 request: EdgeUpsertRequest,
923 guard_endpoints: bool,
924 },
925 Batch {
926 requests: Vec<EdgeUpsertRequest>,
927 guard_endpoints: bool,
928 },
929 CommitAnnotation {
930 edge: Edge,
931 guard: CommitAnnotationGuard,
932 },
933}
934
935#[doc(hidden)]
939pub enum GraphMutationOutcome {
940 Single(GuardedEdgeUpsertOutcome),
941 Batch(GuardedEdgeBatchUpsertOutcome),
942 CommitAnnotation(CommitAnnotationInsertOutcome),
943}
944
945#[doc(hidden)]
949pub struct GraphMutationEventOutcome {
950 pub mutation: GraphMutationOutcome,
951 pub retired: Vec<Edge>,
952}
953
954const GRAPH_MUTATION_EVENTS_OP: &str = "compose_graph_mutation_events";
955
956fn count_graph_mutation_events(outcome: &GraphMutationEventOutcome) {
957 let written = match &outcome.mutation {
958 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(_)) => 1,
959 GraphMutationOutcome::Batch(outcome) if outcome.refusal.is_none() => outcome.rows.len(),
960 GraphMutationOutcome::CommitAnnotation(CommitAnnotationInsertOutcome::Created(_)) => 1,
961 _ => 0,
962 };
963 khive_storage::usage::count(
964 khive_storage::usage::UsageUnit::EventRows,
965 (written + outcome.retired.len()) as u64,
966 );
967}
968
969#[doc(hidden)]
973pub async fn compose_graph_mutation_events<F>(
974 backend: &crate::StorageBackend,
975 mutation: GraphMutationRequest,
976 preconditions: GraphMutationPreconditions,
977 retirements: Vec<Edge>,
978 make_events: F,
979) -> StorageResult<GraphMutationEventOutcome>
980where
981 F: FnOnce(&GraphMutationEventOutcome) -> StorageResult<Vec<Event>> + Send + 'static,
982{
983 if backend.is_read_only() {
984 return Err(StorageError::Pool {
985 operation: GRAPH_MUTATION_EVENTS_OP.into(),
986 message: "backend is read-only".into(),
987 });
988 }
989 backend
992 .graph()
993 .map_err(|error| map_sqlite_err(error, GRAPH_MUTATION_EVENTS_OP))?;
994 backend
995 .events()
996 .map_err(|error| map_sqlite_err(error, GRAPH_MUTATION_EVENTS_OP))?;
997 let pool = backend.pool_arc();
998 if let Some(writer_task) = pool.writer_task_for_write(None, GRAPH_MUTATION_EVENTS_OP)? {
999 return writer_task
1000 .send_bounded(move |conn| {
1001 graph_mutation_events_enlisted(
1002 conn,
1003 mutation,
1004 preconditions,
1005 retirements,
1006 make_events,
1007 )
1008 })
1009 .await
1010 .inspect(count_graph_mutation_events)
1011 .inspect_err(|error| khive_storage::usage::account_event_write(Err(error)));
1012 }
1013 pool.record_direct_route(crate::timeout_sink::Site::DirectRouteGraphGeneralWrite);
1014 let is_file_backed = backend.is_file_backed();
1015 tokio::task::spawn_blocking(move || {
1016 if is_file_backed {
1017 let admission = pool.write_admission();
1018 let _lease = admission
1019 .acquire()
1020 .map_err(|error| map_sqlite_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1021 let conn = pool
1022 .open_standalone_writer_for_admitted_operation()
1023 .map_err(|error| map_sqlite_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1024 run_graph_mutation_transaction(&pool, &conn, false, move |conn| {
1025 graph_mutation_events_enlisted(
1026 conn,
1027 mutation,
1028 preconditions,
1029 retirements,
1030 make_events,
1031 )
1032 })
1033 } else {
1034 let guard = pool
1035 .writer_for_admitted_operation()
1036 .map_err(|error| map_sqlite_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1037 run_graph_mutation_transaction(&pool, guard.conn(), true, move |conn| {
1038 graph_mutation_events_enlisted(
1039 conn,
1040 mutation,
1041 preconditions,
1042 retirements,
1043 make_events,
1044 )
1045 })
1046 }
1047 })
1048 .await
1049 .map_err(|error| {
1050 StorageError::driver(StorageCapability::Graph, GRAPH_MUTATION_EVENTS_OP, error)
1051 })?
1052 .inspect(count_graph_mutation_events)
1053 .inspect_err(|error| khive_storage::usage::account_event_write(Err(error)))
1054}
1055
1056fn graph_mutation_conflict(message: &'static str) -> StorageError {
1057 StorageError::Conflict {
1058 capability: StorageCapability::Graph,
1059 operation: GRAPH_MUTATION_EVENTS_OP.into(),
1060 message: message.into(),
1061 }
1062}
1063
1064fn edge_snapshot_matches(actual: &Edge, expected: &Edge) -> bool {
1065 actual.id == expected.id
1066 && actual.namespace == expected.namespace
1067 && actual.source_id == expected.source_id
1068 && actual.target_id == expected.target_id
1069 && actual.relation == expected.relation
1070 && actual.weight.to_bits() == expected.weight.to_bits()
1071 && actual.metadata == expected.metadata
1072 && actual.target_backend == expected.target_backend
1073 && actual.created_at == expected.created_at
1074 && actual.updated_at == expected.updated_at
1075 && actual.deleted_at == expected.deleted_at
1076}
1077
1078fn check_graph_mutation_preconditions(
1079 conn: &rusqlite::Connection,
1080 preconditions: &GraphMutationPreconditions,
1081 retirements: &[Edge],
1082) -> StorageResult<()> {
1083 if let Some(document) = &preconditions.document {
1084 let matches: bool = conn
1085 .query_row(
1086 "SELECT EXISTS(SELECT 1 FROM entities WHERE id=?1 AND namespace=?2 \
1087 AND deleted_at IS NULL AND json_type(properties, '$.blob_ref')='text' \
1088 AND json_extract(properties, '$.blob_ref')=?3 COLLATE BINARY)",
1089 rusqlite::params![
1090 document.id.to_string(),
1091 &document.namespace,
1092 &document.expected_blob_ref,
1093 ],
1094 |row| row.get(0),
1095 )
1096 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1097 if !matches {
1098 return Err(graph_mutation_conflict(
1099 "source document body reference changed",
1100 ));
1101 }
1102 }
1103 for guard in &preconditions.edges {
1104 let actual = edge_by_natural_key_parts_including_deleted(
1105 conn,
1106 &guard.namespace,
1107 guard.source_id,
1108 guard.target_id,
1109 guard.relation,
1110 )
1111 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1112 let matches = match (actual.as_ref(), guard.expected.as_ref()) {
1113 (None, None) => true,
1114 (Some(actual), Some(expected)) => edge_snapshot_matches(actual, expected),
1115 _ => false,
1116 };
1117 if !matches {
1118 return Err(graph_mutation_conflict("edge ownership snapshot changed"));
1119 }
1120 }
1121 for expected in retirements {
1122 let actual = edge_by_natural_key_including_deleted(conn, expected)
1123 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1124 if expected.deleted_at.is_some()
1125 || !actual
1126 .as_ref()
1127 .is_some_and(|actual| edge_snapshot_matches(actual, expected))
1128 {
1129 return Err(graph_mutation_conflict("edge retirement snapshot changed"));
1130 }
1131 }
1132 Ok(())
1133}
1134
1135fn graph_mutation_events_enlisted<F>(
1137 conn: &rusqlite::Connection,
1138 mutation: GraphMutationRequest,
1139 preconditions: GraphMutationPreconditions,
1140 retirements: Vec<Edge>,
1141 make_events: F,
1142) -> StorageResult<GraphMutationEventOutcome>
1143where
1144 F: FnOnce(&GraphMutationEventOutcome) -> StorageResult<Vec<Event>>,
1145{
1146 if !retirements.is_empty() && !matches!(&mutation, GraphMutationRequest::Batch { .. }) {
1147 return Err(StorageError::InvalidInput {
1148 capability: StorageCapability::Graph,
1149 operation: GRAPH_MUTATION_EVENTS_OP.into(),
1150 message: "retirements require a batch mutation".into(),
1151 });
1152 }
1153 check_graph_mutation_preconditions(conn, &preconditions, &retirements)?;
1154 let mutation = match mutation {
1155 GraphMutationRequest::Single {
1156 request,
1157 guard_endpoints,
1158 } => GraphMutationOutcome::Single(
1159 observed_edge_upsert(conn, &request, guard_endpoints)
1160 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?,
1161 ),
1162 GraphMutationRequest::Batch {
1163 requests,
1164 guard_endpoints,
1165 } => {
1166 let outcome = observed_edge_batch_upsert(conn, &requests, guard_endpoints)
1167 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1168 if outcome.refusal.is_none() && outcome.rows.len() != requests.len() {
1169 return Err(graph_mutation_conflict(
1170 "edge result count differs from request count",
1171 ));
1172 }
1173 GraphMutationOutcome::Batch(outcome)
1174 }
1175 GraphMutationRequest::CommitAnnotation { edge, guard } => {
1176 if edge.relation != EdgeRelation::Annotates || edge.deleted_at.is_some() {
1177 return Err(StorageError::InvalidInput {
1178 capability: StorageCapability::Graph,
1179 operation: GRAPH_MUTATION_EVENTS_OP.into(),
1180 message: "expected a live annotates edge".into(),
1181 });
1182 }
1183 GraphMutationOutcome::CommitAnnotation(
1184 conditional_commit_annotation_insert(conn, edge, &guard)
1185 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?,
1186 )
1187 }
1188 };
1189 let written = match &mutation {
1190 GraphMutationOutcome::Single(GuardedEdgeUpsertOutcome::Written(_)) => 1,
1191 GraphMutationOutcome::Batch(outcome) if outcome.refusal.is_none() => outcome.rows.len(),
1192 GraphMutationOutcome::CommitAnnotation(CommitAnnotationInsertOutcome::Created(_)) => 1,
1193 _ => {
1194 return Ok(GraphMutationEventOutcome {
1195 mutation,
1196 retired: Vec::new(),
1197 })
1198 }
1199 };
1200 let mut retired = Vec::with_capacity(retirements.len());
1201 for edge in retirements {
1202 let statement =
1203 edge_soft_delete_statement(Uuid::from(edge.id), Utc::now().timestamp_micros());
1204 let mut stmt = conn
1205 .prepare(&statement.sql)
1206 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1207 bind_params(&mut stmt, &statement.params)
1208 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?;
1209 if stmt
1210 .raw_execute()
1211 .map_err(|error| map_err(error, GRAPH_MUTATION_EVENTS_OP))?
1212 != 1
1213 {
1214 return Err(graph_mutation_conflict(
1215 "edge retirement changed during mutation",
1216 ));
1217 }
1218 retired.push(edge);
1219 }
1220 let outcome = GraphMutationEventOutcome { mutation, retired };
1221 let expected_events = written + outcome.retired.len();
1222 if expected_events == 0 {
1223 return Ok(outcome);
1224 }
1225 let events = make_events(&outcome)?;
1226 if events.len() != expected_events {
1227 return Err(graph_mutation_conflict(
1228 "event result count differs from mutation count",
1229 ));
1230 }
1231 for event in &events {
1232 super::event::append_event_in_transaction(conn, event).map_err(|error| {
1233 StorageError::driver(StorageCapability::Events, GRAPH_MUTATION_EVENTS_OP, error)
1234 })?;
1235 }
1236 Ok(outcome)
1237}
1238
1239pub struct SqlGraphStore {
1241 pool: Arc<ConnectionPool>,
1242 index_repair: Option<super::index_repair::IndexRepairContext>,
1243 is_file_backed: bool,
1244 namespace: String,
1248 writer_task: Option<WriterTaskHandle>,
1249}
1250
1251impl SqlGraphStore {
1252 pub fn new_scoped(
1260 pool: Arc<ConnectionPool>,
1261 is_file_backed: bool,
1262 namespace: impl Into<String>,
1263 ) -> Self {
1264 let writer_task = pool.writer_task_handle().ok().flatten();
1270
1271 Self {
1272 pool,
1273 index_repair: None,
1274 is_file_backed,
1275 namespace: namespace.into(),
1276 writer_task,
1277 }
1278 }
1279
1280 pub(crate) fn with_index_repair(
1281 mut self,
1282 repair: super::index_repair::IndexRepairContext,
1283 ) -> Self {
1284 self.index_repair = Some(repair);
1285 self
1286 }
1287
1288 async fn with_indexed_reader<F, R>(&self, op: &'static str, read: F) -> Result<R, StorageError>
1289 where
1290 F: FnMut(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
1291 R: Send + 'static,
1292 {
1293 super::index_repair::run_indexed_read(
1294 Arc::clone(&self.pool),
1295 self.index_repair.clone(),
1296 StorageCapability::Graph,
1297 op,
1298 read,
1299 )
1300 .await
1301 }
1302
1303 fn current_writer_task(
1304 &self,
1305 operation: &'static str,
1306 ) -> Result<Option<WriterTaskHandle>, StorageError> {
1307 self.pool
1308 .writer_task_for_write(self.writer_task.as_ref(), operation)
1309 }
1310
1311 async fn with_writer<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
1318 where
1319 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
1320 R: Send + 'static,
1321 {
1322 if let Some(writer_task) = self.current_writer_task(op)? {
1323 return writer_task
1324 .send_bounded(move |conn| f(conn).map_err(|e| map_err(e, op)))
1325 .await;
1326 }
1327
1328 self.pool
1329 .record_direct_route(crate::timeout_sink::Site::DirectRouteGraphGeneralWrite);
1330 let pool = Arc::clone(&self.pool);
1331 let db = crate::timeout_sink::db_label(&pool);
1332 let is_file_backed = self.is_file_backed;
1333 tokio::task::spawn_blocking(move || {
1334 let result =
1335 pool.execute_direct_transaction(StorageCapability::Graph, op, move |conn| {
1336 f(conn).map_err(|error| map_err(error, op))
1337 });
1338 if is_file_backed {
1339 if let Err(error) = &result {
1340 crate::timeout_sink::maybe_emit_busy_storage_error(
1341 &db,
1342 crate::timeout_sink::Site::StandaloneGraph,
1343 error,
1344 );
1345 }
1346 }
1347 result
1348 })
1349 .await
1350 .map_err(|e| StorageError::driver(StorageCapability::Graph, op, e))?
1351 }
1352
1353 async fn observed_edge_write(
1354 &self,
1355 operation: &'static str,
1356 request: EdgeUpsertRequest,
1357 guard_endpoints: bool,
1358 ) -> Result<GuardedEdgeUpsertOutcome, StorageError> {
1359 if let Some(writer_task) = self.current_writer_task(operation)? {
1360 return writer_task
1361 .send_bounded(move |conn| {
1362 observed_edge_upsert(conn, &request, guard_endpoints)
1363 .map_err(|error| map_err(error, operation))
1364 })
1365 .await;
1366 }
1367
1368 self.with_writer(operation, move |conn| {
1369 observed_edge_upsert(conn, &request, guard_endpoints)
1370 })
1371 .await
1372 }
1373
1374 async fn observed_edge_batch_write(
1375 &self,
1376 operation: &'static str,
1377 requests: Vec<EdgeUpsertRequest>,
1378 guard_endpoints: bool,
1379 ) -> Result<GuardedEdgeBatchUpsertOutcome, StorageError> {
1380 if let Some(writer_task) = self.current_writer_task(operation)? {
1381 return writer_task
1382 .send_bounded(move |conn| {
1383 observed_edge_batch_upsert(conn, &requests, guard_endpoints)
1384 .map_err(|error| map_err(error, operation))
1385 })
1386 .await;
1387 }
1388
1389 self.with_writer(operation, move |conn| {
1390 observed_edge_batch_upsert(conn, &requests, guard_endpoints)
1391 })
1392 .await
1393 }
1394
1395 async fn with_reader<F, R>(&self, op: &'static str, f: F) -> Result<R, StorageError>
1396 where
1397 F: FnOnce(&rusqlite::Connection) -> Result<R, rusqlite::Error> + Send + 'static,
1398 R: Send + 'static,
1399 {
1400 super::run_pooled_store_read(
1401 Arc::clone(&self.pool),
1402 StorageCapability::Graph,
1403 op,
1404 move |conn| f(conn).map_err(|error| map_err(error, op)),
1405 )
1406 .await
1407 }
1408}
1409
1410fn report_graph_usage(queries: &std::sync::atomic::AtomicU64, rows: &std::sync::atomic::AtomicU64) {
1426 khive_storage::usage::count(
1427 khive_storage::usage::UsageUnit::DbRoundTrips,
1428 queries.load(std::sync::atomic::Ordering::Relaxed),
1429 );
1430 khive_storage::usage::count(
1431 khive_storage::usage::UsageUnit::GraphHops,
1432 rows.load(std::sync::atomic::Ordering::Relaxed),
1433 );
1434}
1435
1436fn read_edge(row: &rusqlite::Row<'_>) -> Result<Edge, rusqlite::Error> {
1437 let namespace: String = row.get(0)?;
1438 let id_str: String = row.get(1)?;
1439 let source_str: String = row.get(2)?;
1440 let target_str: String = row.get(3)?;
1441 let relation_str: String = row.get(4)?;
1442 let weight: f64 = row.get(5)?;
1443 let created_micros: i64 = row.get(6)?;
1444 let updated_micros: i64 = row.get(7)?;
1445 let deleted_micros: Option<i64> = row.get(8)?;
1446 let metadata_str: Option<String> = row.get(9)?;
1447 let target_backend: Option<String> = row.get(10)?;
1448
1449 let id = parse_uuid(&id_str)?;
1450 let source_id = parse_uuid(&source_str)?;
1451 let target_id = parse_uuid(&target_str)?;
1452 let created_at = micros_to_datetime(created_micros);
1453 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
1454 rusqlite::Error::FromSqlConversionFailure(4, rusqlite::types::Type::Text, Box::new(e))
1455 })?;
1456 let metadata = match metadata_str {
1457 Some(s) => {
1458 let v = serde_json::from_str(&s).map_err(|e| {
1459 rusqlite::Error::FromSqlConversionFailure(
1460 9,
1461 rusqlite::types::Type::Text,
1462 Box::new(e),
1463 )
1464 })?;
1465 Some(v)
1466 }
1467 None => None,
1468 };
1469
1470 Ok(Edge {
1471 id: id.into(),
1472 namespace,
1473 source_id,
1474 target_id,
1475 relation,
1476 weight,
1477 created_at,
1478 updated_at: micros_to_datetime(updated_micros),
1479 deleted_at: deleted_micros.map(micros_to_datetime),
1480 metadata,
1481 target_backend,
1482 })
1483}
1484
1485fn parse_uuid(s: &str) -> Result<Uuid, rusqlite::Error> {
1486 Uuid::parse_str(s).map_err(|e| {
1487 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
1488 })
1489}
1490
1491fn neighbor_extra_clause(
1497 query: &NeighborQuery,
1498 start_param_idx: usize,
1499 after: Option<&NeighborCursor>,
1500 neighbor_kinds: Option<&[String]>,
1501) -> (String, String, Vec<Box<dyn rusqlite::types::ToSql>>) {
1502 let mut conditions: Vec<String> = Vec::new();
1503 let mut extra_params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
1504 let mut param_idx = start_param_idx;
1505
1506 if let Some(ref rels) = query.relations {
1507 if !rels.is_empty() {
1508 let placeholders: Vec<String> = rels
1509 .iter()
1510 .map(|r| {
1511 extra_params.push(Box::new(r.to_string()));
1512 let p = format!("?{}", param_idx);
1513 param_idx += 1;
1514 p
1515 })
1516 .collect();
1517 conditions.push(format!("relation IN ({})", placeholders.join(",")));
1518 }
1519 }
1520
1521 if let Some(min_w) = query.min_weight {
1522 extra_params.push(Box::new(min_w));
1523 conditions.push(format!("weight >= ?{}", param_idx));
1524 param_idx += 1;
1525 }
1526
1527 if let Some(cursor) = after {
1528 extra_params.push(Box::new(cursor.weight));
1529 let weight_idx = param_idx;
1530 param_idx += 1;
1531 extra_params.push(Box::new(cursor.node_id.to_string()));
1532 let node_idx = param_idx;
1533 param_idx += 1;
1534 extra_params.push(Box::new(cursor.edge_id.to_string()));
1535 let edge_idx = param_idx;
1536 param_idx += 1;
1537 conditions.push(format!(
1538 "(weight < ?{weight_idx} OR (weight = ?{weight_idx} AND node_id > ?{node_idx}) OR (weight = ?{weight_idx} AND node_id = ?{node_idx} AND edge_id > ?{edge_idx}))"
1539 ));
1540 }
1541
1542 if let Some(kinds) = neighbor_kinds.filter(|kinds| !kinds.is_empty()) {
1543 let placeholders: Vec<String> = kinds
1544 .iter()
1545 .map(|kind| {
1546 extra_params.push(Box::new(kind.clone()));
1547 let p = format!("?{param_idx}");
1548 param_idx += 1;
1549 p
1550 })
1551 .collect();
1552 let entity_placeholders = placeholders.join(",");
1553 let note_placeholders: Vec<String> = kinds
1554 .iter()
1555 .map(|kind| {
1556 extra_params.push(Box::new(kind.clone()));
1557 let p = format!("?{param_idx}");
1558 param_idx += 1;
1559 p
1560 })
1561 .collect();
1562 conditions.push(format!(
1563 "(EXISTS (SELECT 1 FROM entities AS neighbor_entities WHERE neighbor_entities.id = node_id AND neighbor_entities.namespace = ?1 AND neighbor_entities.deleted_at IS NULL AND neighbor_entities.kind IN ({entity_placeholders})) OR EXISTS (SELECT 1 FROM notes AS neighbor_notes WHERE neighbor_notes.id = node_id AND neighbor_notes.namespace = ?1 AND neighbor_notes.deleted_at IS NULL AND neighbor_notes.kind IN ({})))",
1564 note_placeholders.join(",")
1565 ));
1566 }
1567
1568 let where_extra = if conditions.is_empty() {
1569 String::new()
1570 } else {
1571 format!(" WHERE {}", conditions.join(" AND "))
1572 };
1573
1574 let limit_clause = if let Some(lim) = query.limit {
1575 extra_params.push(Box::new(lim as i64));
1576 format!(" LIMIT ?{}", param_idx)
1577 } else {
1578 String::new()
1579 };
1580
1581 (where_extra, limit_clause, extra_params)
1582}
1583
1584#[cfg(test)]
1592static NEIGHBOR_SELECT_COUNT: std::sync::atomic::AtomicUsize =
1593 std::sync::atomic::AtomicUsize::new(0);
1594
1595#[cfg(test)]
1596fn count_neighbor_select() {
1597 NEIGHBOR_SELECT_COUNT.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1598}
1599
1600#[cfg(not(test))]
1601fn count_neighbor_select() {}
1602
1603#[cfg(test)]
1604pub(crate) fn reset_neighbor_select_count() {
1605 NEIGHBOR_SELECT_COUNT.store(0, std::sync::atomic::Ordering::Relaxed);
1606}
1607
1608#[cfg(test)]
1609pub(crate) fn neighbor_select_count() -> usize {
1610 NEIGHBOR_SELECT_COUNT.load(std::sync::atomic::Ordering::Relaxed)
1611}
1612
1613fn micros_to_datetime(micros: i64) -> DateTime<Utc> {
1614 Utc.timestamp_micros(micros)
1615 .single()
1616 .unwrap_or_else(Utc::now)
1617}
1618
1619fn edge_order_clause(sort: &[SortOrder<EdgeSortField>]) -> String {
1626 if sort.is_empty() {
1627 return " ORDER BY created_at DESC, id DESC".to_string();
1628 }
1629 let mut parts: Vec<String> = sort
1630 .iter()
1631 .map(|s| {
1632 let dir = match s.direction {
1633 SortDirection::Asc => "ASC",
1634 SortDirection::Desc => "DESC",
1635 };
1636 format!("{} {}", edge_sort_col(&s.field), dir)
1637 })
1638 .collect();
1639 let dir = match sort.last().map(|s| &s.direction) {
1640 Some(SortDirection::Asc) => "ASC",
1641 _ => "DESC",
1642 };
1643 parts.push(format!("id {dir}"));
1644 format!(" ORDER BY {}", parts.join(", "))
1645}
1646
1647fn endpoint_base_case(column: &str) -> String {
1665 format!(
1666 "CASE WHEN EXISTS (SELECT 1 FROM entities be WHERE be.id = graph_edges.{column}) \
1667 THEN 'entity' \
1668 WHEN EXISTS (SELECT 1 FROM notes bn WHERE bn.id = graph_edges.{column}) \
1669 THEN 'note' ELSE 'none' END"
1670 )
1671}
1672
1673fn fold_endpoint_base_row(counts: &mut EdgeEndpointBaseCounts, source: &str, target: &str, n: u64) {
1677 let slot = match (source, target) {
1678 ("entity", "entity") => &mut counts.entity_entity,
1679 ("entity", "note") => &mut counts.entity_note,
1680 ("note", "entity") => &mut counts.note_entity,
1681 ("note", "note") => &mut counts.note_note,
1682 _ => &mut counts.unresolved,
1683 };
1684 *slot = slot.saturating_add(n);
1685}
1686
1687const LIVE_ENDPOINTS_CONDITION: &str = "NOT EXISTS (SELECT 1 FROM entities le \
1700 WHERE le.id = graph_edges.source_id AND le.deleted_at IS NOT NULL) \
1701 AND NOT EXISTS (SELECT 1 FROM entities le \
1702 WHERE le.id = graph_edges.target_id AND le.deleted_at IS NOT NULL) \
1703 AND NOT EXISTS (SELECT 1 FROM notes ln \
1704 WHERE ln.id = graph_edges.source_id AND ln.deleted_at IS NOT NULL) \
1705 AND NOT EXISTS (SELECT 1 FROM notes ln \
1706 WHERE ln.id = graph_edges.target_id AND ln.deleted_at IS NOT NULL)";
1707
1708fn with_live_endpoints(where_clause: &str) -> String {
1711 format!("{where_clause} AND {LIVE_ENDPOINTS_CONDITION}")
1712}
1713
1714fn build_edge_filter_sql(
1715 namespace: &str,
1716 filter: &EdgeFilter,
1717) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
1718 build_edge_filter_sql_for_namespaces(&[namespace.to_string()], filter)
1719}
1720
1721fn build_edge_filter_sql_for_namespaces(
1722 namespaces: &[String],
1723 filter: &EdgeFilter,
1724) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
1725 let params: Vec<Box<dyn rusqlite::types::ToSql>> = namespaces
1726 .iter()
1727 .map(|namespace| -> Box<dyn rusqlite::types::ToSql> { Box::new(namespace.clone()) })
1728 .collect();
1729 let namespace_condition = match namespaces.len() {
1730 0 => "0".to_string(),
1731 1 => "namespace = ?1".to_string(),
1732 _ => {
1733 let placeholders: Vec<String> =
1734 (1..=namespaces.len()).map(|i| format!("?{i}")).collect();
1735 format!("namespace IN ({})", placeholders.join(", "))
1736 }
1737 };
1738 build_edge_filter_conditions(namespace_condition, params, filter)
1739}
1740
1741fn build_edge_filter_sql_for_namespaces_json(
1756 namespaces_json: &str,
1757 filter: &EdgeFilter,
1758) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
1759 let params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![Box::new(namespaces_json.to_string())];
1760 let namespace_condition = "namespace IN (SELECT value FROM json_each(?1))".to_string();
1761 build_edge_filter_conditions(namespace_condition, params, filter)
1762}
1763
1764fn build_edge_filter_conditions(
1765 namespace_condition: String,
1766 mut params: Vec<Box<dyn rusqlite::types::ToSql>>,
1767 filter: &EdgeFilter,
1768) -> (String, Vec<Box<dyn rusqlite::types::ToSql>>) {
1769 let mut conditions = vec![namespace_condition, "deleted_at IS NULL".to_string()];
1770
1771 if !filter.ids.is_empty() {
1772 let placeholders: Vec<String> = filter
1773 .ids
1774 .iter()
1775 .map(|id| {
1776 params.push(Box::new(id.to_string()));
1777 format!("?{}", params.len())
1778 })
1779 .collect();
1780 conditions.push(format!("id IN ({})", placeholders.join(",")));
1781 }
1782
1783 if !filter.source_ids.is_empty() {
1784 let placeholders: Vec<String> = filter
1785 .source_ids
1786 .iter()
1787 .map(|id| {
1788 params.push(Box::new(id.to_string()));
1789 format!("?{}", params.len())
1790 })
1791 .collect();
1792 conditions.push(format!("source_id IN ({})", placeholders.join(",")));
1793 }
1794
1795 if !filter.target_ids.is_empty() {
1796 let placeholders: Vec<String> = filter
1797 .target_ids
1798 .iter()
1799 .map(|id| {
1800 params.push(Box::new(id.to_string()));
1801 format!("?{}", params.len())
1802 })
1803 .collect();
1804 conditions.push(format!("target_id IN ({})", placeholders.join(",")));
1805 }
1806
1807 if !filter.relations.is_empty() {
1808 let placeholders: Vec<String> = filter
1809 .relations
1810 .iter()
1811 .map(|r| {
1812 params.push(Box::new(r.to_string()));
1813 format!("?{}", params.len())
1814 })
1815 .collect();
1816 conditions.push(format!("relation IN ({})", placeholders.join(",")));
1817 }
1818
1819 if let Some(min_w) = filter.min_weight {
1820 params.push(Box::new(min_w));
1821 conditions.push(format!("weight >= ?{}", params.len()));
1822 }
1823
1824 if let Some(max_w) = filter.max_weight {
1825 params.push(Box::new(max_w));
1826 conditions.push(format!("weight <= ?{}", params.len()));
1827 }
1828
1829 if let Some(ref time_range) = filter.created_at {
1830 if let Some(start) = time_range.start {
1831 params.push(Box::new(start.timestamp_micros()));
1832 conditions.push(format!("created_at >= ?{}", params.len()));
1833 }
1834 if let Some(end) = time_range.end {
1835 params.push(Box::new(end.timestamp_micros()));
1836 conditions.push(format!("created_at < ?{}", params.len()));
1837 }
1838 }
1839
1840 let clause = format!(" WHERE {}", conditions.join(" AND "));
1841 (clause, params)
1842}
1843
1844fn edge_sort_col(field: &EdgeSortField) -> &'static str {
1845 match field {
1846 EdgeSortField::CreatedAt => "created_at",
1847 EdgeSortField::Weight => "weight",
1848 EdgeSortField::Relation => "relation",
1849 }
1850}
1851
1852fn canonical_edge_endpoints(
1861 relation: EdgeRelation,
1862 source_id: Uuid,
1863 target_id: Uuid,
1864) -> (Uuid, Uuid) {
1865 relation.canonical_endpoints(source_id, target_id)
1866}
1867
1868fn edge_endpoints_exist(
1876 conn: &rusqlite::Connection,
1877 source_id: Uuid,
1878 target_id: Uuid,
1879) -> Result<MissingEndpoints, rusqlite::Error> {
1880 let src_exists = endpoint_exists_clause("?1");
1881 let tgt_exists = endpoint_exists_clause("?2");
1882 let sql = format!("SELECT ({src_exists}), ({tgt_exists})");
1883 conn.query_row(
1884 &sql,
1885 rusqlite::params![source_id.to_string(), target_id.to_string()],
1886 |row| {
1887 let src_exists: bool = row.get(0)?;
1888 let tgt_exists: bool = row.get(1)?;
1889 Ok(MissingEndpoints {
1890 source: !src_exists,
1891 target: !tgt_exists,
1892 })
1893 },
1894 )
1895}
1896
1897fn edge_by_natural_key_including_deleted(
1898 conn: &rusqlite::Connection,
1899 edge: &Edge,
1900) -> Result<Option<Edge>, rusqlite::Error> {
1901 edge_by_natural_key_parts_including_deleted(
1902 conn,
1903 &edge.namespace,
1904 edge.source_id,
1905 edge.target_id,
1906 edge.relation,
1907 )
1908}
1909
1910fn edge_by_natural_key_parts_including_deleted(
1911 conn: &rusqlite::Connection,
1912 namespace: &str,
1913 source_id: Uuid,
1914 target_id: Uuid,
1915 relation: EdgeRelation,
1916) -> Result<Option<Edge>, rusqlite::Error> {
1917 let (source_id, target_id) = canonical_edge_endpoints(relation, source_id, target_id);
1918 conn.query_row(
1919 "SELECT namespace, id, source_id, target_id, relation, weight, \
1920 created_at, updated_at, deleted_at, metadata, target_backend \
1921 FROM graph_edges \
1922 WHERE namespace = ?1 AND source_id = ?2 AND target_id = ?3 AND relation = ?4",
1923 rusqlite::params![
1924 namespace,
1925 source_id.to_string(),
1926 target_id.to_string(),
1927 relation.as_str(),
1928 ],
1929 read_edge,
1930 )
1931 .optional()
1932}
1933
1934fn conditional_commit_annotation_insert(
1939 conn: &rusqlite::Connection,
1940 edge: Edge,
1941 guard: &CommitAnnotationGuard,
1942) -> Result<CommitAnnotationInsertOutcome, rusqlite::Error> {
1943 let cursor_matches: bool = conn.query_row(
1947 "SELECT (SELECT EXISTS(SELECT 1 FROM git_mirror_cursor \
1948 WHERE project_id=?1 AND kind='commits' \
1949 AND typeof(cursor_value)='text' \
1950 AND CAST(cursor_value AS BLOB)=?2 AND updated_at=?3)) \
1951 AND (SELECT EXISTS(SELECT 1 FROM git_mirror_cursor \
1952 WHERE project_id=?1 AND kind='commits_checkpoint' \
1953 AND typeof(cursor_value)='text' \
1954 AND CAST(cursor_value AS BLOB)=?4 AND updated_at=?5))",
1955 rusqlite::params![
1956 edge.target_id.to_string(),
1957 guard.commits.value,
1958 guard.commits.updated_at,
1959 guard.checkpoint.value,
1960 guard.checkpoint.updated_at,
1961 ],
1962 |row| row.get(0),
1963 )?;
1964 let metadata = edge
1965 .metadata
1966 .as_ref()
1967 .map(serde_json::to_string)
1968 .transpose()
1969 .map_err(|error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)))?;
1970 let affected = conn.execute(
1971 include_str!("../../sql/commit-annotation-insert.sql"),
1972 rusqlite::params![
1973 edge.namespace,
1974 Uuid::from(edge.id).to_string(),
1975 edge.source_id.to_string(),
1976 edge.target_id.to_string(),
1977 edge.weight,
1978 edge.created_at.timestamp_micros(),
1979 edge.updated_at.timestamp_micros(),
1980 metadata,
1981 guard.expected_sha,
1982 guard.source_identity,
1983 cursor_matches,
1984 ],
1985 )?;
1986 if affected == 1 {
1987 return Ok(CommitAnnotationInsertOutcome::Created(edge));
1988 }
1989 let source_live: bool = conn.query_row(
1990 include_str!("../../sql/commit-annotation-source-live-select.sql"),
1991 rusqlite::params![
1992 edge.source_id.to_string(),
1993 edge.namespace,
1994 guard.expected_sha
1995 ],
1996 |row| row.get(0),
1997 )?;
1998 if !source_live {
1999 return Ok(CommitAnnotationInsertOutcome::SourceChanged);
2000 }
2001 let target_live: bool = conn.query_row(
2002 "SELECT EXISTS(SELECT 1 FROM entities WHERE id=?1 AND namespace=?2 \
2003 AND kind='project' AND deleted_at IS NULL \
2004 AND json_type(properties, '$.repo_slug')='text' \
2005 AND json_extract(properties, '$.repo_slug')=?3 COLLATE BINARY)",
2006 rusqlite::params![
2007 edge.target_id.to_string(),
2008 edge.namespace,
2009 guard.source_identity
2010 ],
2011 |row| row.get(0),
2012 )?;
2013 if !target_live {
2014 return Ok(CommitAnnotationInsertOutcome::TargetChanged);
2015 }
2016 if !cursor_matches {
2017 return Ok(CommitAnnotationInsertOutcome::CursorChanged);
2018 }
2019 match edge_by_natural_key_including_deleted(conn, &edge)? {
2020 Some(existing) if existing.deleted_at.is_some() => {
2021 Ok(CommitAnnotationInsertOutcome::Tombstoned)
2022 }
2023 Some(_) => Ok(CommitAnnotationInsertOutcome::ExistingLive),
2024 None => Err(rusqlite::Error::QueryReturnedNoRows),
2025 }
2026}
2027
2028fn observed_edge_upsert(
2031 conn: &rusqlite::Connection,
2032 request: &EdgeUpsertRequest,
2033 guard_endpoints: bool,
2034) -> Result<GuardedEdgeUpsertOutcome, rusqlite::Error> {
2035 let (source_id, target_id) = canonical_edge_endpoints(
2036 request.edge.relation,
2037 request.edge.source_id,
2038 request.edge.target_id,
2039 );
2040 if guard_endpoints {
2041 #[cfg(test)]
2044 tests::insert_probe_seam::hook((source_id, target_id));
2045 let missing = edge_endpoints_exist(conn, source_id, target_id)?;
2046 if missing.any() {
2047 return Ok(GuardedEdgeUpsertOutcome::Refused(
2048 EdgeUpsertRefusal::MissingEndpoints(missing),
2049 ));
2050 }
2051 }
2052
2053 let previous = edge_by_natural_key_including_deleted(conn, &request.edge)?;
2054 if let Some(edge) = previous.as_ref() {
2055 if edge.deleted_at.is_some() && !request.resurrect {
2056 return Ok(GuardedEdgeUpsertOutcome::Refused(
2057 EdgeUpsertRefusal::ResurrectionRequired { edge: edge.clone() },
2058 ));
2059 }
2060 }
2061
2062 let statement = edge_upsert_statement_with_resurrection(&request.edge, request.resurrect);
2063 let mut stmt = conn.prepare(&statement.sql)?;
2064 bind_params(&mut stmt, &statement.params)?;
2065 let affected = stmt.raw_execute()?;
2066 if affected == 0 {
2067 let edge = edge_by_natural_key_including_deleted(conn, &request.edge)?
2068 .ok_or(rusqlite::Error::QueryReturnedNoRows)?;
2069 return Ok(GuardedEdgeUpsertOutcome::Refused(
2070 EdgeUpsertRefusal::ResurrectionRequired { edge },
2071 ));
2072 }
2073
2074 let edge = edge_by_natural_key_including_deleted(conn, &request.edge)?
2075 .ok_or(rusqlite::Error::QueryReturnedNoRows)?;
2076 let disposition = match previous.as_ref().and_then(|edge| edge.deleted_at) {
2077 None if previous.is_none() => EdgeUpsertDisposition::Created,
2078 None => EdgeUpsertDisposition::Updated,
2079 Some(_) => EdgeUpsertDisposition::Resurrected,
2080 };
2081 Ok(GuardedEdgeUpsertOutcome::Written(EdgeUpsertResult {
2082 edge,
2083 disposition,
2084 previous,
2085 }))
2086}
2087
2088fn observed_edge_batch_upsert(
2091 conn: &rusqlite::Connection,
2092 requests: &[EdgeUpsertRequest],
2093 guard_endpoints: bool,
2094) -> Result<GuardedEdgeBatchUpsertOutcome, rusqlite::Error> {
2095 for (entry_index, request) in requests.iter().enumerate() {
2099 let (source_id, target_id) = canonical_edge_endpoints(
2100 request.edge.relation,
2101 request.edge.source_id,
2102 request.edge.target_id,
2103 );
2104 if guard_endpoints {
2105 let missing = edge_endpoints_exist(conn, source_id, target_id)?;
2106 if missing.any() {
2107 return Ok(GuardedEdgeBatchUpsertOutcome {
2108 rows: Vec::new(),
2109 refusal: Some(GuardedEdgeBatchRefusal {
2110 entry_index,
2111 reason: EdgeUpsertRefusal::MissingEndpoints(missing),
2112 }),
2113 });
2114 }
2115 }
2116 if let Some(edge) = edge_by_natural_key_including_deleted(conn, &request.edge)? {
2117 if edge.deleted_at.is_some() && !request.resurrect {
2118 return Ok(GuardedEdgeBatchUpsertOutcome {
2119 rows: Vec::new(),
2120 refusal: Some(GuardedEdgeBatchRefusal {
2121 entry_index,
2122 reason: EdgeUpsertRefusal::ResurrectionRequired { edge },
2123 }),
2124 });
2125 }
2126 }
2127 }
2128
2129 let mut rows = Vec::with_capacity(requests.len());
2130 for request in requests {
2131 match observed_edge_upsert(conn, request, false)? {
2132 GuardedEdgeUpsertOutcome::Written(row) => rows.push(row),
2133 GuardedEdgeUpsertOutcome::Refused(_) => {
2134 return Err(rusqlite::Error::ExecuteReturnedResults)
2135 }
2136 }
2137 }
2138 Ok(GuardedEdgeBatchUpsertOutcome {
2139 rows,
2140 refusal: None,
2141 })
2142}
2143
2144fn traversal_neighbor_sql(
2149 direction: Direction,
2150 relation_count: usize,
2151 has_min_weight: bool,
2152) -> String {
2153 let (node_column, endpoint_column, index) = match direction {
2154 Direction::Out => ("target_id", "source_id", "idx_graph_edges_ns_src_rel"),
2155 Direction::In => ("source_id", "target_id", "idx_graph_edges_ns_tgt_rel"),
2156 Direction::Both => unreachable!("Direction::Both is split into indexed Out/In seeks"),
2157 };
2158 let mut sql = format!(
2159 "SELECT {node_column}, id, weight \
2160 FROM graph_edges INDEXED BY {index} \
2161 WHERE namespace = ?1 AND {endpoint_column} = ?2 AND deleted_at IS NULL"
2162 );
2163 if relation_count > 0 {
2164 let placeholders = (0..relation_count)
2165 .map(|offset| format!("?{}", 4 + offset))
2166 .collect::<Vec<_>>()
2167 .join(",");
2168 sql.push_str(&format!(" AND relation IN ({placeholders})"));
2169 }
2170 if has_min_weight {
2171 sql.push_str(&format!(" AND weight >= ?{}", 4 + relation_count));
2172 }
2173 sql.push_str(" LIMIT ?3");
2174 sql
2175}
2176
2177fn traversal_timeout_error(budget: &TraversalExecutionBudget) -> StorageError {
2178 StorageError::Timeout {
2179 operation: format!(
2180 "traverse ({}ms execution budget)",
2181 budget.max_duration().as_millis()
2182 )
2183 .into(),
2184 }
2185}
2186
2187fn traversal_work_error(budget: &TraversalExecutionBudget) -> StorageError {
2188 StorageError::InvalidInput {
2189 capability: StorageCapability::Graph,
2190 operation: "traverse".into(),
2191 message: format!(
2192 "traversal work budget exceeded after {} adjacency rows; \
2193 narrow roots, depth, relations, or result limit",
2194 budget.work_limit()
2195 ),
2196 }
2197}
2198
2199#[derive(Clone, Copy)]
2200struct TraversalFrontierNode {
2201 node_id: Uuid,
2202 depth: usize,
2203 total_weight: f64,
2204}
2205
2206#[allow(clippy::too_many_arguments)]
2207fn run_bounded_traversal(
2208 conn: &rusqlite::Connection,
2209 roots: Vec<Uuid>,
2210 opts: TraversalOptions,
2211 include_roots: bool,
2212 namespace: String,
2213 origin: khive_storage::tx_registry::TxOrigin,
2214 budget: TraversalExecutionBudget,
2215 counted_rows: &std::sync::atomic::AtomicU64,
2216 counted_queries: &std::sync::atomic::AtomicU64,
2217) -> Result<Vec<GraphPath>, StorageError> {
2218 let progress_timed_out = Arc::new(std::sync::atomic::AtomicBool::new(false));
2219 let callback_timed_out = Arc::clone(&progress_timed_out);
2220 let callback_budget = budget.clone();
2221 #[cfg(test)]
2222 let progress_seam_root = roots.first().copied();
2223 conn.progress_handler(
2224 1_000,
2225 Some(move || {
2226 if crate::read_cancellation::current_read_should_interrupt() {
2227 return true;
2228 }
2229 #[cfg(test)]
2230 if tests::traverse_progress_seam::hook(progress_seam_root) {
2231 callback_timed_out.store(true, std::sync::atomic::Ordering::Relaxed);
2232 return true;
2233 }
2234 let expired = callback_budget.is_expired();
2235 if expired {
2236 callback_timed_out.store(true, std::sync::atomic::Ordering::Relaxed);
2237 }
2238 expired
2239 }),
2240 )
2241 .map_err(|e| map_err(e, "traverse_progress_handler"))?;
2242
2243 let result = (|| {
2244 let result_limit = opts.effective_limit() as usize;
2245 let relation_count = opts.relations.as_ref().map_or(0, Vec::len);
2246 let directions = match opts.direction {
2247 Direction::Out => vec![Direction::Out],
2248 Direction::In => vec![Direction::In],
2249 Direction::Both => vec![Direction::Out, Direction::In],
2250 };
2251 let statements = directions
2252 .into_iter()
2253 .map(|direction| {
2254 traversal_neighbor_sql(direction, relation_count, opts.min_weight.is_some())
2255 })
2256 .collect::<Vec<_>>();
2257 let map_sql_error = |error| {
2258 if progress_timed_out.load(std::sync::atomic::Ordering::Relaxed) {
2259 traversal_timeout_error(&budget)
2260 } else {
2261 map_err(error, "traverse")
2262 }
2263 };
2264
2265 let mut all_paths = Vec::with_capacity(roots.len());
2266 for root_id in roots {
2267 let mut seen = HashSet::new();
2268 seen.insert(root_id);
2269 let mut frontier = VecDeque::from([TraversalFrontierNode {
2270 node_id: root_id,
2271 depth: 0,
2272 total_weight: 0.0,
2273 }]);
2274 let mut nodes = Vec::with_capacity(result_limit + usize::from(include_roots));
2275 if include_roots {
2276 nodes.push(PathNode {
2277 node_id: root_id,
2278 via_edge: None,
2279 depth: 0,
2280 name: None,
2281 kind: None,
2282 properties: None,
2283 weight: 0.0,
2284 });
2285 }
2286 let mut non_root_count = 0usize;
2287
2288 'root_walk: while non_root_count < result_limit {
2289 let Some(current) = frontier.pop_front() else {
2290 break;
2291 };
2292 if current.depth >= opts.max_depth {
2293 continue;
2294 }
2295 if budget.is_expired() {
2296 return Err(traversal_timeout_error(&budget));
2297 }
2298
2299 for sql in &statements {
2300 let row_cap = budget.remaining_work().saturating_add(1);
2301 let mut params: Vec<Box<dyn rusqlite::types::ToSql>> = vec![
2302 Box::new(namespace.clone()),
2303 Box::new(current.node_id.to_string()),
2304 Box::new(row_cap as i64),
2305 ];
2306 if let Some(relations) = &opts.relations {
2307 params.extend(relations.iter().map(|relation| {
2308 Box::new(relation.to_string()) as Box<dyn rusqlite::types::ToSql>
2309 }));
2310 }
2311 if let Some(min_weight) = opts.min_weight {
2312 params.push(Box::new(min_weight));
2313 }
2314 let param_refs = params
2315 .iter()
2316 .map(|param| param.as_ref())
2317 .collect::<Vec<&dyn rusqlite::types::ToSql>>();
2318
2319 let _snapshot = khive_storage::tx_registry::register_scoped(
2320 Some("graph_traverse_read".to_string()),
2321 origin.clone(),
2322 );
2323 let mut stmt = conn.prepare(sql).map_err(&map_sql_error)?;
2324 counted_queries.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2325 let mut rows = stmt.query(param_refs.as_slice()).map_err(&map_sql_error)?;
2326 while let Some(row) = rows.next().map_err(&map_sql_error)? {
2327 #[cfg(test)]
2330 tests::traverse_snapshot_seam::hook(current.node_id);
2331 counted_rows.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2332 if budget.is_expired() {
2333 return Err(traversal_timeout_error(&budget));
2334 }
2335 if !budget.try_consume_row() {
2336 return Err(traversal_work_error(&budget));
2337 }
2338 let node_str: String = row.get(0).map_err(&map_sql_error)?;
2339 let edge_str: String = row.get(1).map_err(&map_sql_error)?;
2340 let edge_weight: f64 = row.get(2).map_err(&map_sql_error)?;
2341 let node_id = parse_uuid(&node_str).map_err(&map_sql_error)?;
2342 if !seen.insert(node_id) {
2343 continue;
2344 }
2345 let via_edge = parse_uuid(&edge_str).map_err(&map_sql_error)?;
2346 let depth = current.depth + 1;
2347 let total_weight = current.total_weight + edge_weight;
2348 nodes.push(PathNode {
2349 node_id,
2350 via_edge: Some(via_edge),
2351 depth,
2352 name: None,
2353 kind: None,
2354 properties: None,
2355 weight: total_weight,
2356 });
2357 non_root_count += 1;
2358 if depth < opts.max_depth {
2359 frontier.push_back(TraversalFrontierNode {
2360 node_id,
2361 depth,
2362 total_weight,
2363 });
2364 }
2365 if non_root_count == result_limit {
2366 break 'root_walk;
2367 }
2368 }
2369 }
2370 }
2371
2372 if !nodes.is_empty() {
2373 let total_weight = nodes.iter().map(|node| node.weight).fold(0.0_f64, f64::max);
2374 all_paths.push(GraphPath {
2375 root_id,
2376 nodes,
2377 total_weight,
2378 });
2379 }
2380 }
2381 Ok(all_paths)
2382 })();
2383
2384 conn.progress_handler(0, None::<fn() -> bool>)
2385 .map_err(|e| map_err(e, "traverse_progress_handler_clear"))?;
2386 result
2387}
2388
2389impl SqlGraphStore {
2390 async fn query_neighbors_page(
2391 &self,
2392 operation: &'static str,
2393 node_id: Uuid,
2394 query: NeighborQuery,
2395 after: Option<NeighborCursor>,
2396 neighbor_kinds: Option<Vec<String>>,
2397 ) -> Result<Vec<NeighborHit>, StorageError> {
2398 count_neighbor_select();
2399
2400 let namespace = self.namespace.clone();
2401 let node_str = node_id.to_string();
2402 let counted_queries = Arc::new(std::sync::atomic::AtomicU64::new(0));
2403 let counted_rows = Arc::new(std::sync::atomic::AtomicU64::new(0));
2404 let closure_queries = Arc::clone(&counted_queries);
2405 let closure_rows = Arc::clone(&counted_rows);
2406 let result = self
2407 .with_reader(operation, move |conn| {
2408 let base_out = "SELECT target_id AS node_id, id AS edge_id, relation, weight \
2409 FROM graph_edges \
2410 WHERE namespace = ?1 AND source_id = ?2 AND deleted_at IS NULL";
2411 let base_in = "SELECT source_id AS node_id, id AS edge_id, relation, weight \
2412 FROM graph_edges \
2413 WHERE namespace = ?1 AND target_id = ?2 AND deleted_at IS NULL";
2414 let sql = match query.direction {
2415 Direction::Out => base_out.to_string(),
2416 Direction::In => base_in.to_string(),
2417 Direction::Both => format!("{} UNION ALL {}", base_out, base_in),
2418 };
2419 let (where_extra, limit_clause, extra_params) =
2420 neighbor_extra_clause(&query, 3, after.as_ref(), neighbor_kinds.as_deref());
2421 let full_sql = format!(
2422 "SELECT node_id, edge_id, relation, weight FROM ({}){} \
2423 ORDER BY weight DESC, node_id ASC, edge_id ASC{}",
2424 sql, where_extra, limit_clause
2425 );
2426 let mut stmt = conn.prepare(&full_sql)?;
2427 let mut all_params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
2428 all_params.push(Box::new(namespace.clone()));
2429 all_params.push(Box::new(node_str.clone()));
2430 all_params.extend(extra_params);
2431 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
2432 all_params.iter().map(|p| p.as_ref()).collect();
2433
2434 closure_queries.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2435 let rows = stmt.query_map(param_refs.as_slice(), |row| {
2436 let nid_str: String = row.get(0)?;
2437 let eid_str: String = row.get(1)?;
2438 let relation_str: String = row.get(2)?;
2439 let weight: f64 = row.get(3)?;
2440 Ok((nid_str, eid_str, relation_str, weight))
2441 })?;
2442 let mut hits = Vec::new();
2443 for row in rows {
2444 let (nid_str, eid_str, relation_str, weight) = row?;
2445 closure_rows.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
2446 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
2447 rusqlite::Error::FromSqlConversionFailure(
2448 2,
2449 rusqlite::types::Type::Text,
2450 Box::new(e),
2451 )
2452 })?;
2453 hits.push(NeighborHit {
2454 node_id: parse_uuid(&nid_str)?,
2455 edge_id: parse_uuid(&eid_str)?,
2456 relation,
2457 weight,
2458 name: None,
2459 kind: None,
2460 entity_type: None,
2461 });
2462 }
2463 Ok(hits)
2464 })
2465 .await;
2466
2467 report_graph_usage(&counted_queries, &counted_rows);
2468 result
2469 }
2470}
2471
2472#[async_trait]
2473impl GraphStore for SqlGraphStore {
2474 async fn latest_annotating_note(
2475 &self,
2476 node_id: Uuid,
2477 kind: &str,
2478 tag: &str,
2479 ) -> Result<Option<(Uuid, i64)>, StorageError> {
2480 let namespace = self.namespace.clone();
2481 let node_id = node_id.to_string();
2482 let kind = kind.to_owned();
2483 let tag = tag.to_owned();
2484 self.with_indexed_reader("latest_annotating_note", move |conn| {
2485 conn.query_row(
2486 LATEST_ANNOTATING_NOTE_SQL,
2487 rusqlite::params![namespace, node_id, kind, tag],
2488 |row| {
2489 let id: String = row.get(0)?;
2490 Ok((parse_uuid(&id)?, row.get(1)?))
2491 },
2492 )
2493 .optional()
2494 })
2495 .await
2496 }
2497
2498 async fn latest_annotating_note_with_property(
2499 &self,
2500 node_id: Uuid,
2501 kind: &str,
2502 tag: &str,
2503 property_key: &str,
2504 property_value: &str,
2505 ) -> Result<Option<(Uuid, i64)>, StorageError> {
2506 let namespace = self.namespace.clone();
2507 let node_id = node_id.to_string();
2508 let kind = kind.to_owned();
2509 let tag = tag.to_owned();
2510 let property_key = property_key.to_owned();
2511 let property_value = property_value.to_owned();
2512 self.with_indexed_reader("latest_annotating_note_with_property", move |conn| {
2513 conn.query_row(
2514 LATEST_ANNOTATING_NOTE_WITH_PROPERTY_SQL,
2515 rusqlite::params![namespace, node_id, kind, tag, property_key, property_value],
2516 |row| {
2517 let id: String = row.get(0)?;
2518 Ok((parse_uuid(&id)?, row.get(1)?))
2519 },
2520 )
2521 .optional()
2522 })
2523 .await
2524 }
2525
2526 async fn upsert_edge(&self, edge: Edge) -> Result<(), StorageError> {
2527 self.upsert_edge_observed(EdgeUpsertRequest {
2528 edge,
2529 resurrect: false,
2530 })
2531 .await
2532 .map(|_| ())
2533 }
2534
2535 async fn upsert_edge_observed(
2536 &self,
2537 request: EdgeUpsertRequest,
2538 ) -> Result<EdgeUpsertResult, StorageError> {
2539 match self
2540 .observed_edge_write("upsert_edge_observed", request, false)
2541 .await?
2542 {
2543 GuardedEdgeUpsertOutcome::Written(result) => Ok(result),
2544 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::ResurrectionRequired { edge }) => {
2545 Err(resurrection_required_error("upsert_edge_observed", &edge))
2546 }
2547 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::MissingEndpoints(_)) => {
2548 Err(StorageError::Conflict {
2549 capability: StorageCapability::Graph,
2550 operation: "upsert_edge_observed".into(),
2551 message: "unguarded edge upsert reported a missing-endpoint refusal".into(),
2552 })
2553 }
2554 }
2555 }
2556
2557 async fn insert_edge_if_absent(&self, edge: Edge) -> Result<bool, StorageError> {
2558 let statement = edge_insert_if_absent_statement(&edge);
2559 self.with_writer("insert_edge_if_absent", move |conn| {
2560 let mut stmt = conn.prepare(&statement.sql)?;
2561 bind_params(&mut stmt, &statement.params)?;
2562 Ok(stmt.raw_execute()? > 0)
2563 })
2564 .await
2565 }
2566
2567 async fn insert_commit_annotation_if_absent(
2568 &self,
2569 edge: Edge,
2570 guard: CommitAnnotationGuard,
2571 ) -> Result<CommitAnnotationInsertOutcome, StorageError> {
2572 const OP: &str = "insert_commit_annotation_if_absent";
2573 if edge.relation != EdgeRelation::Annotates || edge.deleted_at.is_some() {
2574 return Err(StorageError::InvalidInput {
2575 capability: StorageCapability::Graph,
2576 operation: OP.into(),
2577 message: "expected a live annotates edge".into(),
2578 });
2579 }
2580 if let Some(writer_task) = self.current_writer_task(OP)? {
2581 return writer_task
2582 .send_bounded(move |conn| {
2583 conditional_commit_annotation_insert(conn, edge, &guard)
2584 .map_err(|error| map_err(error, OP))
2585 })
2586 .await;
2587 }
2588 let origin = self.pool.origin();
2589 self.with_writer(OP, move |conn| {
2590 conn.execute_batch("BEGIN IMMEDIATE")?;
2591 let _tx_handle =
2592 khive_storage::tx_registry::register_scoped(Some(OP.to_string()), origin);
2593 let outcome = match conditional_commit_annotation_insert(conn, edge, &guard) {
2594 Ok(outcome) => outcome,
2595 Err(error) => {
2596 let _ = conn.execute_batch("ROLLBACK");
2597 return Err(error);
2598 }
2599 };
2600 if let Err(error) = conn.execute_batch("COMMIT") {
2601 let _ = conn.execute_batch("ROLLBACK");
2602 return Err(error);
2603 }
2604 Ok(outcome)
2605 })
2606 .await
2607 }
2608
2609 async fn replace_edge_if_unchanged(
2610 &self,
2611 edge: Edge,
2612 expected_updated_at: DateTime<Utc>,
2613 expected_deleted_at: Option<DateTime<Utc>>,
2614 ) -> Result<bool, StorageError> {
2615 let statement =
2616 edge_replace_if_unchanged_statement(&edge, expected_updated_at, expected_deleted_at);
2617 self.with_writer("replace_edge_if_unchanged", move |conn| {
2618 let mut stmt = conn.prepare(&statement.sql)?;
2619 bind_params(&mut stmt, &statement.params)?;
2620 Ok(stmt.raw_execute()? > 0)
2621 })
2622 .await
2623 }
2624
2625 async fn upsert_edges(&self, edges: Vec<Edge>) -> Result<BatchWriteSummary, StorageError> {
2626 let attempted = edges.len() as u64;
2627 let requests = edges
2628 .into_iter()
2629 .map(|edge| EdgeUpsertRequest {
2630 edge,
2631 resurrect: false,
2632 })
2633 .collect();
2634 let outcome = self
2635 .observed_edge_batch_write("upsert_edges", requests, false)
2636 .await?;
2637 if let Some(refusal) = outcome.refusal {
2638 return match refusal.reason {
2639 EdgeUpsertRefusal::ResurrectionRequired { edge } => {
2640 Err(resurrection_required_error("upsert_edges", &edge))
2641 }
2642 EdgeUpsertRefusal::MissingEndpoints(_) => Err(StorageError::Conflict {
2643 capability: StorageCapability::Graph,
2644 operation: "upsert_edges".into(),
2645 message: "unguarded edge batch reported a missing-endpoint refusal".into(),
2646 }),
2647 };
2648 }
2649 Ok(BatchWriteSummary {
2650 attempted,
2651 affected: outcome.rows.len() as u64,
2652 ..BatchWriteSummary::default()
2653 })
2654 }
2655
2656 async fn upsert_edge_guarded(&self, edge: Edge) -> Result<GuardedWriteOutcome, StorageError> {
2657 match self
2658 .upsert_edge_guarded_observed(EdgeUpsertRequest {
2659 edge,
2660 resurrect: false,
2661 })
2662 .await?
2663 {
2664 GuardedEdgeUpsertOutcome::Written(_) => Ok(GuardedWriteOutcome::Written),
2665 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::MissingEndpoints(missing)) => {
2666 Ok(GuardedWriteOutcome::Refused(missing))
2667 }
2668 GuardedEdgeUpsertOutcome::Refused(EdgeUpsertRefusal::ResurrectionRequired { edge }) => {
2669 Err(resurrection_required_error("upsert_edge_guarded", &edge))
2670 }
2671 }
2672 }
2673
2674 async fn upsert_edge_guarded_observed(
2675 &self,
2676 request: EdgeUpsertRequest,
2677 ) -> Result<GuardedEdgeUpsertOutcome, StorageError> {
2678 self.observed_edge_write("upsert_edge_guarded_observed", request, true)
2679 .await
2680 }
2681
2682 async fn upsert_edges_guarded(
2683 &self,
2684 edges: Vec<Edge>,
2685 ) -> Result<GuardedBatchOutcome, StorageError> {
2686 let attempted = edges.len() as u64;
2687 let requests = edges
2688 .iter()
2689 .cloned()
2690 .map(|edge| EdgeUpsertRequest {
2691 edge,
2692 resurrect: false,
2693 })
2694 .collect();
2695 let outcome = self.upsert_edges_guarded_observed(requests).await?;
2696 match outcome.refusal {
2697 None => Ok(GuardedBatchOutcome {
2698 summary: BatchWriteSummary {
2699 attempted,
2700 affected: outcome.rows.len() as u64,
2701 ..BatchWriteSummary::default()
2702 },
2703 refused: None,
2704 }),
2705 Some(refusal) => match refusal.reason {
2706 EdgeUpsertRefusal::MissingEndpoints(missing) => {
2707 let index = refusal.entry_index;
2708 let edge = &edges[index];
2709 let (source_id, target_id) =
2710 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
2711 let message = format!(
2712 "batch entry {index}: edge endpoint no longer exists at write time: source \
2713 {source_id} or target {target_id}"
2714 );
2715 let mut summary = BatchWriteSummary {
2716 attempted,
2717 ..BatchWriteSummary::default()
2718 };
2719 summary.first_error = message.clone();
2722 let refusal = GuardedBatchRefusal {
2723 entry_index: index,
2724 missing,
2725 };
2726 for (failed_index, failed_edge) in edges.iter().enumerate() {
2727 refusal.record_failure(&mut summary, failed_index, failed_edge, &message);
2728 }
2729 Ok(GuardedBatchOutcome {
2730 summary,
2731 refused: Some(refusal),
2732 })
2733 }
2734 EdgeUpsertRefusal::ResurrectionRequired { edge } => {
2735 Err(resurrection_required_error("upsert_edges_guarded", &edge))
2736 }
2737 },
2738 }
2739 }
2740
2741 async fn upsert_edges_guarded_observed(
2742 &self,
2743 requests: Vec<EdgeUpsertRequest>,
2744 ) -> Result<GuardedEdgeBatchUpsertOutcome, StorageError> {
2745 self.observed_edge_batch_write("upsert_edges_guarded_observed", requests, true)
2746 .await
2747 }
2748
2749 async fn get_edge(&self, id: LinkId) -> Result<Option<Edge>, StorageError> {
2750 let id_str = Uuid::from(id).to_string();
2751
2752 self.with_reader("get_edge", move |conn| {
2753 let mut stmt = conn.prepare(
2754 "SELECT namespace, id, source_id, target_id, relation, weight, \
2755 created_at, updated_at, deleted_at, metadata, target_backend \
2756 FROM graph_edges WHERE id = ?1 AND deleted_at IS NULL",
2757 )?;
2758 let mut rows = stmt.query(rusqlite::params![id_str])?;
2759 match rows.next()? {
2760 Some(row) => Ok(Some(read_edge(row)?)),
2761 None => Ok(None),
2762 }
2763 })
2764 .await
2765 }
2766
2767 async fn get_edge_including_deleted(&self, id: LinkId) -> Result<Option<Edge>, StorageError> {
2768 let id_str = Uuid::from(id).to_string();
2769
2770 self.with_reader("get_edge_including_deleted", move |conn| {
2771 let mut stmt = conn.prepare(
2772 "SELECT namespace, id, source_id, target_id, relation, weight, \
2773 created_at, updated_at, deleted_at, metadata, target_backend \
2774 FROM graph_edges WHERE id = ?1",
2775 )?;
2776 let mut rows = stmt.query(rusqlite::params![id_str])?;
2777 match rows.next()? {
2778 Some(row) => Ok(Some(read_edge(row)?)),
2779 None => Ok(None),
2780 }
2781 })
2782 .await
2783 }
2784
2785 async fn edge_sequence(&self, id: Uuid) -> Result<Option<i64>, StorageError> {
2786 let id = id.to_string();
2787 self.with_reader("edge_sequence", move |conn| {
2788 conn.query_row(
2789 "SELECT seq FROM graph_edges_seq WHERE edge_id = ?1",
2790 rusqlite::params![id],
2791 |row| row.get(0),
2792 )
2793 .optional()
2794 })
2795 .await
2796 }
2797
2798 async fn edge_sequences(&self, ids: &[Uuid]) -> Result<Vec<(Uuid, i64)>, StorageError> {
2799 if ids.is_empty() {
2800 return Ok(Vec::new());
2801 }
2802 let ids = ids.to_vec();
2803 self.with_reader("edge_sequences", move |conn| {
2804 const CHUNK: usize = 900;
2805 let mut resolved = Vec::with_capacity(ids.len());
2806 for chunk in ids.chunks(CHUNK) {
2807 let placeholders = (1..=chunk.len())
2808 .map(|index| format!("?{index}"))
2809 .collect::<Vec<_>>()
2810 .join(", ");
2811 let sql = format!(
2812 "SELECT edge_id, seq FROM graph_edges_seq WHERE edge_id IN ({placeholders})"
2813 );
2814 let strings = chunk.iter().map(Uuid::to_string).collect::<Vec<_>>();
2815 let params = strings
2816 .iter()
2817 .map(|id| id as &dyn rusqlite::types::ToSql)
2818 .collect::<Vec<_>>();
2819 let mut stmt = conn.prepare(&sql)?;
2820 let rows = stmt.query_map(params.as_slice(), |row| {
2821 let id: String = row.get(0)?;
2822 Ok((parse_uuid(&id)?, row.get::<_, i64>(1)?))
2823 })?;
2824 resolved.extend(rows.collect::<Result<Vec<_>, _>>()?);
2825 }
2826 Ok(resolved)
2827 })
2828 .await
2829 }
2830
2831 async fn get_edge_by_natural_key_including_deleted(
2832 &self,
2833 namespace: &str,
2834 source_id: Uuid,
2835 target_id: Uuid,
2836 relation: EdgeRelation,
2837 ) -> Result<Option<Edge>, StorageError> {
2838 let namespace = namespace.to_string();
2839 let source_str = source_id.to_string();
2840 let target_str = target_id.to_string();
2841 let relation_str = relation.to_string();
2842
2843 self.with_reader("get_edge_by_natural_key_including_deleted", move |conn| {
2844 let mut stmt = conn.prepare(
2845 "SELECT namespace, id, source_id, target_id, relation, weight, \
2846 created_at, updated_at, deleted_at, metadata, target_backend \
2847 FROM graph_edges \
2848 WHERE namespace = ?1 AND source_id = ?2 AND target_id = ?3 AND relation = ?4",
2849 )?;
2850 let mut rows = stmt.query(rusqlite::params![
2851 namespace,
2852 source_str,
2853 target_str,
2854 relation_str
2855 ])?;
2856 match rows.next()? {
2857 Some(row) => Ok(Some(read_edge(row)?)),
2858 None => Ok(None),
2859 }
2860 })
2861 .await
2862 }
2863
2864 async fn get_edges(&self, ids: &[LinkId]) -> Result<Vec<Edge>, StorageError> {
2865 if ids.is_empty() {
2866 return Ok(Vec::new());
2867 }
2868 const CHUNK: usize = 900;
2870 let id_strs: Vec<String> = ids.iter().map(|id| Uuid::from(*id).to_string()).collect();
2871
2872 let mut result: Vec<Edge> = Vec::with_capacity(ids.len());
2873 for chunk in id_strs.chunks(CHUNK) {
2874 let chunk_owned: Vec<String> = chunk.to_vec();
2875 let edges = self
2876 .with_reader("get_edges", move |conn| {
2877 let placeholders: Vec<String> =
2878 (1..=chunk_owned.len()).map(|i| format!("?{}", i)).collect();
2879 let sql = format!(
2880 "SELECT namespace, id, source_id, target_id, relation, weight, \
2881 created_at, updated_at, deleted_at, metadata, target_backend \
2882 FROM graph_edges WHERE id IN ({}) AND deleted_at IS NULL",
2883 placeholders.join(",")
2884 );
2885 let mut stmt = conn.prepare(&sql)?;
2886 let params: Vec<&dyn rusqlite::types::ToSql> = chunk_owned
2887 .iter()
2888 .map(|s| s as &dyn rusqlite::types::ToSql)
2889 .collect();
2890 let rows = stmt.query_map(params.as_slice(), read_edge)?;
2891 let mut edges = Vec::new();
2892 for row in rows {
2893 edges.push(row?);
2894 }
2895 Ok(edges)
2896 })
2897 .await?;
2898 result.extend(edges);
2899 }
2900 Ok(result)
2901 }
2902
2903 async fn get_edge_read_outcomes(
2904 &self,
2905 ids: &[LinkId],
2906 ) -> Result<Vec<Result<Option<Edge>, StorageError>>, StorageError> {
2907 const CHUNK: usize = 900;
2908 let mut outcomes = Vec::with_capacity(ids.len());
2909 for chunk in ids.chunks(CHUNK) {
2910 let requested: Vec<String> =
2911 chunk.iter().map(|id| Uuid::from(*id).to_string()).collect();
2912 let expected = requested.len();
2913 let rows = self
2914 .with_reader("get_edge_read_outcomes", move |conn| {
2915 let requested_json = serde_json::to_string(&requested)
2916 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
2917 let mut stmt = conn.prepare(
2918 "SELECT e.namespace, e.id, e.source_id, e.target_id, e.relation, e.weight, \
2919 e.created_at, e.updated_at, e.deleted_at, e.metadata, e.target_backend, \
2920 requested.key \
2921 FROM json_each(?1) AS requested \
2922 LEFT JOIN graph_edges AS e \
2923 ON e.id = requested.value AND e.deleted_at IS NULL \
2924 ORDER BY CAST(requested.key AS INTEGER) ASC",
2925 )?;
2926 let mut rows = stmt.query([requested_json])?;
2927 let mut outcomes = Vec::with_capacity(expected);
2928 while let Some(row) = rows.next()? {
2929 let ordinal: i64 = row.get(11)?;
2930 let ordinal =
2931 usize::try_from(ordinal).map_err(|_| rusqlite::Error::InvalidQuery)?;
2932 if ordinal != outcomes.len() {
2933 return Err(rusqlite::Error::InvalidQuery);
2934 }
2935 let outcome = match row.get_ref(1)? {
2936 rusqlite::types::ValueRef::Null => Ok(None),
2937 _ => read_edge(row).map(Some).map_err(|e| map_err(e, "get_edge")),
2938 };
2939 outcomes.push(outcome);
2940 }
2941 if outcomes.len() != expected {
2942 return Err(rusqlite::Error::InvalidQuery);
2943 }
2944 Ok(outcomes)
2945 })
2946 .await?;
2947 outcomes.extend(rows);
2948 }
2949 Ok(outcomes)
2950 }
2951
2952 async fn batch_neighbors(
2953 &self,
2954 sources: &[Uuid],
2955 query: NeighborQuery,
2956 ) -> Result<Vec<(Uuid, NeighborHit)>, StorageError> {
2957 use khive_storage::types::Direction;
2958
2959 if sources.is_empty() {
2960 return Ok(Vec::new());
2961 }
2962 let mut seen_sources = HashSet::with_capacity(sources.len());
2963 let unique_sources: Vec<Uuid> = sources
2964 .iter()
2965 .copied()
2966 .filter(|source| seen_sources.insert(*source))
2967 .collect();
2968 const CHUNK_SIZE: usize = 880;
2969
2970 let namespace = self.namespace.clone();
2971 let mut result: Vec<(Uuid, NeighborHit)> = Vec::new();
2972 let counted_queries = Arc::new(std::sync::atomic::AtomicU64::new(0));
2975 let counted_rows = Arc::new(std::sync::atomic::AtomicU64::new(0));
2976
2977 for chunk in unique_sources.chunks(CHUNK_SIZE) {
2978 let chunk_owned: Vec<Uuid> = chunk.to_vec();
2979 let query_clone = query.clone();
2980 let ns = namespace.clone();
2981 let closure_queries = Arc::clone(&counted_queries);
2982 let closure_rows = Arc::clone(&counted_rows);
2983
2984 let chunk_result = self
2985 .with_reader("batch_neighbors", move |conn| {
2986 let src_strs: Vec<String> = chunk_owned.iter().map(|u| u.to_string()).collect();
2987
2988 let sources_json = serde_json::to_string(&src_strs).map_err(|error| {
2989 rusqlite::Error::ToSqlConversionFailure(Box::new(error))
2990 })?;
2991
2992 let build_inner_sql =
2993 |direction_out: bool,
2994 q: &NeighborQuery|
2995 -> (String, Vec<String>, Option<f64>) {
2996 let (filter_col, node_col) = if direction_out {
2997 ("source_id", "target_id")
2998 } else {
2999 ("target_id", "source_id")
3000 };
3001
3002 let mut rel_params: Vec<String> = Vec::new();
3003 let mut conditions: Vec<String> = Vec::new();
3004 let mut param_idx = 3;
3005
3006 if let Some(ref rels) = q.relations {
3007 if !rels.is_empty() {
3008 let ps: Vec<String> = rels
3009 .iter()
3010 .map(|r| {
3011 rel_params.push(r.to_string());
3012 let p = format!("?{param_idx}");
3013 param_idx += 1;
3014 p
3015 })
3016 .collect();
3017 conditions
3018 .push(format!("edges.relation IN ({})", ps.join(",")));
3019 }
3020 }
3021
3022 let min_weight_val = if let Some(min_w) = q.min_weight {
3025 conditions.push(format!("edges.weight >= ?{param_idx}"));
3026 Some(min_w)
3027 } else {
3028 None
3029 };
3030
3031 let where_extra = if conditions.is_empty() {
3032 String::new()
3033 } else {
3034 format!(" AND {}", conditions.join(" AND "))
3035 };
3036
3037 let sql = format!(
3038 "SELECT requested.origin_id, edges.{node_col} AS node_id, \
3039 edges.id AS edge_id, edges.relation, edges.weight \
3040 FROM requested CROSS JOIN graph_edges AS edges \
3041 ON edges.{filter_col} = requested.origin_id \
3042 WHERE edges.namespace = ?1 \
3043 AND edges.deleted_at IS NULL{where_extra}",
3044 );
3045 (sql, rel_params, min_weight_val)
3046 };
3047
3048 let mut all_params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
3049 all_params.push(Box::new(ns.to_string()));
3050 all_params.push(Box::new(sources_json));
3051
3052 let (combined_inner, rel_params, min_weight_val) = match query_clone.direction {
3053 Direction::Out => build_inner_sql(true, &query_clone),
3054 Direction::In => build_inner_sql(false, &query_clone),
3055 Direction::Both => {
3056 let (out_sql, rel_params, min_weight_val) =
3057 build_inner_sql(true, &query_clone);
3058 let (in_sql, _, _) = build_inner_sql(false, &query_clone);
3059 (
3060 format!("{out_sql} UNION ALL {in_sql}"),
3061 rel_params,
3062 min_weight_val,
3063 )
3064 }
3065 };
3066
3067 for relation in rel_params {
3068 all_params.push(Box::new(relation));
3069 }
3070 if let Some(min_weight) = min_weight_val {
3071 all_params.push(Box::new(min_weight));
3072 }
3073 let limit_param_idx = all_params.len() + 1;
3074
3075 let full_sql = if let Some(lim) = query_clone.limit {
3083 all_params.push(Box::new(lim as i64));
3084 format!(
3085 "WITH requested(origin_id) AS (\
3086 SELECT value FROM json_each(?2)\
3087 ) SELECT origin_id, node_id, edge_id, relation, weight \
3088 FROM (SELECT *, ROW_NUMBER() OVER (PARTITION BY origin_id \
3089 ORDER BY weight DESC, node_id ASC) AS rn \
3090 FROM ({combined_inner})) WHERE rn <= ?{limit_param_idx}",
3091 )
3092 } else {
3093 format!(
3094 "WITH requested(origin_id) AS (\
3095 SELECT value FROM json_each(?2)\
3096 ) SELECT origin_id, node_id, edge_id, relation, weight \
3097 FROM ({combined_inner})",
3098 )
3099 };
3100
3101 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3102 all_params.iter().map(|p| p.as_ref()).collect();
3103
3104 let mut stmt = conn.prepare(&full_sql)?;
3105 closure_queries.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
3106 let rows = stmt.query_map(param_refs.as_slice(), |row| {
3107 let origin_str: String = row.get(0)?;
3108 let nid_str: String = row.get(1)?;
3109 let eid_str: String = row.get(2)?;
3110 let relation_str: String = row.get(3)?;
3111 let weight: f64 = row.get(4)?;
3112 Ok((origin_str, nid_str, eid_str, relation_str, weight))
3113 })?;
3114
3115 let mut pairs = Vec::new();
3116 for row in rows {
3117 let (origin_str, nid_str, eid_str, relation_str, weight) = row?;
3118 closure_rows.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
3119 let origin = parse_uuid(&origin_str)?;
3120 let node_id = parse_uuid(&nid_str)?;
3121 let edge_id = parse_uuid(&eid_str)?;
3122 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
3123 rusqlite::Error::FromSqlConversionFailure(
3124 3,
3125 rusqlite::types::Type::Text,
3126 Box::new(e),
3127 )
3128 })?;
3129 pairs.push((
3130 origin,
3131 NeighborHit {
3132 node_id,
3133 edge_id,
3134 relation,
3135 weight,
3136 name: None,
3137 kind: None,
3138 entity_type: None,
3139 },
3140 ));
3141 }
3142 Ok(pairs)
3143 })
3144 .await;
3145 let pairs = match chunk_result {
3146 Ok(pairs) => pairs,
3147 Err(e) => {
3148 report_graph_usage(&counted_queries, &counted_rows);
3149 return Err(e);
3150 }
3151 };
3152 result.extend(pairs);
3153 }
3154 report_graph_usage(&counted_queries, &counted_rows);
3155
3156 let requested: HashSet<Uuid> = unique_sources.iter().copied().collect();
3157 let mut grouped: HashMap<Uuid, Vec<NeighborHit>> =
3158 HashMap::with_capacity(unique_sources.len());
3159 for (origin, hit) in result {
3160 if !requested.contains(&origin) {
3161 return Err(StorageError::Internal(format!(
3162 "batch_neighbors returned unrequested origin {origin}"
3163 )));
3164 }
3165 grouped.entry(origin).or_default().push(hit);
3166 }
3167
3168 for hits in grouped.values_mut() {
3169 hits.sort_by(|a, b| {
3170 b.weight
3171 .partial_cmp(&a.weight)
3172 .unwrap_or(std::cmp::Ordering::Equal)
3173 .then(a.node_id.cmp(&b.node_id))
3174 .then(a.edge_id.cmp(&b.edge_id))
3175 });
3176 }
3177
3178 let mut ordered = Vec::new();
3179 for &source in sources {
3180 if let Some(hits) = grouped.get(&source) {
3181 ordered.extend(hits.iter().cloned().map(|hit| (source, hit)));
3182 }
3183 }
3184 Ok(ordered)
3185 }
3186
3187 async fn delete_edge(&self, id: LinkId, mode: DeleteMode) -> Result<bool, StorageError> {
3188 let id = Uuid::from(id);
3189 let statement = match mode {
3190 DeleteMode::Soft => {
3191 edge_soft_delete_statement(id, chrono::Utc::now().timestamp_micros())
3192 }
3193 DeleteMode::Hard => edge_hard_delete_statement(id),
3194 };
3195 self.with_writer("delete_edge", move |conn| {
3196 let mut stmt = conn.prepare(&statement.sql)?;
3197 bind_params(&mut stmt, &statement.params)?;
3198 Ok(stmt.raw_execute()? > 0)
3199 })
3200 .await
3201 }
3202
3203 async fn query_edges(
3204 &self,
3205 filter: EdgeFilter,
3206 sort: Vec<SortOrder<EdgeSortField>>,
3207 page: PageRequest,
3208 ) -> Result<Page<Edge>, StorageError> {
3209 let namespace = self.namespace.clone();
3210 let limit_i64 = i64::from(page.limit);
3211 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
3212 capability: StorageCapability::Graph,
3213 operation: "query_edges".into(),
3214 message: format!(
3215 "PageRequest: offset must be <= i64::MAX, got {}",
3216 page.offset
3217 ),
3218 })?;
3219 self.with_reader("query_edges", move |conn| {
3220 let (where_clause, mut all_params) = build_edge_filter_sql(&namespace, &filter);
3221 let order_clause = edge_order_clause(&sort);
3222 all_params.push(Box::new(limit_i64));
3223 all_params.push(Box::new(offset_i64));
3224
3225 let limit_idx = all_params.len() - 1;
3226 let offset_idx = all_params.len();
3227
3228 let data_sql = format!(
3229 "SELECT namespace, id, source_id, target_id, relation, weight, \
3230 created_at, updated_at, deleted_at, metadata, target_backend \
3231 FROM graph_edges{}{} LIMIT ?{} OFFSET ?{}",
3232 where_clause, order_clause, limit_idx, offset_idx,
3233 );
3234
3235 let mut stmt = conn.prepare(&data_sql)?;
3236 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3237 all_params.iter().map(|p| p.as_ref()).collect();
3238 let rows = stmt.query_map(param_refs.as_slice(), read_edge)?;
3239
3240 let mut items = Vec::new();
3241 for row in rows {
3242 items.push(row?);
3243 }
3244
3245 Ok(Page { items, total: None })
3246 })
3247 .await
3248 }
3249
3250 async fn count_edges(&self, filter: EdgeFilter) -> Result<u64, StorageError> {
3251 let namespace = self.namespace.clone();
3252 self.with_reader("count_edges", move |conn| {
3253 let (where_clause, params) = build_edge_filter_sql(&namespace, &filter);
3254 let sql = format!(
3255 "SELECT COUNT(*) FROM graph_edges{}",
3256 with_live_endpoints(&where_clause)
3257 );
3258 let mut stmt = conn.prepare(&sql)?;
3259 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3260 params.iter().map(|p| p.as_ref()).collect();
3261 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
3262 Ok(count as u64)
3263 })
3264 .await
3265 }
3266
3267 async fn count_edges_in_namespaces(
3268 &self,
3269 namespaces: &[String],
3270 filter: EdgeFilter,
3271 ) -> Result<u64, StorageError> {
3272 let namespaces: Vec<String> = namespaces
3273 .iter()
3274 .cloned()
3275 .collect::<HashSet<_>>()
3276 .into_iter()
3277 .collect();
3278 self.with_reader("count_edges_in_namespaces", move |conn| {
3279 let mut total = 0;
3280 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
3281 let (where_clause, params) = build_edge_filter_sql_for_namespaces(chunk, &filter);
3282 let sql = format!(
3283 "SELECT COUNT(*) FROM graph_edges{}",
3284 with_live_endpoints(&where_clause)
3285 );
3286 let mut stmt = conn.prepare(&sql)?;
3287 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3288 params.iter().map(|p| p.as_ref()).collect();
3289 let count: i64 = stmt.query_row(param_refs.as_slice(), |row| row.get(0))?;
3290 total += count as u64;
3291 }
3292 Ok(total)
3293 })
3294 .await
3295 }
3296
3297 async fn query_edges_in_namespaces(
3298 &self,
3299 namespaces: &[String],
3300 filter: EdgeFilter,
3301 sort: Vec<SortOrder<EdgeSortField>>,
3302 page: PageRequest,
3303 ) -> Result<Page<Edge>, StorageError> {
3304 let namespaces: Vec<String> = {
3319 let mut seen = HashSet::new();
3320 namespaces
3321 .iter()
3322 .filter(|ns| seen.insert((*ns).clone()))
3323 .cloned()
3324 .collect()
3325 };
3326 let limit_i64 = i64::from(page.limit);
3327 let offset_i64 = i64::try_from(page.offset).map_err(|_| StorageError::InvalidInput {
3328 capability: StorageCapability::Graph,
3329 operation: "query_edges_in_namespaces".into(),
3330 message: format!(
3331 "PageRequest: offset must be <= i64::MAX, got {}",
3332 page.offset
3333 ),
3334 })?;
3335 self.with_reader("query_edges_in_namespaces", move |conn| {
3336 let namespaces_json = serde_json::to_string(&namespaces)
3337 .map_err(|error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)))?;
3338
3339 let (where_clause, mut all_params) =
3340 build_edge_filter_sql_for_namespaces_json(&namespaces_json, &filter);
3341 let order_clause = edge_order_clause(&sort);
3342 all_params.push(Box::new(limit_i64));
3343 all_params.push(Box::new(offset_i64));
3344
3345 let limit_idx = all_params.len() - 1;
3346 let offset_idx = all_params.len();
3347
3348 let data_sql = format!(
3349 "SELECT namespace, id, source_id, target_id, relation, weight, \
3350 created_at, updated_at, deleted_at, metadata, target_backend \
3351 FROM graph_edges{}{} LIMIT ?{} OFFSET ?{}",
3352 where_clause, order_clause, limit_idx, offset_idx,
3353 );
3354
3355 let mut stmt = conn.prepare(&data_sql)?;
3356 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3357 all_params.iter().map(|p| p.as_ref()).collect();
3358 let rows = stmt.query_map(param_refs.as_slice(), read_edge)?;
3359
3360 let mut items = Vec::new();
3361 for row in rows {
3362 items.push(row?);
3363 }
3364
3365 Ok(Page { items, total: None })
3366 })
3367 .await
3368 }
3369
3370 async fn count_edges_by_relation(&self) -> Result<Vec<(EdgeRelation, u64)>, StorageError> {
3371 let namespace = self.namespace.clone();
3372 self.with_reader("count_edges_by_relation", move |conn| {
3373 let sql = format!(
3374 "SELECT relation, COUNT(*) FROM graph_edges \
3375 WHERE namespace = ?1 AND deleted_at IS NULL AND {LIVE_ENDPOINTS_CONDITION} \
3376 GROUP BY relation"
3377 );
3378 let mut stmt = conn.prepare(&sql)?;
3379 let rows = stmt.query_map([&namespace], |row| {
3380 let relation_str: String = row.get(0)?;
3381 let count: i64 = row.get(1)?;
3382 Ok((relation_str, count))
3383 })?;
3384 let mut out = Vec::new();
3385 for row in rows {
3386 let (relation_str, count) = row?;
3387 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
3388 rusqlite::Error::FromSqlConversionFailure(
3389 0,
3390 rusqlite::types::Type::Text,
3391 Box::new(e),
3392 )
3393 })?;
3394 out.push((relation, count as u64));
3395 }
3396 Ok(out)
3397 })
3398 .await
3399 }
3400
3401 async fn count_edges_by_relation_in_namespaces(
3402 &self,
3403 namespaces: &[String],
3404 ) -> Result<Vec<(EdgeRelation, u64)>, StorageError> {
3405 let namespaces: Vec<String> = namespaces
3406 .iter()
3407 .cloned()
3408 .collect::<HashSet<_>>()
3409 .into_iter()
3410 .collect();
3411 self.with_reader("count_edges_by_relation_in_namespaces", move |conn| {
3412 let mut totals = HashMap::new();
3413 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
3414 let (where_clause, params) =
3415 build_edge_filter_sql_for_namespaces(chunk, &EdgeFilter::default());
3416 let sql = format!(
3417 "SELECT relation, COUNT(*) FROM graph_edges{} GROUP BY relation",
3418 with_live_endpoints(&where_clause)
3419 );
3420 let mut stmt = conn.prepare(&sql)?;
3421 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3422 params.iter().map(|p| p.as_ref()).collect();
3423 let rows = stmt.query_map(param_refs.as_slice(), |row| {
3424 let relation_str: String = row.get(0)?;
3425 let count: i64 = row.get(1)?;
3426 Ok((relation_str, count))
3427 })?;
3428 for row in rows {
3429 let (relation_str, count) = row?;
3430 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
3431 rusqlite::Error::FromSqlConversionFailure(
3432 0,
3433 rusqlite::types::Type::Text,
3434 Box::new(e),
3435 )
3436 })?;
3437 *totals.entry(relation).or_insert(0) += count as u64;
3438 }
3439 }
3440 Ok(totals.into_iter().collect())
3441 })
3442 .await
3443 }
3444
3445 async fn count_edges_by_endpoint_base(&self) -> Result<EdgeEndpointBaseCounts, StorageError> {
3446 let namespace = self.namespace.clone();
3447 self.with_reader("count_edges_by_endpoint_base", move |conn| {
3448 let source_case = endpoint_base_case("source_id");
3449 let target_case = endpoint_base_case("target_id");
3450 let sql = format!(
3451 "SELECT {source_case} AS source_base, {target_case} AS target_base, COUNT(*) \
3452 FROM graph_edges \
3453 WHERE namespace = ?1 AND deleted_at IS NULL AND {LIVE_ENDPOINTS_CONDITION} \
3454 GROUP BY source_base, target_base"
3455 );
3456 let mut stmt = conn.prepare(&sql)?;
3457 let rows = stmt.query_map([&namespace], |row| {
3458 let source: String = row.get(0)?;
3459 let target: String = row.get(1)?;
3460 let count: i64 = row.get(2)?;
3461 Ok((source, target, count))
3462 })?;
3463 let mut counts = EdgeEndpointBaseCounts::default();
3464 for row in rows {
3465 let (source, target, count) = row?;
3466 fold_endpoint_base_row(&mut counts, &source, &target, count as u64);
3467 }
3468 Ok(counts)
3469 })
3470 .await
3471 }
3472
3473 async fn count_edges_by_endpoint_base_in_namespaces(
3474 &self,
3475 namespaces: &[String],
3476 ) -> Result<EdgeEndpointBaseCounts, StorageError> {
3477 let namespaces: Vec<String> = namespaces
3478 .iter()
3479 .cloned()
3480 .collect::<HashSet<_>>()
3481 .into_iter()
3482 .collect();
3483 self.with_reader("count_edges_by_endpoint_base_in_namespaces", move |conn| {
3484 let source_case = endpoint_base_case("source_id");
3485 let target_case = endpoint_base_case("target_id");
3486 let mut counts = EdgeEndpointBaseCounts::default();
3487 for chunk in namespaces.chunks(NAMESPACE_COUNT_CHUNK_SIZE) {
3488 let (where_clause, params) =
3489 build_edge_filter_sql_for_namespaces(chunk, &EdgeFilter::default());
3490 let sql = format!(
3491 "SELECT {source_case} AS source_base, {target_case} AS target_base, COUNT(*) \
3492 FROM graph_edges{} GROUP BY source_base, target_base",
3493 with_live_endpoints(&where_clause)
3494 );
3495 let mut stmt = conn.prepare(&sql)?;
3496 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3497 params.iter().map(|p| p.as_ref()).collect();
3498 let rows = stmt.query_map(param_refs.as_slice(), |row| {
3499 let source: String = row.get(0)?;
3500 let target: String = row.get(1)?;
3501 let count: i64 = row.get(2)?;
3502 Ok((source, target, count))
3503 })?;
3504 for row in rows {
3505 let (source, target, count) = row?;
3506 fold_endpoint_base_row(&mut counts, &source, &target, count as u64);
3507 }
3508 }
3509 Ok(counts)
3510 })
3511 .await
3512 }
3513
3514 async fn query_edges_after(
3515 &self,
3516 filter: EdgeFilter,
3517 after: Option<Uuid>,
3518 limit: u32,
3519 ) -> Result<EdgeSeekPage, StorageError> {
3520 let namespace = self.namespace.clone();
3521 let limit_usize = limit as usize;
3522 let probe_limit_i64 = i64::from(limit) + 1;
3523 self.with_reader("query_edges_after", move |conn| {
3524 let (mut where_clause, mut params) = build_edge_filter_sql(&namespace, &filter);
3525 if let Some(cursor) = after {
3526 params.push(Box::new(cursor.to_string()));
3527 where_clause.push_str(&format!(" AND id > ?{}", params.len()));
3528 }
3529 params.push(Box::new(probe_limit_i64));
3530 let limit_idx = params.len();
3531
3532 let data_sql = format!(
3538 "SELECT namespace, id, source_id, target_id, relation, weight, \
3539 created_at, updated_at, deleted_at, metadata, target_backend \
3540 FROM graph_edges{} ORDER BY id ASC LIMIT ?{}",
3541 where_clause, limit_idx,
3542 );
3543
3544 let mut stmt = conn.prepare(&data_sql)?;
3545 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3546 params.iter().map(|p| p.as_ref()).collect();
3547 let rows = stmt.query_map(param_refs.as_slice(), read_edge)?;
3548
3549 let mut items = Vec::new();
3550 for row in rows {
3551 items.push(row?);
3552 }
3553 let has_more = items.len() > limit_usize;
3554 if has_more {
3555 items.truncate(limit_usize);
3556 }
3557 let next_after = if has_more {
3558 items.last().map(|e| Uuid::from(e.id))
3559 } else {
3560 None
3561 };
3562
3563 Ok(EdgeSeekPage { items, next_after })
3564 })
3565 .await
3566 }
3567
3568 async fn query_edges_sequence_after(
3569 &self,
3570 filter: EdgeFilter,
3571 after: Option<SeekCursor>,
3572 limit: u32,
3573 ) -> Result<SeekPage<Edge>, StorageError> {
3574 if limit == 0 {
3575 return Ok(SeekPage::default());
3576 }
3577 let namespace = self.namespace.clone();
3578 let limit_usize = limit as usize;
3579 let probe_limit_i64 = i64::from(limit) + 1;
3580 self.with_reader("query_edges_sequence_after", move |conn| {
3581 let (mut where_clause, mut params) = build_edge_filter_sql(&namespace, &filter);
3582 if let Some(cursor) = after {
3583 params.push(Box::new(cursor.sequence));
3584 where_clause.push_str(&format!(" AND graph_edges_seq.seq > ?{}", params.len()));
3585 }
3586 params.push(Box::new(probe_limit_i64));
3587 let limit_idx = params.len();
3588 let sql = format!(
3591 "SELECT namespace, id, source_id, target_id, relation, weight, \
3592 created_at, updated_at, deleted_at, metadata, target_backend, \
3593 graph_edges_seq.seq \
3594 FROM graph_edges_seq CROSS JOIN graph_edges \
3595 ON graph_edges.id = graph_edges_seq.edge_id{where_clause} \
3596 ORDER BY graph_edges_seq.seq ASC LIMIT ?{limit_idx}"
3597 );
3598 let mut stmt = conn.prepare(&sql)?;
3599 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3600 params.iter().map(|param| param.as_ref()).collect();
3601 let rows = stmt.query_map(param_refs.as_slice(), |row| {
3602 Ok((read_edge(row)?, row.get::<_, i64>(11)?))
3603 })?;
3604 let mut entries = rows.collect::<Result<Vec<_>, _>>()?;
3605 let has_more = entries.len() > limit_usize;
3606 if has_more {
3607 entries.truncate(limit_usize);
3608 }
3609 let next_after = if has_more {
3610 entries.last().map(|(edge, sequence)| SeekCursor {
3611 sequence: *sequence,
3612 id: Uuid::from(edge.id),
3613 })
3614 } else {
3615 None
3616 };
3617 let items = entries.into_iter().map(|(edge, _)| edge).collect();
3618 Ok(SeekPage { items, next_after })
3619 })
3620 .await
3621 }
3622
3623 async fn neighbors(
3624 &self,
3625 node_id: Uuid,
3626 query: NeighborQuery,
3627 ) -> Result<Vec<NeighborHit>, StorageError> {
3628 self.query_neighbors_page("neighbors", node_id, query, None, None)
3629 .await
3630 }
3631
3632 async fn neighbors_page(
3633 &self,
3634 node_id: Uuid,
3635 query: NeighborQuery,
3636 after: Option<NeighborCursor>,
3637 neighbor_kinds: Option<Vec<String>>,
3638 ) -> Result<Vec<NeighborHit>, StorageError> {
3639 self.query_neighbors_page("neighbors_page", node_id, query, after, neighbor_kinds)
3640 .await
3641 }
3642
3643 async fn neighbors_both_directions(
3648 &self,
3649 node_id: Uuid,
3650 query: NeighborQuery,
3651 ) -> Result<Vec<DirectedNeighborHit>, StorageError> {
3652 count_neighbor_select();
3653
3654 let namespace = self.namespace.clone();
3655 let node_str = node_id.to_string();
3656
3657 let counted_queries = Arc::new(std::sync::atomic::AtomicU64::new(0));
3658 let counted_rows = Arc::new(std::sync::atomic::AtomicU64::new(0));
3659 let closure_queries = Arc::clone(&counted_queries);
3660 let closure_rows = Arc::clone(&counted_rows);
3661 let result = self
3662 .with_reader("neighbors_both_directions", move |conn| {
3663 let base_out = "SELECT target_id AS node_id, id AS edge_id, relation, weight, \
3664 'out' AS dir \
3665 FROM graph_edges \
3666 WHERE namespace = ?1 AND source_id = ?2 AND deleted_at IS NULL";
3667 let base_in = "SELECT source_id AS node_id, id AS edge_id, relation, weight, \
3668 'in' AS dir \
3669 FROM graph_edges \
3670 WHERE namespace = ?1 AND target_id = ?2 AND deleted_at IS NULL";
3671 let sql = format!("{} UNION ALL {}", base_out, base_in);
3672
3673 let (where_extra, limit_clause, extra_params) =
3674 neighbor_extra_clause(&query, 3, None, None);
3675
3676 let full_sql = format!(
3685 "SELECT node_id, edge_id, relation, weight, dir FROM ({}){} \
3686 ORDER BY weight DESC, node_id ASC, \
3687 CASE dir WHEN 'out' THEN 0 ELSE 1 END ASC, edge_id ASC{}",
3688 sql, where_extra, limit_clause
3689 );
3690
3691 let mut stmt = conn.prepare(&full_sql)?;
3692
3693 let mut all_params: Vec<Box<dyn rusqlite::types::ToSql>> = Vec::new();
3694 all_params.push(Box::new(namespace.clone()));
3695 all_params.push(Box::new(node_str.clone()));
3696 all_params.extend(extra_params);
3697
3698 let param_refs: Vec<&dyn rusqlite::types::ToSql> =
3699 all_params.iter().map(|p| p.as_ref()).collect();
3700
3701 closure_queries.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
3702 let rows = stmt.query_map(param_refs.as_slice(), |row| {
3703 let nid_str: String = row.get(0)?;
3704 let eid_str: String = row.get(1)?;
3705 let relation_str: String = row.get(2)?;
3706 let weight: f64 = row.get(3)?;
3707 let dir_str: String = row.get(4)?;
3708 Ok((nid_str, eid_str, relation_str, weight, dir_str))
3709 })?;
3710
3711 let mut hits = Vec::new();
3712 for row in rows {
3713 let (nid_str, eid_str, relation_str, weight, dir_str) = row?;
3714 closure_rows.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
3715 let relation = relation_str.parse::<EdgeRelation>().map_err(|e| {
3716 rusqlite::Error::FromSqlConversionFailure(
3717 2,
3718 rusqlite::types::Type::Text,
3719 Box::new(e),
3720 )
3721 })?;
3722 let direction = if dir_str == "out" {
3723 Direction::Out
3724 } else {
3725 Direction::In
3726 };
3727 hits.push(DirectedNeighborHit {
3728 hit: NeighborHit {
3729 node_id: parse_uuid(&nid_str)?,
3730 edge_id: parse_uuid(&eid_str)?,
3731 relation,
3732 weight,
3733 name: None,
3734 kind: None,
3735 entity_type: None,
3736 },
3737 direction,
3738 });
3739 }
3740
3741 Ok(hits)
3742 })
3743 .await;
3744
3745 report_graph_usage(&counted_queries, &counted_rows);
3746 result
3747 }
3748
3749 async fn traverse(&self, request: TraversalRequest) -> Result<Vec<GraphPath>, StorageError> {
3750 request
3751 .validate()
3752 .map_err(|message| StorageError::InvalidInput {
3753 capability: StorageCapability::Graph,
3754 operation: "traverse".into(),
3755 message,
3756 })?;
3757 if request.roots.is_empty() {
3758 return Ok(Vec::new());
3759 }
3760
3761 let budget = request.execution_budget;
3762 let outer_context = khive_storage::capture_request_read_context();
3763 let execution_deadline =
3764 khive_storage::RequestReadDeadline::after(budget.remaining_duration());
3765 khive_storage::scope_request_read_deadline_at(execution_deadline, async {
3766 if outer_context.stop_reason().is_some() {
3767 return Err(StorageError::Timeout {
3768 operation: "traverse".into(),
3769 });
3770 }
3771 if budget.is_expired() {
3772 return Err(traversal_timeout_error(&budget));
3773 }
3774 let mut distinct_roots = HashSet::with_capacity(request.roots.len());
3775 let roots = request
3776 .roots
3777 .iter()
3778 .copied()
3779 .filter(|root| distinct_roots.insert(*root))
3780 .collect::<Vec<_>>();
3781 let opts = request.options;
3782 if self
3783 .index_repair
3784 .as_ref()
3785 .is_some_and(super::index_repair::IndexRepairContext::is_writable)
3786 {
3787 let directions = match opts.direction {
3793 Direction::Out => vec![Direction::Out],
3794 Direction::In => vec![Direction::In],
3795 Direction::Both => vec![Direction::Out, Direction::In],
3796 };
3797 let relation_count = opts.relations.as_ref().map_or(0, Vec::len);
3798 let statements = directions
3799 .into_iter()
3800 .map(|direction| {
3801 traversal_neighbor_sql(direction, relation_count, opts.min_weight.is_some())
3802 })
3803 .collect::<Vec<_>>();
3804 let preflight = self.with_indexed_reader("traverse", move |conn| {
3807 for sql in &statements {
3808 let mut statement = conn.prepare(sql)?;
3809 statement.raw_bind_parameter(3, 0_i64)?;
3810 let _ = statement.raw_query().next()?;
3811 }
3812 Ok(())
3813 });
3814 let preflight_result = tokio::select! {
3815 biased;
3816 _ = outer_context.clone().wait_for_stop() => Err(StorageError::Timeout {
3817 operation: "traverse".into(),
3818 }),
3819 result = tokio::time::timeout_at(execution_deadline.async_at(), preflight) => {
3820 match result {
3821 Ok(result) => result,
3822 Err(_) => Err(traversal_timeout_error(&budget)),
3823 }
3824 },
3825 };
3826 preflight_result.map_err(|error| match error {
3827 StorageError::Timeout { .. }
3828 if budget.is_expired() && outer_context.stop_reason().is_none() =>
3829 {
3830 traversal_timeout_error(&budget)
3831 }
3832 error => error,
3833 })?;
3834 if budget.is_expired() {
3835 return Err(traversal_timeout_error(&budget));
3836 }
3837 }
3838 let include_roots = request.include_roots;
3839 let namespace = self.namespace.clone();
3840 let origin = self.pool.origin();
3841 let closure_budget = budget.clone();
3842 let counted_rows = Arc::new(std::sync::atomic::AtomicU64::new(0));
3848 let counted_queries = Arc::new(std::sync::atomic::AtomicU64::new(0));
3849 let closure_rows = Arc::clone(&counted_rows);
3850 let closure_queries = Arc::clone(&counted_queries);
3851 let result = self
3852 .with_reader("traverse", move |conn| {
3853 Ok(run_bounded_traversal(
3854 conn,
3855 roots,
3856 opts,
3857 include_roots,
3858 namespace,
3859 origin,
3860 closure_budget,
3861 closure_rows.as_ref(),
3862 closure_queries.as_ref(),
3863 ))
3864 })
3865 .await
3866 .and_then(|inner| inner)
3867 .map_err(|error| match error {
3868 StorageError::Timeout { .. }
3869 if budget.is_expired() && outer_context.stop_reason().is_none() =>
3870 {
3871 traversal_timeout_error(&budget)
3872 }
3873 error => error,
3874 });
3875
3876 khive_storage::usage::count(
3885 khive_storage::usage::UsageUnit::DbRoundTrips,
3886 counted_queries.load(std::sync::atomic::Ordering::Relaxed),
3887 );
3888 khive_storage::usage::count(
3889 khive_storage::usage::UsageUnit::GraphHops,
3890 counted_rows.load(std::sync::atomic::Ordering::Relaxed),
3891 );
3892
3893 result
3894 })
3895 .await
3896 }
3897
3898 async fn purge_incident_edges(&self, node_id: Uuid) -> Result<u64, StorageError> {
3899 let statement = purge_incident_edges_statement(node_id);
3903 self.with_writer("purge_incident_edges", move |conn| {
3904 let mut stmt = conn.prepare(&statement.sql)?;
3905 bind_params(&mut stmt, &statement.params)?;
3906 Ok(stmt.raw_execute()? as u64)
3907 })
3908 .await
3909 }
3910}
3911
3912const GRAPH_DDL: &str = include_str!("../../sql/graph-ddl.sql");
3917
3918pub(crate) fn ensure_graph_schema(conn: &rusqlite::Connection) -> Result<(), rusqlite::Error> {
3919 conn.execute_batch(GRAPH_DDL)
3920}
3921
3922#[cfg(test)]
3923#[path = "graph_annotation_tests.rs"]
3924mod annotation_tests;
3925
3926#[cfg(test)]
3927#[path = "graph_tests.rs"]
3928mod tests;
3929
3930#[cfg(test)]
3931#[path = "graph_index_repair_tests.rs"]
3932mod index_repair_tests;
3933
3934#[cfg(test)]
3935#[path = "graph_busy_tests.rs"]
3936mod direct_busy_tests;